dApp Docs/AI Agent 跨链部署与IBC通信指南
Development reference. Not independently verified for production.

AI Agent 跨链部署与IBC通信指南 — MSG Chain

适用链:msg-chain-1 · Bech32 前缀:msg
⚠️ No-Go Disclaimer: MSGChain 主网裁决为 No-Go。本文件所有内容反映的是开发阶段的技术设计,不代表主网未独立核验上线状态。生产部署状态请以白皮书为准:https://msgchain.org/whitepaper/


目录

  1. 概述
  2. 跨链Agent身份
  3. 跨链Agent注册与发现
  4. 跨链A2A通信
  5. 跨链支付结算
  6. 跨链Agent编排
  7. 安全考虑

1. 概述

1.1 为什么AI Agent需要跨链能力

AI Agent 在单一区块链上运行受限于该链的流动性、用户群体和计算资源。跨链能力允许 Agent:

MSG Chain 通过 IBC(Inter-Blockchain Communication)协议连接 Cosmos 生态中的多条链,为 AI Agent 提供原生的跨链基础设施。

1.2 IBC 基础

IBC 是 Cosmos 生态的核心互操作协议。AI Agent 使用 IBC 进行跨链通信时,涉及以下核心概念:

// IBC 核心概念类型定义
interface IbcConnection {
    clientId: string;           // 轻客户端 ID,如 "07-tendermint-0"
    connectionId: string;       // 连接 ID,如 "connection-0"
    counterpartyChainId: string; // 对方链 ID
    state: 'INIT' | 'TRYOPEN' | 'OPEN';
}

interface IbcChannel {
    portId: string;             // 端口 ID,如 "a2a"、"transfer"
    channelId: string;          // 通道 ID,如 "channel-0"
    connectionId: string;       // 所属连接
    ordering: 'UNORDERED' | 'ORDERED';
    state: 'INIT' | 'TRYOPEN' | 'OPEN' | 'CLOSED';
}

interface IbcPacket {
    sequence: number;
    sourcePort: string;
    sourceChannel: string;
    destinationPort: string;
    destinationChannel: string;
    data: Uint8Array;
    timeoutHeight?: Height;
    timeoutTimestamp?: number;
}

1.3 MSG Chain IBC 配置

在 msg-chain-1 上,IBC 的初始化配置如下:

// msg-chain-1 IBC 配置
pub const MSG_CHAIN_IBC_CONFIG: IbcConfig = IbcConfig {
    chain_id: "msg-chain-1",
    bech32_prefix: "msg",
    transfer_port: "transfer",
    a2a_port: "a2a",
    agent_registry_port: "agent-registry",
    ibc_version: "ics-20-v2",
    default_timeout_seconds: 600,
};

1.4 跨链身份概述

AI Agent 在跨链环境中需要一个全局唯一的身份标识。DID(Decentralized Identifier)提供了与链无关的身份层,Agent 可以在每条链上拥有一个本地地址,同时通过 DID 文档将这些地址关联到同一个全局身份。

┌─────────────────────────────────────────────────┐
│             跨链 AI Agent 身份层                  │
├─────────────┬─────────────┬─────────────────────┤
│  msg-chain-1 │   chain-a   │      chain-b        │
│  msg1abc...  │  cosmos1... │    osmo1...         │
└─────────────┴─────────────┴─────────────────────┘
        │              │              │
        └──────────────┴──────────────┘
                   │
            DID Document
            did:msg:agent:abc123

1.5 适用场景

跨链 AI Agent 的典型应用包括:

场景 描述 涉及链数
跨链套利 监控多条链上的 DEX 价差并执行原子交易 2+
跨链预言机 聚合多条链上的数据并提供统一喂价 3+
跨链借贷 在一条链上抵押资产,在另一条链上借款 2
跨链治理 Agent 代表用户在多条链上投票 3+
跨链自动化 根据一条链上的事件触发另一条链上的操作 2

2. 跨链Agent身份

2.1 DID 架构设计

AI Agent 在 MSG Chain 上的身份体系采用 W3C DID 标准,扩展了跨链字段以支持多链地址映射。

// ============================================================
// 跨链 DID 核心类型 (TypeScript)
// ============================================================

// 支持的验证方法类型
type VerificationMethodType = 'EcdsaSecp256k1VerificationKey2019' | 'Ed25519VerificationKey2020';

// 验证方法
interface VerificationMethod {
    id: string;                           // did:msg:agent:<id>#keys-1
    type: VerificationMethodType;
    controller: string;                   // did:msg:agent:<id>
    publicKeyMultibase: string;           // 多格式编码的公钥
}

// 跨链地址映射
interface ChainAddress {
    chainId: string;                      // "msg-chain-1"
    address: string;                      // "msg1..."
    network: 'mainnet' | 'testnet';
    proof?: string;                       // 签名证明,证明地址所有权
}

// 跨链 DID 主接口
interface CrossChainDID {
    id: string;                           // "did:msg:agent:<unique-id>"
    controller: string | string[];
    alsoKnownAs: string[];                // 跨链别名列表
    
    // 多链地址映射 chain_id -> address
    chainAddresses: Map<string, ChainAddress>;
    
    // 验证方法列表
    verificationMethod: VerificationMethod[];
    
    // 认证关系
    authentication: string[];
    assertionMethod: string[];
    keyAgreement: string[];
    capabilityInvocation: string[];
    capabilityDelegation: string[];
    
    // IBC 相关元数据
    ibcEndpoints: IbcEndpoint[];
    
    // 元数据
    created: string;                      // ISO 8601
    updated: string;                      // ISO 8601
    version: number;
}

// IBC 端点
interface IbcEndpoint {
    chainId: string;
    portId: string;
    channelId: string;
    connectionId: string;
}

2.2 DID 文档生成

Agent 首次部署时,生成 DID 文档并在 MSG Chain 上注册。

// ============================================================
// DID 文档生成器 (TypeScript)
// ============================================================

import { Secp256k1, sha256, randomBytes } from './crypto';
import { bech32 } from 'bech32';

class DidGenerator {
    /**
     * 为 Agent 生成唯一的 DID
     */
    static generateDid(): string {
        const entropy = randomBytes(16);
        const hash = sha256(entropy);
        const id = Buffer.from(hash).toString('hex').substring(0, 32);
        return `did:msg:agent:${id}`;
    }

    /**
     * 生成 Agent 跨链密钥对
     */
    static generateKeyPair(): { privateKey: Uint8Array; publicKey: Uint8Array } {
        return Secp256k1.generateKeyPair();
    }

    /**
     * 从公钥生成 MSG Chain 地址
     */
    static generateMsgAddress(publicKey: Uint8Array): string {
        const hash = sha256(publicKey);
        const twentyBytes = hash.slice(0, 20);
        return bech32.encode('msg', bech32.toWords(twentyBytes));
    }

    /**
     * 创建完整的 DID 文档
     */
    static createDidDocument(
        agentId: string,
        keyPair: { privateKey: Uint8Array; publicKey: Uint8Array },
        chains: string[]
    ): CrossChainDID {
        const did = `did:msg:agent:${agentId}`;
        const pubKeyMultibase = Buffer.from(keyPair.publicKey).toString('base64');
        
        const vm: VerificationMethod = {
            id: `${did}#keys-1`,
            type: 'EcdsaSecp256k1VerificationKey2019',
            controller: did,
            publicKeyMultibase: pubKeyMultibase,
        };

        const chainAddresses = new Map<string, ChainAddress>();
        
        for (const chainId of chains) {
            const derivationPath = `m/44'/118'/${chains.indexOf(chainId)}'/0/0`;
            const chainPubKey = this.deriveChildKey(keyPair.privateKey, derivationPath);
            const address = this.generateMsgAddress(chainPubKey);
            
            chainAddresses.set(chainId, {
                chainId,
                address,
                network: 'mainnet',
                proof: this.signAddressClaim(keyPair.privateKey, chainId, address),
            });
        }

        return {
            id: did,
            controller: did,
            alsoKnownAs: [],
            chainAddresses,
            verificationMethod: [vm],
            authentication: [`${did}#keys-1`],
            assertionMethod: [`${did}#keys-1`],
            keyAgreement: [`${did}#keys-1`],
            capabilityInvocation: [`${did}#keys-1`],
            capabilityDelegation: [`${did}#keys-1`],
            ibcEndpoints: [],
            created: new Date().toISOString(),
            updated: new Date().toISOString(),
            version: 1,
        };
    }

    /**
     * 派生子密钥(HD 钱包路径)
     */
    private static deriveChildKey(
        parentKey: Uint8Array,
        path: string
    ): Uint8Array {
        return parentKey;
    }

    /**
     * 签名地址声明,证明地址所有权
     */
    private static signAddressClaim(
        privateKey: Uint8Array,
        chainId: string,
        address: string
    ): string {
        const message = `Claim: ${chainId} -> ${address}`;
        const signature = Secp256k1.sign(sha256(Buffer.from(message)), privateKey);
        return Buffer.from(signature).toString('hex');
    }
}

2.3 跨链密钥管理

跨链密钥管理是 Agent 安全的基础。Agent 在每条链上可以使用相同的根密钥派生不同地址,也可以使用独立的密钥对。

// ============================================================
// 跨链密钥管理器 (TypeScript)
// ============================================================

interface KeyDescriptor {
    chainId: string;
    derivationPath: string;
    purpose: 'authentication' | 'payment' | 'governance' | 'encryption';
}

class CrossChainKeyManager {
    private rootKey: Uint8Array;
    private keysCache: Map<string, { publicKey: Uint8Array; privateKey: Uint8Array }>;

    constructor(rootKey: Uint8Array) {
        this.rootKey = rootKey;
        this.keysCache = new Map();
    }

    /**
     * 为指定链派生密钥
     */
    deriveKey(descriptor: KeyDescriptor): { publicKey: Uint8Array; privateKey: Uint8Array } {
        const cacheKey = `${descriptor.chainId}:${descriptor.purpose}`;
        
        if (this.keysCache.has(cacheKey)) {
            return this.keysCache.get(cacheKey)!;
        }

        const chainCode = sha256(Buffer.from(descriptor.chainId + descriptor.purpose));
        const privateKey = sha256(Buffer.concat([this.rootKey, chainCode]));
        const publicKey = Secp256k1.getPublicKey(privateKey);

        const keyPair = { publicKey, privateKey };
        this.keysCache.set(cacheKey, keyPair);
        return keyPair;
    }

    /**
     * 获取 MSG Chain 上的 Agent 地址
     */
    getMsgAddress(): string {
        const key = this.deriveKey({
            chainId: 'msg-chain-1',
            derivationPath: "m/44'/118'/0'/0/0",
            purpose: 'authentication',
        });
        return DidGenerator.generateMsgAddress(key.publicKey);
    }

    /**
     * 跨链签名 — 在指定链的上下文中签名消息
     */
    signForChain(chainId: string, message: Uint8Array): Uint8Array {
        const key = this.deriveKey({
            chainId,
            derivationPath: "m/44'/118'/0'/0/0",
            purpose: 'authentication',
        });
        return Secp256k1.sign(sha256(message), key.privateKey);
    }

    /**
     * 验证来自跨链 Agent 的签名
     */
    verifyCrossChainSignature(
        did: CrossChainDID,
        chainId: string,
        message: Uint8Array,
        signature: Uint8Array
    ): boolean {
        const chainAddr = did.chainAddresses.get(chainId);
        if (!chainAddr) {
            throw new Error(`No address for chain ${chainId}`);
        }

        const vmId = did.verificationMethod[0]?.id;
        if (!vmId) return false;

        const pubKey = Buffer.from(did.verificationMethod[0].publicKeyMultibase, 'base64');
        return Secp256k1.verify(sha256(message), signature, pubKey);
    }
}

2.4 IBC 启用的 DID 注册表

DID 注册表是一个部署在 MSG Chain 上的智能合约,管理 Agent 的 DID 文档并通过 IBC 与其他链同步。

// ============================================================
// 跨链 DID 注册表合约 (Rust - CosmWasm)
// ============================================================

use cosmwasm_std::{
    entry_point, Binary, Deps, DepsMut, Env, MessageInfo, Response, StdResult,
    IbcMsg, IbcPacket, IbcChannel, IbcEndpoint, IbcOrder, Storage,
};
use serde::{Deserialize, Serialize};

// ---------- 状态定义 ----------

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
pub struct DidDocument {
    pub id: String,
    pub controller: String,
    pub also_known_as: Vec<String>,
    pub verification_method: Vec<VerificationMethod>,
    pub authentication: Vec<String>,
    pub assertion_method: Vec<String>,
    pub key_agreement: Vec<String>,
    pub capability_invocation: Vec<String>,
    pub capability_delegation: Vec<String>,
    pub chain_addresses: Vec<ChainAddressEntry>,
    pub created: String,
    pub updated: String,
    pub version: u64,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
pub struct VerificationMethod {
    pub id: String,
    #[serde(rename = "type")]
    pub vm_type: String,
    pub controller: String,
    pub public_key_multibase: String,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
pub struct ChainAddressEntry {
    pub chain_id: String,
    pub address: String,
    pub network: String,
    pub proof: Option<String>,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
pub struct IbcSyncState {
    pub last_synced_sequence: u64,
    pub remote_chain_id: String,
    pub remote_registry_address: String,
}

const DID_REGISTRY: &str = "did_registry";
const IBC_SYNC_STATE: &str = "ibc_sync";
const CHAIN_INDEX: &str = "chain_index";

// ---------- 实例化 ----------

#[derive(Serialize, Deserialize)]
pub struct InstantiateMsg {
    pub admin: String,
    pub supported_chains: Vec<String>,
}

#[entry_point]
pub fn instantiate(
    deps: DepsMut,
    _env: Env,
    info: MessageInfo,
    msg: InstantiateMsg,
) -> StdResult<Response> {
    let config = RegistryConfig {
        admin: info.sender.to_string(),
        supported_chains: msg.supported_chains,
    };
    deps.storage.set(b"config", &bincode2::serialize(&config)?);
    Ok(Response::new()
        .add_attribute("action", "instantiate")
        .add_attribute("admin", info.sender))
}

// ---------- 执行消息 ----------

#[derive(Serialize, Deserialize)]
pub enum ExecuteMsg {
    RegisterAgent { did_document: DidDocument },
    UpdateAgent { agent_id: String, did_document: DidDocument },
    DeactivateAgent { agent_id: String },
    SyncDidToChain { target_chain: String, agent_id: String },
    AddChainAddress {
        agent_id: String,
        address_entry: ChainAddressEntry,
        proof: Binary,
    },
}

#[entry_point]
pub fn execute(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    msg: ExecuteMsg,
) -> StdResult<Response> {
    match msg {
        ExecuteMsg::RegisterAgent { did_document } => {
            execute_register_agent(deps, env, info, did_document)
        }
        ExecuteMsg::UpdateAgent { agent_id, did_document } => {
            execute_update_agent(deps, env, info, agent_id, did_document)
        }
        ExecuteMsg::DeactivateAgent { agent_id } => {
            execute_deactivate_agent(deps, env, info, agent_id)
        }
        ExecuteMsg::SyncDidToChain { target_chain, agent_id } => {
            execute_sync_did_to_chain(deps, env, info, target_chain, agent_id)
        }
        ExecuteMsg::AddChainAddress { agent_id, address_entry, proof } => {
            execute_add_chain_address(deps, env, info, agent_id, address_entry, proof)
        }
    }
}

fn execute_register_agent(
    deps: DepsMut,
    _env: Env,
    info: MessageInfo,
    doc: DidDocument,
) -> StdResult<Response> {
    let store = deps.storage;
    let key = format!("{}{}", DID_REGISTRY, doc.id);
    if store.get(key.as_bytes()).is_some() {
        return Err(cosmwasm_std::StdError::generic_err(
            format!("DID {} already registered", doc.id)
        ));
    }
    if doc.controller != info.sender.to_string() {
        return Err(cosmwasm_std::StdError::generic_err(
            "Sender must be the DID controller"
        ));
    }
    let serialized = bincode2::serialize(&doc)?;
    store.set(key.as_bytes(), &serialized);
    for addr in &doc.chain_addresses {
        let idx_key = format!("{}{}:{}", CHAIN_INDEX, addr.chain_id, addr.address);
        store.set(idx_key.as_bytes(), doc.id.as_bytes());
    }
    Ok(Response::new()
        .add_attribute("action", "register_agent")
        .add_attribute("did", doc.id)
        .add_attribute("sender", info.sender))
}

fn execute_update_agent(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    agent_id: String,
    new_doc: DidDocument,
) -> StdResult<Response> {
    let store = deps.storage;
    let did = format!("did:msg:agent:{}", agent_id);
    let key = format!("{}{}", DID_REGISTRY, did);
    let existing = store.get(key.as_bytes())
        .ok_or_else(|| cosmwasm_std::StdError::generic_err("Agent not found"))?;
    let mut current: DidDocument = bincode2::deserialize(&existing)?;
    if current.controller != info.sender.to_string() {
        return Err(cosmwasm_std::StdError::generic_err("Unauthorized"));
    }
    new_doc.version = current.version + 1;
    new_doc.created = current.created;
    new_doc.updated = env.block.time.to_string();
    let serialized = bincode2::serialize(&new_doc)?;
    store.set(key.as_bytes(), &serialized);
    Ok(Response::new()
        .add_attribute("action", "update_agent")
        .add_attribute("did", did)
        .add_attribute("version", new_doc.version.to_string()))
}

fn execute_deactivate_agent(
    deps: DepsMut,
    _env: Env,
    info: MessageInfo,
    agent_id: String,
) -> StdResult<Response> {
    let store = deps.storage;
    let did = format!("did:msg:agent:{}", agent_id);
    let key = format!("{}{}", DID_REGISTRY, did);
    let existing = store.get(key.as_bytes())
        .ok_or_else(|| cosmwasm_std::StdError::generic_err("Agent not found"))?;
    let doc: DidDocument = bincode2::deserialize(&existing)?;
    if doc.controller != info.sender.to_string() {
        return Err(cosmwasm_std::StdError::generic_err("Unauthorized"));
    }
    store.remove(key.as_bytes());
    Ok(Response::new()
        .add_attribute("action", "deactivate_agent")
        .add_attribute("did", did))
}

fn execute_sync_did_to_chain(
    deps: DepsMut,
    env: Env,
    _info: MessageInfo,
    target_chain: String,
    agent_id: String,
) -> StdResult<Response> {
    let store = deps.storage;
    let did = format!("did:msg:agent:{}", agent_id);
    let key = format!("{}{}", DID_REGISTRY, did);
    let existing = store.get(key.as_bytes())
        .ok_or_else(|| cosmwasm_std::StdError::generic_err("Agent not found"))?;
    let doc: DidDocument = bincode2::deserialize(&existing)?;
    let packet_data = IbcDidSyncPacket {
        did_document: doc,
        source_chain: "msg-chain-1".to_string(),
        timestamp: env.block.time,
    };
    let serialized_packet = bincode2::serialize(&packet_data)?;
    let ibc_msg = IbcMsg::SendPacket {
        channel_id: format!("channel-{}", get_channel_for_chain(&target_chain)),
        data: Binary::from(serialized_packet),
        timeout: IbcTimeout::with_timestamp(env.block.time.plus_seconds(600)),
    };
    Ok(Response::new()
        .add_message(ibc_msg)
        .add_attribute("action", "sync_did")
        .add_attribute("agent_id", agent_id)
        .add_attribute("target_chain", target_chain))
}

fn execute_add_chain_address(
    deps: DepsMut,
    _env: Env,
    info: MessageInfo,
    agent_id: String,
    address_entry: ChainAddressEntry,
    _proof: Binary,
) -> StdResult<Response> {
    let store = deps.storage;
    let did = format!("did:msg:agent:{}", agent_id);
    let key = format!("{}{}", DID_REGISTRY, did);
    let existing = store.get(key.as_bytes())
        .ok_or_else(|| cosmwasm_std::StdError::generic_err("Agent not found"))?;
    let mut doc: DidDocument = bincode2::deserialize(&existing)?;
    if doc.controller != info.sender.to_string() {
        return Err(cosmwasm_std::StdError::generic_err("Unauthorized"));
    }
    if address_entry.proof.is_none() {
        return Err(cosmwasm_std::StdError::generic_err("Proof required"));
    }
    doc.chain_addresses.push(address_entry.clone());
    doc.version += 1;
    let idx_key = format!("{}{}:{}", CHAIN_INDEX, address_entry.chain_id, address_entry.address);
    store.set(idx_key.as_bytes(), did.as_bytes());
    let serialized = bincode2::serialize(&doc)?;
    store.set(key.as_bytes(), &serialized);
    Ok(Response::new()
        .add_attribute("action", "add_chain_address")
        .add_attribute("agent_id", agent_id)
        .add_attribute("chain", &address_entry.chain_id)
        .add_attribute("address", &address_entry.address))
}

// ---------- 查询消息 ----------

#[derive(Serialize, Deserialize)]
pub enum QueryMsg {
    GetAgent { did: String },
    GetAgentByAddress { chain_id: String, address: String },
    GetChainAddresses { agent_id: String },
    ListAgents { start_after: Option<String>, limit: Option<u32> },
}

#[entry_point]
pub fn query(deps: Deps, _env: Env, msg: QueryMsg) -> StdResult<Binary> {
    match msg {
        QueryMsg::GetAgent { did } => query_get_agent(deps, did),
        QueryMsg::GetAgentByAddress { chain_id, address } => {
            query_get_agent_by_address(deps, chain_id, address)
        }
        QueryMsg::GetChainAddresses { agent_id } => {
            query_get_chain_addresses(deps, agent_id)
        }
        QueryMsg::ListAgents { start_after, limit } => {
            query_list_agents(deps, start_after, limit)
        }
    }
}

fn query_get_agent(deps: Deps, did: String) -> StdResult<Binary> {
    let key = format!("{}{}", DID_REGISTRY, did);
    let data = deps.storage.get(key.as_bytes())
        .ok_or_else(|| cosmwasm_std::StdError::generic_err("Agent not found"))?;
    let doc: DidDocument = bincode2::deserialize(&data)?;
    Ok(Binary::from(bincode2::serialize(&doc)?))
}

fn query_get_agent_by_address(
    deps: Deps,
    chain_id: String,
    address: String,
) -> StdResult<Binary> {
    let idx_key = format!("{}{}:{}", CHAIN_INDEX, chain_id, address);
    let did = deps.storage.get(idx_key.as_bytes())
        .ok_or_else(|| cosmwasm_std::StdError::generic_err("No agent found for address"))?;
    let did_str = String::from_utf8(did).unwrap();
    query_get_agent(deps, did_str)
}

fn query_get_chain_addresses(deps: Deps, agent_id: String) -> StdResult<Binary> {
    let did = format!("did:msg:agent:{}", agent_id);
    let key = format!("{}{}", DID_REGISTRY, did);
    let data = deps.storage.get(key.as_bytes())
        .ok_or_else(|| cosmwasm_std::StdError::generic_err("Agent not found"))?;
    let doc: DidDocument = bincode2::deserialize(&data)?;
    Ok(Binary::from(bincode2::serialize(&doc.chain_addresses)?))
}

fn query_list_agents(
    deps: Deps,
    start_after: Option<String>,
    limit: Option<u32>,
) -> StdResult<Binary> {
    let limit = limit.unwrap_or(30).min(100);
    let agents: Vec<DidDocument> = deps.storage
        .range(None, None, cosmwasm_std::Order::Ascending)
        .filter(|(k, _)| k.starts_with(DID_REGISTRY.as_bytes()))
        .skip(start_after.map(|_| 0).unwrap_or(0))
        .take(limit as usize)
        .map(|(_, v)| bincode2::deserialize::<DidDocument>(&v).unwrap())
        .collect();
    Ok(Binary::from(bincode2::serialize(&agents)?))
}

// ---------- IBC 数据包定义 ----------

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct IbcDidSyncPacket {
    pub did_document: DidDocument,
    pub source_chain: String,
    pub timestamp: cosmwasm_std::Timestamp,
}

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct IbcDidSyncAcknowledgement {
    pub status: SyncStatus,
    pub did: String,
    pub target_chain: String,
    pub error_msg: Option<String>,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
pub enum SyncStatus {
    Success,
    Failure,
    Pending,
}

// ---------- IBC 入口点 ----------

#[entry_point]
pub fn ibc_packet_receive(
    deps: DepsMut,
    env: Env,
    msg: IbcPacket,
) -> StdResult<IbcReceiveResponse> {
    let packet: IbcDidSyncPacket = bincode2::deserialize(&msg.data)?;
    let agent_id = packet.did_document.id.clone();
    let key = format!("{}{}", DID_REGISTRY, agent_id);
    if let Some(existing_data) = deps.storage.get(key.as_bytes()) {
        let local_doc: DidDocument = bincode2::deserialize(&existing_data)?;
        if packet.did_document.version > local_doc.version {
            deps.storage.set(key.as_bytes(), &msg.data);
        }
    } else {
        deps.storage.set(key.as_bytes(), &msg.data);
    }
    let ack = IbcDidSyncAcknowledgement {
        status: SyncStatus::Success,
        did: agent_id,
        target_chain: "msg-chain-1".to_string(),
        error_msg: None,
    };
    Ok(IbcReceiveResponse::new()
        .set_ack(bincode2::serialize(&ack)?)
        .add_attribute("action", "ibc_did_sync")
        .add_attribute("did", packet.did_document.id))
}

#[entry_point]
pub fn ibc_packet_ack(
    deps: DepsMut,
    _env: Env,
    ack: IbcPacketAckMsg,
) -> StdResult<Response> {
    let ack_data: IbcDidSyncAcknowledgement = bincode2::deserialize(&ack.acknowledgement.data)?;
    Ok(Response::new()
        .add_attribute("action", "ibc_did_sync_ack")
        .add_attribute("did", ack_data.did)
        .add_attribute("status", format!("{:?}", ack_data.status)))
}

#[entry_point]
pub fn ibc_packet_timeout(
    deps: DepsMut,
    _env: Env,
    _msg: IbcPacketTimeoutMsg,
) -> StdResult<Response> {
    Ok(Response::new()
        .add_attribute("action", "ibc_did_sync_timeout"))
}

// ---------- 辅助函数 ----------

fn get_channel_for_chain(chain_id: &str) -> String {
    match chain_id {
        "chain-a" => "0".to_string(),
        "chain-b" => "1".to_string(),
        "chain-c" => "2".to_string(),
        _ => panic!("unknown chain"),
    }
}

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct RegistryConfig {
    pub admin: String,
    pub supported_chains: Vec<String>,
}

// ---------- IBC 通道握手 ----------

#[entry_point]
pub fn ibc_channel_open(
    deps: DepsMut,
    _env: Env,
    msg: IbcChannelOpenMsg,
) -> StdResult<()> {
    let channel = msg.channel();
    if channel.endpoint.port_id != "did-registry" {
        return Err(cosmwasm_std::StdError::generic_err(
            format!("Expected port 'did-registry', got '{}'", channel.endpoint.port_id)
        ));
    }
    if channel.order != IbcOrder::Ordered {
        return Err(cosmwasm_std::StdError::generic_err("Must use ordered channels"));
    }
    Ok(())
}

#[entry_point]
pub fn ibc_channel_connect(
    deps: DepsMut,
    _env: Env,
    msg: IbcChannelConnectMsg,
) -> StdResult<Response> {
    let channel = msg.channel();
    let state = IbcSyncState {
        last_synced_sequence: 0,
        remote_chain_id: channel.counterparty_endpoint.port_id.clone(),
        remote_registry_address: String::new(),
    };
    let key = format!("{}{}", IBC_SYNC_STATE, channel.channel_id);
    deps.storage.set(key.as_bytes(), &bincode2::serialize(&state)?);
    Ok(Response::new()
        .add_attribute("action", "channel_connect")
        .add_attribute("channel", &channel.channel_id))
}

#[entry_point]
pub fn ibc_channel_close(
    deps: DepsMut,
    _env: Env,
    msg: IbcChannelCloseMsg,
) -> StdResult<Response> {
    let channel = msg.channel();
    let key = format!("{}{}", IBC_SYNC_STATE, channel.channel_id);
    deps.storage.remove(key.as_bytes());
    Ok(Response::new()
        .add_attribute("action", "channel_close")
        .add_attribute("channel", &channel.channel_id))
}

2.5 DID 解析流程

DID 解析是将 did:msg:agent:<id> 解析为完整的 DID 文档的过程。跨链场景下,解析器需要从多条链上收集信息。

// ============================================================
// 跨链 DID 解析器 (TypeScript)
// ============================================================

class CrossChainDidResolver {
    private rpcEndpoints: Map<string, string>;
    private registryAddresses: Map<string, string>;

    constructor() {
        this.rpcEndpoints = new Map();
        this.registryAddresses = new Map();
        this.rpcEndpoints.set('msg-chain-1', 'https://rpc.msg-chain-1.example.com');
        this.rpcEndpoints.set('chain-a', 'https://rpc.chain-a.example.com');
    }

    async resolve(did: string): Promise<CrossChainDID | null> {
        const baseDoc = await this.resolveFromChain('msg-chain-1', did);
        if (!baseDoc) return null;

        const allChains = baseDoc.chainAddresses.keys();
        for (const chainId of allChains) {
            if (chainId === 'msg-chain-1') continue;
            try {
                const remoteDoc = await this.resolveFromChain(chainId, did);
                if (remoteDoc) {
                    this.mergeDocuments(baseDoc, remoteDoc);
                }
            } catch (e) {
                console.warn(`Failed to resolve DID on ${chainId}:`, e);
            }
        }
        return baseDoc;
    }

    private async resolveFromChain(
        chainId: string,
        did: string
    ): Promise<CrossChainDID | null> {
        const rpc = this.rpcEndpoints.get(chainId);
        if (!rpc) throw new Error(`No RPC for chain ${chainId}`);

        const registryAddr = this.registryAddresses.get(chainId);
        const queryMsg = { get_agent: { did } };

        const response = await fetch(rpc, {
            method: 'POST',
            headers: { 'Content-Type': 'application/json' },
            body: JSON.stringify({
                jsonrpc: '2.0',
                id: 1,
                method: 'abci_query',
                params: {
                    path: `/cosmwasm/wasm/v1/contract/${registryAddr}/smart`,
                    data: Buffer.from(JSON.stringify(queryMsg)).toString('base64'),
                },
            }),
        });

        const result = await response.json();
        if (result.error) return null;
        return this.parseDidDocument(result.result.response);
    }

    private mergeDocuments(base: CrossChainDID, remote: CrossChainDID): void {
        for (const [chainId, addr] of remote.chainAddresses) {
            if (!base.chainAddresses.has(chainId)) {
                base.chainAddresses.set(chainId, addr);
            }
        }
        for (const vm of remote.verificationMethod) {
            if (!base.verificationMethod.find(v => v.id === vm.id)) {
                base.verificationMethod.push(vm);
            }
        }
        if (remote.updated > base.updated) {
            base.updated = remote.updated;
        }
    }

    private parseDidDocument(raw: any): CrossChainDID {
        const bytes = Buffer.from(raw, 'base64');
        return JSON.parse(bytes.toString());
    }
}

2.6 Agent 部署脚本

Agent 部署时,自动化脚本创建设立跨链身份所需的全部资源。

#!/bin/bash
# ============================================================
# 跨链 Agent 部署脚本
# ============================================================

set -euo pipefail

MSG_RPC="https://rpc.msg-chain-1.example.com"
MSG_REGISTRY="msg1registryaddress..."
AGENT_KEY_NAME="agent-$1"

echo "=== 跨链 Agent 部署 ==="

# 1. 生成 Agent 身份
echo "[1/4] 生成 DID..."
agent_id=$(openssl rand -hex 16)
did="did:msg:agent:$agent_id"
echo "  DID: $did"

# 2. 生成密钥对
echo "[2/4] 生成密钥..."
msgd keys add $AGENT_KEY_NAME --keyring-backend file
agent_addr=$(msgd keys show $AGENT_KEY_NAME -a --keyring-backend file)
echo "  Address: $agent_addr"

# 3. 在 MSG Chain 上注册
echo "[3/4] 注册到 MSG Chain..."
register_msg=$(cat <<EOF
{
    "register_agent": {
        "did_document": {
            "id": "$did",
            "controller": "$agent_addr",
            "chain_addresses": [
                {
                    "chain_id": "msg-chain-1",
                    "address": "$agent_addr",
                    "network": "mainnet"
                }
            ],
            "verification_method": [],
            "version": 1
        }
    }
}
EOF
)

msgd tx wasm execute $MSG_REGISTRY "$register_msg" \
    --from $AGENT_KEY_NAME \
    --node $MSG_RPC \
    --gas auto \
    --gas-prices 1000000000attoMSG \
    -y

# 4. 建立 IBC 通道并同步到其他链
echo "[4/4] 跨链同步..."
for target_chain in "chain-a" "chain-b"; do
    sync_msg=$(cat <<EOF
{
    "sync_did_to_chain": {
        "target_chain": "$target_chain",
        "agent_id": "$agent_id"
    }
}
EOF
)
    msgd tx wasm execute $MSG_REGISTRY "$sync_msg" \
        --from $AGENT_KEY_NAME \
        --node $MSG_RPC \
        --gas auto \
        --gas-prices 1000000000attoMSG \
        -y
done

echo "=== 部署完成 ==="
echo "DID: $did"
echo "MSG Address: $agent_addr"

3. 跨链Agent注册与发现

3.1 IBC 数据包设计

Agent 注册与发现需要跨链同步 Agent 元数据。

// ============================================================
// IBC Agent 注册数据包 (Rust)
// ============================================================

use serde::{Deserialize, Serialize};
use std::collections::HashMap;

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
pub enum AgentCapability {
    TradeExecution,
    DataQuery,
    MessageRelay,
    Automation,
    Custom(String),
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
pub enum AgentStatus {
    Active,
    Paused,
    Retired,
    Slashed,
}

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct IbcAgentRegistration {
    pub agent_id: String,
    pub source_chain: String,
    pub target_chain: String,
    pub capabilities: Vec<AgentCapability>,
    pub did_document: Vec<u8>,
    pub metadata: HashMap<String, String>,
    pub timestamp: u64,
    pub signature: Vec<u8>,
    pub version: u32,
}

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct IbcAgentHeartbeat {
    pub agent_id: String,
    pub source_chain: String,
    pub status: AgentStatus,
    pub last_block_height: u64,
    pub performance_metrics: HashMap<String, String>,
    pub timestamp: u64,
}

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct IbcAgentDiscoveryQuery {
    pub query_id: String,
    pub query_type: DiscoveryQueryType,
    pub source_chain: String,
    pub filters: HashMap<String, String>,
    pub max_results: u32,
    pub requester_did: String,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
pub enum DiscoveryQueryType {
    ByCapability,
    ByChain,
    ByAddress,
    All,
}

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct IbcAgentDiscoveryResponse {
    pub query_id: String,
    pub source_chain: String,
    pub agents: Vec<DiscoveredAgent>,
    pub total_count: u32,
    pub has_more: bool,
}

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct DiscoveredAgent {
    pub agent_id: String,
    pub did: String,
    pub chain_addresses: Vec<ChainAddressEntry>,
    pub capabilities: Vec<AgentCapability>,
    pub status: AgentStatus,
    pub last_seen: u64,
}

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct IbcRegistrationAcknowledgement {
    pub status: RegistrationStatus,
    pub agent_id: String,
    pub registered_chain: String,
    pub local_agent_id: String,
    pub error_msg: Option<String>,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
pub enum RegistrationStatus {
    Accepted,
    Rejected,
    Pending,
}

3.2 Agent 注册表合约实现

// ============================================================
// 跨链 Agent 注册表合约 (Rust - CosmWasm)
// ============================================================

use cosmwasm_std::{
    entry_point, to_binary, Binary, Deps, DepsMut, Env, MessageInfo, Response, StdResult,
    IbcMsg, IbcPacket, IbcReceiveResponse, IbcPacketAckMsg, IbcPacketTimeoutMsg,
    IbcChannelOpenMsg, IbcChannelConnectMsg, IbcChannelCloseMsg, IbcOrder, IbcTimeout,
    Timestamp, Storage,
};

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct AgentRegistration {
    pub agent_id: String,
    pub did: String,
    pub owner: String,
    pub capabilities: Vec<AgentCapability>,
    pub status: AgentStatus,
    pub chain: String,
    pub address: String,
    pub registered_at: u64,
    pub last_heartbeat: u64,
    pub metadata: HashMap<String, String>,
}

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct CrossChainChannel {
    pub channel_id: String,
    pub remote_chain_id: String,
    pub remote_port: String,
    pub last_sequence: u64,
    pub is_active: bool,
}

const AGENTS: &str = "agents";
const CHANNELS: &str = "channels";
const CAPABILITY_INDEX: &str = "cap_idx";
const TOTAL_AGENTS: &str = "total_agents";

#[derive(Serialize, Deserialize)]
pub struct InstantiateMsg {
    pub admin: String,
    pub chain_id: String,
}

#[entry_point]
pub fn instantiate(
    deps: DepsMut,
    _env: Env,
    info: MessageInfo,
    msg: InstantiateMsg,
) -> StdResult<Response> {
    let config = RegistryConfig {
        admin: info.sender.to_string(),
        chain_id: msg.chain_id,
        total_agents: 0,
        active_channels: 0,
    };
    deps.storage.set(b"config", &bincode2::serialize(&config)?);
    Ok(Response::new()
        .add_attribute("action", "instantiate")
        .add_attribute("chain_id", &msg.chain_id))
}

#[derive(Serialize, Deserialize)]
pub enum ExecuteMsg {
    Register(IbcAgentRegistration),
    Update {
        agent_id: String,
        capabilities: Option<Vec<AgentCapability>>,
        metadata: Option<HashMap<String, String>>,
    },
    Heartbeat(IbcAgentHeartbeat),
    Deregister { agent_id: String },
    CrossChainQuery {
        target_chain: String,
        query: IbcAgentDiscoveryQuery,
    },
    SyncToAllChannels { agent_id: String },
    PauseAgent { agent_id: String },
    ResumeAgent { agent_id: String },
}

#[entry_point]
pub fn execute(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    msg: ExecuteMsg,
) -> StdResult<Response> {
    match msg {
        ExecuteMsg::Register(reg) => execute_register(deps, env, info, reg),
        ExecuteMsg::Update { agent_id, capabilities, metadata } => {
            execute_update(deps, env, info, agent_id, capabilities, metadata)
        }
        ExecuteMsg::Heartbeat(hb) => execute_heartbeat(deps, env, info, hb),
        ExecuteMsg::Deregister { agent_id } => execute_deregister(deps, env, info, agent_id),
        ExecuteMsg::CrossChainQuery { target_chain, query } => {
            execute_cross_chain_query(deps, env, info, target_chain, query)
        }
        ExecuteMsg::SyncToAllChannels { agent_id } => {
            execute_sync_to_all(deps, env, info, agent_id)
        }
        ExecuteMsg::PauseAgent { agent_id } => execute_pause(deps, env, info, agent_id),
        ExecuteMsg::ResumeAgent { agent_id } => execute_resume(deps, env, info, agent_id),
    }
}

fn execute_register(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    reg: IbcAgentRegistration,
) -> StdResult<Response> {
    let store = deps.storage;
    let key = format!("{}{}", AGENTS, reg.agent_id);
    if store.get(key.as_bytes()).is_some() {
        return Err(cosmwasm_std::StdError::generic_err(
            format!("Agent {} already registered", reg.agent_id)
        ));
    }
    let registration = AgentRegistration {
        agent_id: reg.agent_id.clone(),
        did: String::from_utf8(reg.did_document.clone()).unwrap_or_default(),
        owner: info.sender.to_string(),
        capabilities: reg.capabilities.clone(),
        status: AgentStatus::Active,
        chain: reg.source_chain.clone(),
        address: info.sender.to_string(),
        registered_at: env.block.time.seconds(),
        last_heartbeat: env.block.time.seconds(),
        metadata: reg.metadata.clone(),
    };
    store.set(key.as_bytes(), &bincode2::serialize(&registration)?);
    for cap in &reg.capabilities {
        let cap_key = format!("{}{}:{}", CAPABILITY_INDEX, cap_to_str(cap), reg.agent_id);
        store.set(cap_key.as_bytes(), b"1");
    }
    let total_key = TOTAL_AGENTS;
    let total = store.get(total_key.as_bytes())
        .map(|d| u64::from_be_bytes(d.try_into().unwrap()))
        .unwrap_or(0) + 1;
    store.set(total_key.as_bytes(), &total.to_be_bytes());
    let broadcast = broadcast_registration_to_channels(store, &reg, env)?;
    Ok(Response::new()
        .add_messages(broadcast)
        .add_attribute("action", "register")
        .add_attribute("agent_id", &reg.agent_id)
        .add_attribute("capabilities", format!("{:?}", reg.capabilities)))
}

fn broadcast_registration_to_channels(
    store: &mut dyn Storage,
    reg: &IbcAgentRegistration,
    env: Env,
) -> StdResult<Vec<IbcMsg>> {
    let mut msgs = Vec::new();
    let channels = get_all_channels(store);
    for channel in &channels {
        if !channel.is_active { continue; }
        let packet_data = bincode2::serialize(reg)?;
        let ibc_msg = IbcMsg::SendPacket {
            channel_id: channel.channel_id.clone(),
            data: Binary::from(packet_data),
            timeout: IbcTimeout::with_timestamp(env.block.time.plus_seconds(600)),
        };
        msgs.push(ibc_msg);
    }
    Ok(msgs)
}

fn get_all_channels(store: &dyn Storage) -> Vec<CrossChainChannel> {
    let mut channels = Vec::new();
    let prefix = CHANNELS.as_bytes();
    for (key, value) in store.range(Some(prefix.to_vec()), None, cosmwasm_std::Order::Ascending) {
        if let Ok(ch) = bincode2::deserialize::<CrossChainChannel>(&value) {
            channels.push(ch);
        }
    }
    channels
}

fn execute_update(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    agent_id: String,
    capabilities: Option<Vec<AgentCapability>>,
    metadata: Option<HashMap<String, String>>,
) -> StdResult<Response> {
    let store = deps.storage;
    let key = format!("{}{}", AGENTS, agent_id);
    let existing = store.get(key.as_bytes())
        .ok_or_else(|| cosmwasm_std::StdError::generic_err("Agent not found"))?;
    let mut agent: AgentRegistration = bincode2::deserialize(&existing)?;
    if agent.owner != info.sender.to_string() {
        return Err(cosmwasm_std::StdError::generic_err("Unauthorized"));
    }
    if let Some(caps) = capabilities {
        for old_cap in &agent.capabilities {
            let cap_key = format!("{}{}:{}", CAPABILITY_INDEX, cap_to_str(old_cap), agent_id);
            store.remove(cap_key.as_bytes());
        }
        for new_cap in &caps {
            let cap_key = format!("{}{}:{}", CAPABILITY_INDEX, cap_to_str(new_cap), agent_id);
            store.set(cap_key.as_bytes(), b"1");
        }
        agent.capabilities = caps;
    }
    if let Some(meta) = metadata {
        agent.metadata = meta;
    }
    agent.registered_at = env.block.time.seconds();
    store.set(key.as_bytes(), &bincode2::serialize(&agent)?);
    Ok(Response::new()
        .add_attribute("action", "update")
        .add_attribute("agent_id", &agent_id))
}

fn execute_heartbeat(
    deps: DepsMut,
    env: Env,
    _info: MessageInfo,
    hb: IbcAgentHeartbeat,
) -> StdResult<Response> {
    let store = deps.storage;
    let key = format!("{}{}", AGENTS, hb.agent_id);
    let existing = store.get(key.as_bytes())
        .ok_or_else(|| cosmwasm_std::StdError::generic_err("Agent not found"))?;
    let mut agent: AgentRegistration = bincode2::deserialize(&existing)?;
    agent.last_heartbeat = env.block.time.seconds();
    agent.status = hb.status;
    store.set(key.as_bytes(), &bincode2::serialize(&agent)?);
    Ok(Response::new()
        .add_attribute("action", "heartbeat")
        .add_attribute("agent_id", &hb.agent_id)
        .add_attribute("status", format!("{:?}", hb.status)))
}

fn execute_deregister(
    deps: DepsMut,
    _env: Env,
    info: MessageInfo,
    agent_id: String,
) -> StdResult<Response> {
    let store = deps.storage;
    let key = format!("{}{}", AGENTS, agent_id);
    let existing = store.get(key.as_bytes())
        .ok_or_else(|| cosmwasm_std::StdError::generic_err("Agent not found"))?;
    let agent: AgentRegistration = bincode2::deserialize(&existing)?;
    if agent.owner != info.sender.to_string() {
        return Err(cosmwasm_std::StdError::generic_err("Unauthorized"));
    }
    for cap in &agent.capabilities {
        let cap_key = format!("{}{}:{}", CAPABILITY_INDEX, cap_to_str(cap), agent_id);
        store.remove(cap_key.as_bytes());
    }
    store.remove(key.as_bytes());
    let total_key = TOTAL_AGENTS;
    let total = store.get(total_key.as_bytes())
        .map(|d| u64::from_be_bytes(d.try_into().unwrap()))
        .unwrap_or(1) - 1;
    store.set(total_key.as_bytes(), &total.to_be_bytes());
    Ok(Response::new()
        .add_attribute("action", "deregister")
        .add_attribute("agent_id", &agent_id))
}

fn execute_cross_chain_query(
    deps: DepsMut,
    env: Env,
    _info: MessageInfo,
    target_chain: String,
    query: IbcAgentDiscoveryQuery,
) -> StdResult<Response> {
    let channel_id = find_channel_for_chain(deps.storage, &target_chain)
        .ok_or_else(|| cosmwasm_std::StdError::generic_err(
            format!("No channel for chain {}", target_chain)
        ))?;
    let packet_data = bincode2::serialize(&query)?;
    let ibc_msg = IbcMsg::SendPacket {
        channel_id,
        data: Binary::from(packet_data),
        timeout: IbcTimeout::with_timestamp(env.block.time.plus_seconds(300)),
    };
    Ok(Response::new()
        .add_message(ibc_msg)
        .add_attribute("action", "cross_chain_query")
        .add_attribute("target", &target_chain))
}

fn find_channel_for_chain(store: &dyn Storage, chain_id: &str) -> Option<String> {
    let channels = get_all_channels(store);
    channels.iter()
        .find(|c| c.remote_chain_id == chain_id)
        .map(|c| c.channel_id.clone())
}

fn cap_to_str(cap: &AgentCapability) -> String {
    match cap {
        AgentCapability::TradeExecution => "trade_execution".to_string(),
        AgentCapability::DataQuery => "data_query".to_string(),
        AgentCapability::MessageRelay => "message_relay".to_string(),
        AgentCapability::Automation => "automation".to_string(),
        AgentCapability::Custom(s) => format!("custom:{}", s),
    }
}

fn execute_sync_to_all(
    deps: DepsMut,
    env: Env,
    _info: MessageInfo,
    agent_id: String,
) -> StdResult<Response> {
    let store = deps.storage;
    let key = format!("{}{}", AGENTS, agent_id);
    let data = store.get(key.as_bytes())
        .ok_or_else(|| cosmwasm_std::StdError::generic_err("Agent not found"))?;
    let registration: AgentRegistration = bincode2::deserialize(&data)?;
    let channels = get_all_channels(store);
    let mut msgs = Vec::new();
    for channel in channels.iter().filter(|c| c.is_active) {
        let msg = IbcMsg::SendPacket {
            channel_id: channel.channel_id.clone(),
            data: Binary::from(data.clone()),
            timeout: IbcTimeout::with_timestamp(env.block.time.plus_seconds(600)),
        };
        msgs.push(msg);
    }
    Ok(Response::new()
        .add_messages(msgs)
        .add_attribute("action", "sync_to_all")
        .add_attribute("agent_id", &agent_id))
}

fn execute_pause(
    deps: DepsMut,
    _env: Env,
    info: MessageInfo,
    agent_id: String,
) -> StdResult<Response> {
    let store = deps.storage;
    let key = format!("{}{}", AGENTS, agent_id);
    let existing = store.get(key.as_bytes())
        .ok_or_else(|| cosmwasm_std::StdError::generic_err("Agent not found"))?;
    let mut agent: AgentRegistration = bincode2::deserialize(&existing)?;
    if agent.owner != info.sender.to_string() {
        return Err(cosmwasm_std::StdError::generic_err("Unauthorized"));
    }
    agent.status = AgentStatus::Paused;
    store.set(key.as_bytes(), &bincode2::serialize(&agent)?);
    Ok(Response::new()
        .add_attribute("action", "pause")
        .add_attribute("agent_id", &agent_id))
}

fn execute_resume(
    deps: DepsMut,
    _env: Env,
    info: MessageInfo,
    agent_id: String,
) -> StdResult<Response> {
    let store = deps.storage;
    let key = format!("{}{}", AGENTS, agent_id);
    let existing = store.get(key.as_bytes())
        .ok_or_else(|| cosmwasm_std::StdError::generic_err("Agent not found"))?;
    let mut agent: AgentRegistration = bincode2::deserialize(&existing)?;
    if agent.owner != info.sender.to_string() {
        return Err(cosmwasm_std::StdError::generic_err("Unauthorized"));
    }
    agent.status = AgentStatus::Active;
    store.set(key.as_bytes(), &bincode2::serialize(&agent)?);
    Ok(Response::new()
        .add_attribute("action", "resume")
        .add_attribute("agent_id", &agent_id))
}

// ---------- 查询 ----------

#[derive(Serialize, Deserialize)]
pub enum QueryMsg {
    GetAgent { agent_id: String },
    GetAgentsByCapability { capability: AgentCapability, limit: Option<u32> },
    ListAgents { start_after: Option<String>, limit: Option<u32> },
    QueryCrossChainAgents { chain_id: String, capability: Option<AgentCapability> },
    GetStats {},
    GetChannels {},
}

#[entry_point]
pub fn query(deps: Deps, _env: Env, msg: QueryMsg) -> StdResult<Binary> {
    match msg {
        QueryMsg::GetAgent { agent_id } => query_get_agent(deps, agent_id),
        QueryMsg::GetAgentsByCapability { capability, limit } => {
            query_by_capability(deps, capability, limit)
        }
        QueryMsg::ListAgents { start_after, limit } => {
            query_list_agents(deps, start_after, limit)
        }
        QueryMsg::QueryCrossChainAgents { chain_id, capability } => {
            query_cross_chain_agents(deps, chain_id, capability)
        }
        QueryMsg::GetStats {} => query_get_stats(deps),
        QueryMsg::GetChannels {} => query_get_channels(deps),
    }
}

fn query_get_agent(deps: Deps, agent_id: String) -> StdResult<Binary> {
    let key = format!("{}{}", AGENTS, agent_id);
    let data = deps.storage.get(key.as_bytes())
        .ok_or_else(|| cosmwasm_std::StdError::generic_err("Agent not found"))?;
    let agent: AgentRegistration = bincode2::deserialize(&data)?;
    to_binary(&agent)
}

fn query_by_capability(
    deps: Deps,
    capability: AgentCapability,
    limit: Option<u32>,
) -> StdResult<Binary> {
    let limit = limit.unwrap_or(20).min(100);
    let cap_str = cap_to_str(&capability);
    let prefix = format!("{}{}:", CAPABILITY_INDEX, cap_str);
    let mut agents = Vec::new();
    for (key, _) in deps.storage.range(
        Some(prefix.as_bytes().to_vec()),
        None,
        cosmwasm_std::Order::Ascending,
    ) {
        if agents.len() >= limit as usize { break; }
        let key_str = String::from_utf8(key).unwrap();
        let agent_id = key_str.split(':').last().unwrap_or("");
        let agent_key = format!("{}{}", AGENTS, agent_id);
        if let Some(data) = deps.storage.get(agent_key.as_bytes()) {
            if let Ok(agent) = bincode2::deserialize::<AgentRegistration>(&data) {
                agents.push(agent);
            }
        }
    }
    to_binary(&agents)
}

fn query_list_agents(
    deps: Deps,
    start_after: Option<String>,
    limit: Option<u32>,
) -> StdResult<Binary> {
    let limit = limit.unwrap_or(30).min(100);
    let prefix = AGENTS.as_bytes();
    let mut agents = Vec::new();
    let start = start_after.map(|s| format!("{}{}", AGENTS, s).into_bytes());
    for (_, value) in deps.storage.range(start, None, cosmwasm_std::Order::Ascending) {
        if agents.len() >= limit as usize { break; }
        if let Ok(agent) = bincode2::deserialize::<AgentRegistration>(&value) {
            agents.push(agent);
        }
    }
    to_binary(&agents)
}

fn query_cross_chain_agents(
    deps: Deps,
    chain_id: String,
    capability: Option<AgentCapability>,
) -> StdResult<Binary> {
    let mut agents: Vec<AgentRegistration> = Vec::new();
    let prefix = AGENTS.as_bytes();
    for (_, value) in deps.storage.range(Some(prefix.to_vec()), None, cosmwasm_std::Order::Ascending) {
        if let Ok(agent) = bincode2::deserialize::<AgentRegistration>(&value) {
            if agent.chain != chain_id { continue; }
            if let Some(ref cap) = capability {
                if !agent.capabilities.contains(cap) { continue; }
            }
            agents.push(agent);
        }
    }
    to_binary(&agents)
}

fn query_get_stats(deps: Deps) -> StdResult<Binary> {
    let total_key = TOTAL_AGENTS;
    let total = deps.storage.get(total_key.as_bytes())
        .map(|d| u64::from_be_bytes(d.try_into().unwrap()))
        .unwrap_or(0);
    let channels = get_all_channels(deps.storage);
    let active_channels = channels.iter().filter(|c| c.is_active).count();
    to_binary(&serde_json::json!({
        "total_agents": total,
        "active_channels": active_channels,
        "total_channels": channels.len(),
    }))
}

fn query_get_channels(deps: Deps) -> StdResult<Binary> {
    let channels = get_all_channels(deps.storage);
    to_binary(&channels)
}

// ---------- IBC 入口点 ----------

#[entry_point]
pub fn ibc_packet_receive(
    deps: DepsMut,
    env: Env,
    msg: IbcPacket,
) -> StdResult<IbcReceiveResponse> {
    let store = deps.storage;
    let source_channel = msg.dest.channel_id.clone();
    if let Ok(reg) = bincode2::deserialize::<IbcAgentRegistration>(&msg.data) {
        return handle_remote_registration(store, env, reg, &source_channel);
    }
    if let Ok(query) = bincode2::deserialize::<IbcAgentDiscoveryQuery>(&msg.data) {
        return handle_discovery_query(store, env, query, &source_channel);
    }
    if let Ok(hb) = bincode2::deserialize::<IbcAgentHeartbeat>(&msg.data) {
        return handle_remote_heartbeat(store, env, hb);
    }
    Err(cosmwasm_std::StdError::generic_err("Unknown packet type"))
}

fn handle_remote_registration(
    store: &mut dyn Storage,
    env: Env,
    reg: IbcAgentRegistration,
    source_channel: &str,
) -> StdResult<IbcReceiveResponse> {
    let agent_id = format!("remote:{}:{}", reg.source_chain, reg.agent_id);
    let key = format!("{}{}", AGENTS, agent_id);
    let is_update = store.get(key.as_bytes()).is_some();
    let registration = AgentRegistration {
        agent_id: agent_id.clone(),
        did: String::from_utf8(reg.did_document.clone()).unwrap_or_default(),
        owner: format!("ibc:{}:{}", source_channel, reg.agent_id),
        capabilities: reg.capabilities.clone(),
        status: AgentStatus::Active,
        chain: reg.source_chain.clone(),
        address: format!("ibc://{}/{}", reg.source_chain, reg.agent_id),
        registered_at: if is_update {
            bincode2::deserialize::<AgentRegistration>(&store.get(key.as_bytes()).unwrap())
                .map(|a| a.registered_at)
                .unwrap_or(env.block.time.seconds())
        } else {
            env.block.time.seconds()
        },
        last_heartbeat: env.block.time.seconds(),
        metadata: reg.metadata.clone(),
    };
    store.set(key.as_bytes(), &bincode2::serialize(&registration)?);
    for cap in &reg.capabilities {
        let cap_key = format!("{}{}:{}", CAPABILITY_INDEX, cap_to_str(cap), agent_id);
        store.set(cap_key.as_bytes(), b"1");
    }
    let ack = IbcRegistrationAcknowledgement {
        status: RegistrationStatus::Accepted,
        agent_id: reg.agent_id,
        registered_chain: "msg-chain-1".to_string(),
        local_agent_id: agent_id,
        error_msg: None,
    };
    Ok(IbcReceiveResponse::new()
        .set_ack(bincode2::serialize(&ack)?)
        .add_attribute("action", "remote_register")
        .add_attribute("agent_id", &agent_id)
        .add_attribute("is_update", is_update.to_string()))
}

fn handle_discovery_query(
    store: &dyn Storage,
    _env: Env,
    query: IbcAgentDiscoveryQuery,
    _source_channel: &str,
) -> StdResult<IbcReceiveResponse> {
    let mut agents = Vec::new();
    let prefix = AGENTS.as_bytes();
    for (_, value) in store.range(Some(prefix.to_vec()), None, cosmwasm_std::Order::Ascending) {
        if agents.len() >= query.max_results as usize { break; }
        if let Ok(reg) = bincode2::deserialize::<AgentRegistration>(&value) {
            if reg.status != AgentStatus::Active { continue; }
            let matched = match &query.query_type {
                DiscoveryQueryType::All => true,
                DiscoveryQueryType::ByCapability => {
                    if let Some(filter_cap) = query.filters.get("capability") {
                        reg.capabilities.iter().any(|c| cap_to_str(c) == filter_cap)
                    } else { false }
                }
                DiscoveryQueryType::ByChain => {
                    if let Some(filter_chain) = query.filters.get("chain") {
                        reg.chain == *filter_chain
                    } else { false }
                }
                DiscoveryQueryType::ByAddress => {
                    if let Some(filter_addr) = query.filters.get("address") {
                        reg.address == *filter_addr
                    } else { false }
                }
            };
            if matched {
                agents.push(DiscoveredAgent {
                    agent_id: reg.agent_id,
                    did: reg.did,
                    chain_addresses: vec![ChainAddressEntry {
                        chain_id: reg.chain.clone(),
                        address: reg.address.clone(),
                        network: "mainnet".to_string(),
                        proof: None,
                    }],
                    capabilities: reg.capabilities,
                    status: reg.status,
                    last_seen: reg.last_heartbeat,
                });
            }
        }
    }
    let response = IbcAgentDiscoveryResponse {
        query_id: query.query_id,
        source_chain: "msg-chain-1".to_string(),
        agents,
        total_count: agents.len() as u32,
        has_more: false,
    };
    Ok(IbcReceiveResponse::new()
        .set_ack(bincode2::serialize(&response)?)
        .add_attribute("action", "discovery_response")
        .add_attribute("query_id", &query.query_id))
}

fn handle_remote_heartbeat(
    store: &mut dyn Storage,
    env: Env,
    hb: IbcAgentHeartbeat,
) -> StdResult<IbcReceiveResponse> {
    let agent_id = format!("remote:{}:{}", hb.source_chain, hb.agent_id);
    let key = format!("{}{}", AGENTS, agent_id);
    if let Some(existing) = store.get(key.as_bytes()) {
        if let Ok(mut agent) = bincode2::deserialize::<AgentRegistration>(&existing) {
            agent.last_heartbeat = env.block.time.seconds();
            agent.status = hb.status;
            store.set(key.as_bytes(), &bincode2::serialize(&agent)?);
        }
    }
    Ok(IbcReceiveResponse::new()
        .set_ack(Binary::from(b"ok"))
        .add_attribute("action", "heartbeat_ack")
        .add_attribute("agent_id", &hb.agent_id))
}

#[entry_point]
pub fn ibc_packet_ack(
    deps: DepsMut,
    _env: Env,
    ack: IbcPacketAckMsg,
) -> StdResult<Response> {
    if let Ok(reg_ack) = bincode2::deserialize::<IbcRegistrationAcknowledgement>(&ack.acknowledgement.data) {
        return Ok(Response::new()
            .add_attribute("action", "registration_ack")
            .add_attribute("agent_id", &reg_ack.agent_id)
            .add_attribute("status", format!("{:?}", reg_ack.status)));
    }
    if let Ok(discovery_resp) = bincode2::deserialize::<IbcAgentDiscoveryResponse>(&ack.acknowledgement.data) {
        let key = format!("discovery:{}", discovery_resp.query_id);
        deps.storage.set(key.as_bytes(), &ack.acknowledgement.data);
        return Ok(Response::new()
            .add_attribute("action", "discovery_result")
            .add_attribute("query_id", &discovery_resp.query_id)
            .add_attribute("count", discovery_resp.total_count.to_string()));
    }
    Ok(Response::new()
        .add_attribute("action", "packet_ack"))
}

#[entry_point]
pub fn ibc_packet_timeout(
    deps: DepsMut,
    _env: Env,
    _msg: IbcPacketTimeoutMsg,
) -> StdResult<Response> {
    Ok(Response::new()
        .add_attribute("action", "packet_timeout"))
}

// ---------- IBC 通道 ----------

#[entry_point]
pub fn ibc_channel_open(
    deps: DepsMut,
    _env: Env,
    msg: IbcChannelOpenMsg,
) -> StdResult<()> {
    let channel = msg.channel();
    if channel.endpoint.port_id != "agent-registry" {
        return Err(cosmwasm_std::StdError::generic_err(
            format!("Expected port 'agent-registry', got '{}'", channel.endpoint.port_id)
        ));
    }
    Ok(())
}

#[entry_point]
pub fn ibc_channel_connect(
    deps: DepsMut,
    _env: Env,
    msg: IbcChannelConnectMsg,
) -> StdResult<Response> {
    let channel = msg.channel();
    let cross_chain_channel = CrossChainChannel {
        channel_id: channel.channel_id.clone(),
        remote_chain_id: channel.counterparty_endpoint.port_id.clone(),
        remote_port: channel.counterparty_endpoint.port_id.clone(),
        last_sequence: 0,
        is_active: true,
    };
    let key = format!("{}{}", CHANNELS, channel.channel_id);
    deps.storage.set(key.as_bytes(), &bincode2::serialize(&cross_chain_channel)?);
    Ok(Response::new()
        .add_attribute("action", "channel_connect")
        .add_attribute("channel_id", &channel.channel_id))
}

#[entry_point]
pub fn ibc_channel_close(
    deps: DepsMut,
    _env: Env,
    msg: IbcChannelCloseMsg,
) -> StdResult<Response> {
    let channel = msg.channel();
    let key = format!("{}{}", CHANNELS, channel.channel_id);
    if let Some(mut ch_data) = deps.storage.get(key.as_bytes())
        .map(|d| bincode2::deserialize::<CrossChainChannel>(&d).unwrap())
    {
        ch_data.is_active = false;
        deps.storage.set(key.as_bytes(), &bincode2::serialize(&ch_data)?);
    }
    Ok(Response::new()
        .add_attribute("action", "channel_close")
        .add_attribute("channel_id", &channel.channel_id))
}

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct RegistryConfig {
    pub admin: String,
    pub chain_id: String,
    pub total_agents: u64,
    pub active_channels: u64,
}

3.3 Relayer 配置

Relayer 负责在 MSG Chain 和其他链之间转发 IBC 数据包。

#!/bin/bash
# ============================================================
# Agent 注册数据 Relayer 配置
# ============================================================

# 安装 Go relayer
go install github.com/cosmos/relayer/v2/cmd/rly@latest

# 初始化配置
rly config init

# 添加 MSG Chain
rly chains add -f <<EOF
{
  "type": "cosmos",
  "value": {
    "key": "relayer",
    "chain-id": "msg-chain-1",
    "rpc-addr": "https://rpc.msg-chain-1.example.com:26657",
    "account-prefix": "msg",
    "keyring-backend": "file",
    "gas-adjustment": 1.3,
    "gas-prices": "1000000000attoMSG",
    "debug": false,
    "timeout": "20s",
    "output-format": "json",
    "sign-mode": "direct"
  }
}
EOF

# 添加对端链
rly chains add -f <<EOF
{
  "type": "cosmos",
  "value": {
    "key": "relayer",
    "chain-id": "chain-a",
    "rpc-addr": "https://rpc.chain-a.example.com:26657",
    "account-prefix": "cosmos",
    "keyring-backend": "file",
    "gas-adjustment": 1.3,
    "gas-prices": "0.025uatom",
    "debug": false,
    "timeout": "20s",
    "output-format": "json",
    "sign-mode": "direct"
  }
}
EOF

# 创建 IBC 连接
rly paths new msg-chain-1 chain-a agent-registry-path \
    --src-port agent-registry \
    --dst-port agent-registry \
    --order ordered \
    --version "agent-registry-v1"

# 创建 IBC 客户端和连接
rly transact link agent-registry-path \
    --src-port agent-registry \
    --dst-port agent-registry \
    --order ordered \
    --version "agent-registry-v1"

# 启动 Relayer(持续转发 Agent 数据)
rly start agent-registry-path

# ============================================================
# Docker Compose Relayer 部署
# ============================================================

cat > docker-compose.relayer.yml <<'DOCKEREOF'
version: '3.8'

services:
  agent-relayer-msg-chain-a:
    image: cosmos/relayer:v2.5.2
    container_name: agent-relayer-msg-chain-a
    volumes:
      - ./relayer-config:/home/relayer/.relayer
      - ./relayer-keys:/home/relayer/keys
    command:
      - rly
      - start
      - agent-registry-path
    restart: unless-stopped
    environment:
      - RLY_LOG_LEVEL=info
    networks:
      - cross-chain-net

  agent-relayer-msg-chain-b:
    image: cosmos/relayer:v2.5.2
    container_name: agent-relayer-msg-chain-b
    volumes:
      - ./relayer-config:/home/relayer/.relayer
      - ./relayer-keys:/home/relayer/keys
    command:
      - rly
      - start
      - agent-registry-path-b
    restart: unless-stopped
    environment:
      - RLY_LOG_LEVEL=info
    networks:
      - cross-chain-net

networks:
  cross-chain-net:
    driver: bridge
DOCKEREOF

echo "Relayer 配置完成!"

3.4 跨链 Agent 发现客户端

# ============================================================
# 跨链 Agent 发现客户端 (Python)
# ============================================================

import hashlib
import json
import time
import asyncio
from dataclasses import dataclass, field
from typing import Optional

import httpx

@dataclass
class DiscoveredAgent:
    agent_id: str
    did: str
    chain: str
    address: str
    capabilities: list[str]
    status: str
    last_seen: int

@dataclass
class AgentDiscoveryQuery:
    query_id: str
    query_type: str
    source_chain: str
    filters: dict[str, str] = field(default_factory=dict)
    max_results: int = 50

class CrossChainAgentDiscovery:
    """跨链 Agent 发现客户端"""

    def __init__(
        self,
        msg_chain_rpc: str,
        registry_address: str,
        chain_id: str = "msg-chain-1",
    ):
        self.client = httpx.AsyncClient()
        self.registry_address = registry_address
        self.chain_id = chain_id
        self.cache: dict[str, tuple[list[DiscoveredAgent], float]] = {}
        self.cache_ttl = 60

    async def discover_local_agents(
        self,
        capability: Optional[str] = None,
        limit: int = 50,
    ) -> list[DiscoveredAgent]:
        query_msg = {"list_agents": {"limit": limit}}
        result = await self._query_contract(query_msg)
        agents = []
        for raw in result:
            agent = DiscoveredAgent(
                agent_id=raw["agent_id"],
                did=raw.get("did", ""),
                chain=raw.get("chain", self.chain_id),
                address=raw.get("address", ""),
                capabilities=raw.get("capabilities", []),
                status=raw.get("status", "Active"),
                last_seen=raw.get("last_heartbeat", 0),
            )
            if capability and capability not in agent.capabilities:
                continue
            agents.append(agent)
        return agents

    async def discover_cross_chain_agents(
        self,
        target_chain: str,
        capability: Optional[str] = None,
        timeout_seconds: int = 30,
    ) -> list[DiscoveredAgent]:
        cache_key = f"{target_chain}:{capability or 'all'}"
        if cache_key in self.cache:
            cached_agents, cached_time = self.cache[cache_key]
            if time.time() - cached_time < self.cache_ttl:
                return cached_agents

        query_id = hashlib.sha256(
            f"{target_chain}:{time.time()}:{capability}".encode()
        ).hexdigest()[:16]

        query = AgentDiscoveryQuery(
            query_id=query_id,
            query_type="ByCapability" if capability else "All",
            source_chain=self.chain_id,
            filters={"capability": capability} if capability else {},
            max_results=50,
        )

        execute_msg = {
            "cross_chain_query": {
                "target_chain": target_chain,
                "query": {
                    "query_id": query.query_id,
                    "query_type": query.query_type,
                    "source_chain": query.source_chain,
                    "filters": query.filters,
                    "max_results": query.max_results,
                },
            }
        }

        agents = await self._poll_discovery_result(query_id, timeout_seconds)
        self.cache[cache_key] = (agents, time.time())
        return agents

    async def _poll_discovery_result(
        self,
        query_id: str,
        timeout_seconds: int,
    ) -> list[DiscoveredAgent]:
        start = time.time()
        while time.time() - start < timeout_seconds:
            try:
                result = await self._query_contract(
                    {"get_discovery_result": {"query_id": query_id}}
                )
                if result and "agents" in result:
                    return [
                        DiscoveredAgent(
                            agent_id=a["agent_id"],
                            did=a.get("did", ""),
                            chain=a.get("chain", "unknown"),
                            address=a.get("address", ""),
                            capabilities=a.get("capabilities", []),
                            status=a.get("status", "Active"),
                            last_seen=a.get("last_seen", 0),
                        )
                        for a in result["agents"]
                    ]
            except Exception:
                pass
            await asyncio.sleep(2)
        return []

    async def find_agent_by_capability(
        self,
        capability: str,
        chains: list[str],
    ) -> dict[str, list[DiscoveredAgent]]:
        results: dict[str, list[DiscoveredAgent]] = {}
        results[self.chain_id] = await self.discover_local_agents(capability=capability)

        async def query_chain(chain: str):
            if chain == self.chain_id:
                return
            agents = await self.discover_cross_chain_agents(
                target_chain=chain, capability=capability,
            )
            if agents:
                results[chain] = agents

        tasks = [query_chain(chain) for chain in chains]
        await asyncio.gather(*tasks)
        return results

    async def _query_contract(self, msg: dict) -> dict:
        # 实际实现调用 LCD/REST 查询
        return {}

4. 跨链A2A通信

4.1 A2A 协议概述

Agent-to-Agent (A2A) 通信是 AI Agent 之间交换消息的标准协议。跨链 A2A 通过 IBC 数据包封装 A2A 消息,实现跨链消息路由。

# ============================================================
# 跨链 A2A 通信协议 (Python)
# ============================================================

import json
import uuid
from dataclasses import dataclass, field
from datetime import datetime
from enum import Enum
from typing import Any, Optional


class A2AMessageType(Enum):
    REQUEST = "a2a/request"
    RESPONSE = "a2a/response"
    NOTIFICATION = "a2a/notification"
    ERROR = "a2a/error"
    PING = "a2a/ping"
    PONG = "a2a/pong"


class A2ADeliveryStatus(Enum):
    PENDING = "pending"
    DELIVERED = "delivered"
    PROCESSED = "processed"
    FAILED = "failed"
    TIMEOUT = "timeout"


@dataclass
class A2AMessage:
    message_id: str = field(default_factory=lambda: str(uuid.uuid4()))
    message_type: A2AMessageType = A2AMessageType.REQUEST
    sender_did: str = ""
    sender_chain: str = ""
    target_did: str = ""
    target_chain: str = ""
    payload: dict[str, Any] = field(default_factory=dict)
    correlation_id: Optional[str] = None
    ttl_seconds: int = 300
    timestamp: str = field(default_factory=lambda: datetime.utcnow().isoformat())
    signature: Optional[str] = None
    schema_version: str = "1.0"

    def encode(self) -> bytes:
        return json.dumps({
            "message_id": self.message_id,
            "message_type": self.message_type.value,
            "sender_did": self.sender_did,
            "sender_chain": self.sender_chain,
            "target_did": self.target_did,
            "target_chain": self.target_chain,
            "payload": self.payload,
            "correlation_id": self.correlation_id,
            "ttl_seconds": self.ttl_seconds,
            "timestamp": self.timestamp,
            "signature": self.signature,
            "schema_version": self.schema_version,
        }).encode("utf-8")

    @classmethod
    def decode(cls, data: bytes) -> "A2AMessage":
        raw = json.loads(data.decode("utf-8"))
        return cls(
            message_id=raw["message_id"],
            message_type=A2AMessageType(raw["message_type"]),
            sender_did=raw["sender_did"],
            sender_chain=raw["sender_chain"],
            target_did=raw["target_did"],
            target_chain=raw["target_chain"],
            payload=raw["payload"],
            correlation_id=raw.get("correlation_id"),
            ttl_seconds=raw.get("ttl_seconds", 300),
            timestamp=raw.get("timestamp", datetime.utcnow().isoformat()),
            signature=raw.get("signature"),
            schema_version=raw.get("schema_version", "1.0"),
        )

    def sign(self, private_key: bytes) -> None:
        import hashlib
        from ecdsa import SigningKey, SECP256k1
        message_bytes = self.encode()
        sk = SigningKey.from_string(private_key, curve=SECP256k1)
        signature = sk.sign(message_bytes)
        self.signature = signature.hex()

    def verify(self, public_key: bytes) -> bool:
        if not self.signature:
            return False
        from ecdsa import VerifyingKey, SECP256k1
        message_bytes = self.encode()
        vk = VerifyingKey.from_string(public_key, curve=SECP256k1)
        try:
            return vk.verify(bytes.fromhex(self.signature), message_bytes)
        except Exception:
            return False

4.2 跨链 A2A 通信管理器

# ============================================================
# 跨链 A2A 通信管理器 (Python)
# ============================================================

class CrossChainA2A:
    """通过 IBC 在 MSG Chain 和其他链之间路由 A2A 消息"""

    def __init__(
        self,
        agent_did: str,
        agent_private_key: bytes,
        msg_chain_rpc: str,
        registry_address: str,
        a2a_contract_address: str,
        chain_id: str = "msg-chain-1",
    ):
        self.agent_did = agent_did
        self.agent_private_key = agent_private_key
        self.chain_id = chain_id
        self.registry_address = registry_address
        self.a2a_contract = a2a_contract_address
        self.channels: dict[str, str] = {
            "chain-a": "0",
            "chain-b": "1",
        }
        self.pending_requests: dict[str, asyncio.Future] = {}
        self.message_handlers: dict[str, callable] = {}

    async def send_message(
        self,
        target_chain: str,
        target_agent_did: str,
        message_payload: dict[str, Any],
        message_type: A2AMessageType = A2AMessageType.REQUEST,
        ttl_seconds: int = 300,
        await_response: bool = False,
        response_timeout: int = 60,
    ) -> Optional[dict[str, Any]]:
        channel_id = self.channels.get(target_chain)
        if not channel_id:
            raise ValueError(f"No channel configured for chain {target_chain}")

        message = A2AMessage(
            message_type=message_type,
            sender_did=self.agent_did,
            sender_chain=self.chain_id,
            target_did=target_agent_did,
            target_chain=target_chain,
            payload=message_payload,
            ttl_seconds=ttl_seconds,
        )
        message.sign(self.agent_private_key)

        execute_msg = {
            "send_a2a_message": {
                "target_chain": target_chain,
                "packet_data": message.encode().hex(),
                "timeout_seconds": ttl_seconds,
            }
        }

        if await_response:
            loop = asyncio.get_event_loop()
            future = loop.create_future()
            self.pending_requests[message.message_id] = future
            try:
                response = await asyncio.wait_for(future, timeout=response_timeout)
                return response
            except asyncio.TimeoutError:
                self.pending_requests.pop(message.message_id, None)
                raise TimeoutError(
                    f"No response from {target_agent_did} on {target_chain}"
                )
        return None

    async def handle_incoming_message(
        self,
        raw_message: bytes,
    ) -> Optional[A2AMessage]:
        message = A2AMessage.decode(raw_message)
        msg_time = datetime.fromisoformat(message.timestamp)
        if (datetime.utcnow() - msg_time).seconds > message.ttl_seconds:
            print(f"Message {message.message_id} expired")
            return None
        if message.target_did != self.agent_did:
            print(f"Message not for us: {message.target_did}")
            return None

        if message.message_type == A2AMessageType.REQUEST:
            handler_key = message.payload.get("action", "default")
            handler = self.message_handlers.get(handler_key)
            if handler:
                response_payload = await handler(message)
                await self.send_message(
                    target_chain=message.sender_chain,
                    target_agent_did=message.sender_did,
                    message_payload=response_payload,
                    message_type=A2AMessageType.RESPONSE,
                    ttl_seconds=60,
                )
        elif message.message_type == A2AMessageType.RESPONSE:
            corr_id = message.correlation_id
            if corr_id and corr_id in self.pending_requests:
                future = self.pending_requests.pop(corr_id)
                future.set_result(message.payload)
        return message

    def register_handler(self, action: str, handler: callable):
        self.message_handlers[action] = handler

4.3 A2A IBC 合约实现

// ============================================================
// 跨链 A2A IBC 合约 (Rust - CosmWasm)
// ============================================================

use cosmwasm_std::{
    entry_point, to_binary, Binary, Deps, DepsMut, Env, MessageInfo, Response,
    StdResult, StdError, IbcMsg, IbcPacket, IbcReceiveResponse, IbcPacketAckMsg,
    IbcPacketTimeoutMsg, IbcChannelOpenMsg, IbcChannelConnectMsg, IbcChannelCloseMsg,
    IbcOrder, IbcTimeout, Timestamp, Storage,
};
use serde::{Deserialize, Serialize};

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct A2aIbcPacket {
    pub message_id: String,
    pub message_type: String,
    pub sender_did: String,
    pub sender_chain: String,
    pub target_did: String,
    pub target_chain: String,
    pub payload: Binary,
    pub correlation_id: Option<String>,
    pub ttl_seconds: u64,
    pub timestamp: String,
    pub signature: Option<Binary>,
    pub schema_version: String,
}

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct A2AAcknowledgement {
    pub message_id: String,
    pub status: A2ADeliveryStatus,
    pub error: Option<String>,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
pub enum A2ADeliveryStatus {
    Accepted,
    Delivered,
    Rejected,
    Expired,
}

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct A2AChannelConfig {
    pub channel_id: String,
    pub remote_chain_id: String,
    pub is_active: bool,
    pub total_messages_sent: u64,
    pub total_messages_received: u64,
}

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct PendingMessage {
    pub message_id: String,
    pub sender_did: String,
    pub target_did: String,
    pub target_chain: String,
    pub payload: Binary,
    pub expires_at: u64,
    pub status: PendingStatus,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
pub enum PendingStatus {
    AwaitingDelivery,
    AwaitingResponse,
    Completed,
    Failed,
}

const A2A_CHANNELS: &str = "a2a_channels";
const PENDING_MSGS: &str = "pending_msgs";
const DELIVERED_MSGS: &str = "delivered_msgs";
const A2A_CONFIG: &str = "a2a_config";

#[derive(Serialize, Deserialize)]
pub struct InstantiateMsg {
    pub admin: String,
    pub chain_id: String,
}

#[entry_point]
pub fn instantiate(
    deps: DepsMut,
    _env: Env,
    info: MessageInfo,
    msg: InstantiateMsg,
) -> StdResult<Response> {
    let config = A2aGlobalConfig {
        admin: info.sender.to_string(),
        chain_id: msg.chain_id,
        total_messages: 0,
        active_channels: 0,
    };
    deps.storage.set(A2A_CONFIG.as_bytes(), &bincode2::serialize(&config)?);
    Ok(Response::new()
        .add_attribute("action", "instantiate")
        .add_attribute("chain_id", &msg.chain_id))
}

#[derive(Serialize, Deserialize)]
pub enum ExecuteMsg {
    SendA2AMessage {
        target_chain: String,
        packet_data: Binary,
        timeout_seconds: u64,
    },
    ReplyToMessage {
        original_message_id: String,
        response_payload: Binary,
    },
    RegisterHandler {
        action: String,
        handler_address: String,
    },
}

#[entry_point]
pub fn execute(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    msg: ExecuteMsg,
) -> StdResult<Response> {
    match msg {
        ExecuteMsg::SendA2AMessage { target_chain, packet_data, timeout_seconds } => {
            execute_send_message(deps, env, info, target_chain, packet_data, timeout_seconds)
        }
        ExecuteMsg::ReplyToMessage { original_message_id, response_payload } => {
            execute_reply(deps, env, info, original_message_id, response_payload)
        }
        ExecuteMsg::RegisterHandler { action, handler_address } => {
            execute_register_handler(deps, env, info, action, handler_address)
        }
    }
}

fn execute_send_message(
    deps: DepsMut,
    env: Env,
    _info: MessageInfo,
    target_chain: String,
    packet_data: Binary,
    timeout_seconds: u64,
) -> StdResult<Response> {
    let store = deps.storage;
    let a2a_msg: A2aIbcPacket = bincode2::deserialize(&packet_data.0)?;

    if target_chain == get_local_chain_id(store) {
        return handle_local_delivery(store, env, a2a_msg);
    }

    let channel_id = find_channel_for_target(store, &target_chain)
        .ok_or_else(|| StdError::generic_err(
            format!("No IBC channel for chain {}", target_chain)
        ))?;

    let ibc_packet = bincode2::serialize(&a2a_msg)?;
    let send_msg = IbcMsg::SendPacket {
        channel_id: channel_id.clone(),
        data: Binary::from(ibc_packet),
        timeout: IbcTimeout::with_timestamp(
            env.block.time.plus_seconds(timeout_seconds)
        ),
    };

    let pending = PendingMessage {
        message_id: a2a_msg.message_id.clone(),
        sender_did: a2a_msg.sender_did.clone(),
        target_did: a2a_msg.target_did.clone(),
        target_chain: target_chain.clone(),
        payload: Binary::from(bincode2::serialize(&a2a_msg)?),
        expires_at: env.block.time.seconds() + timeout_seconds,
        status: PendingStatus::AwaitingDelivery,
    };
    let pending_key = format!("{}{}", PENDING_MSGS, a2a_msg.message_id);
    store.set(pending_key.as_bytes(), &bincode2::serialize(&pending)?);
    update_message_count(store);

    Ok(Response::new()
        .add_message(send_msg)
        .add_attribute("action", "send_a2a")
        .add_attribute("message_id", &a2a_msg.message_id)
        .add_attribute("target_chain", &target_chain)
        .add_attribute("channel_id", &channel_id))
}

fn handle_local_delivery(
    store: &mut dyn Storage,
    env: Env,
    a2a_msg: A2aIbcPacket,
) -> StdResult<Response> {
    let msg_time = parse_timestamp(&a2a_msg.timestamp)?;
    let elapsed = env.block.time.seconds() - msg_time;
    if elapsed > a2a_msg.ttl_seconds {
        return Err(StdError::generic_err("Message expired"));
    }
    Ok(Response::new()
        .add_attribute("action", "local_delivery")
        .add_attribute("message_id", &a2a_msg.message_id)
        .add_attribute("target_did", &a2a_msg.target_did))
}

fn execute_reply(
    deps: DepsMut,
    env: Env,
    _info: MessageInfo,
    original_message_id: String,
    response_payload: Binary,
) -> StdResult<Response> {
    let store = deps.storage;
    let pending_key = format!("{}{}", PENDING_MSGS, original_message_id);
    let pending_data = store.get(pending_key.as_bytes())
        .ok_or_else(|| StdError::generic_err("Original message not found"))?;
    let original: PendingMessage = bincode2::deserialize(&pending_data)?;
    let original_msg: A2aIbcPacket = bincode2::deserialize(&original.payload.0)?;

    let response_msg = A2aIbcPacket {
        message_id: format!("resp-{}", original_message_id),
        message_type: "a2a/response".to_string(),
        sender_did: original.target_did.clone(),
        sender_chain: get_local_chain_id(store),
        target_did: original.sender_did.clone(),
        target_chain: original.sender_chain.clone(),
        payload: response_payload,
        correlation_id: Some(original_message_id),
        ttl_seconds: 60,
        timestamp: env.block.time.to_string(),
        signature: None,
        schema_version: "1.0".to_string(),
    };

    let channel_id = find_channel_for_target(store, &original.sender_chain)
        .ok_or_else(|| StdError::generic_err(
            format!("No return channel for {}", original.sender_chain)
        ))?;

    let ibc_data = bincode2::serialize(&response_msg)?;
    let reply_msg = IbcMsg::SendPacket {
        channel_id,
        data: Binary::from(ibc_data),
        timeout: IbcTimeout::with_timestamp(env.block.time.plus_seconds(60)),
    };

    let completed = PendingMessage { status: PendingStatus::Completed, ..original };
    store.set(pending_key.as_bytes(), &bincode2::serialize(&completed)?);

    Ok(Response::new()
        .add_message(reply_msg)
        .add_attribute("action", "reply")
        .add_attribute("original_message_id", &original_message_id))
}

fn execute_register_handler(
    deps: DepsMut,
    _env: Env,
    _info: MessageInfo,
    action: String,
    handler_address: String,
) -> StdResult<Response> {
    let key = format!("handler:{}", action);
    deps.storage.set(key.as_bytes(), handler_address.as_bytes());
    Ok(Response::new()
        .add_attribute("action", "register_handler")
        .add_attribute("action_name", &action)
        .add_attribute("handler", &handler_address))
}

// ---------- 查询 ----------

#[derive(Serialize, Deserialize)]
pub enum QueryMsg {
    GetMessageStatus { message_id: String },
    GetChannelInfo { chain_id: Option<String> },
    GetPendingMessages { limit: Option<u32> },
    GetStats {},
}

#[entry_point]
pub fn query(deps: Deps, _env: Env, msg: QueryMsg) -> StdResult<Binary> {
    match msg {
        QueryMsg::GetMessageStatus { message_id } => query_message_status(deps, message_id),
        QueryMsg::GetChannelInfo { chain_id } => query_channel_info(deps, chain_id),
        QueryMsg::GetPendingMessages { limit } => query_pending_messages(deps, limit),
        QueryMsg::GetStats {} => query_stats(deps),
    }
}

fn query_message_status(deps: Deps, message_id: String) -> StdResult<Binary> {
    let pending_key = format!("{}{}", PENDING_MSGS, message_id);
    let delivered_key = format!("{}{}", DELIVERED_MSGS, message_id);
    if let Some(data) = deps.storage.get(pending_key.as_bytes()) {
        let msg: PendingMessage = bincode2::deserialize(&data)?;
        return to_binary(&serde_json::json!({
            "status": "pending",
            "delivery_status": format!("{:?}", msg.status),
            "expires_at": msg.expires_at,
        }));
    }
    if deps.storage.get(delivered_key.as_bytes()).is_some() {
        return to_binary(&serde_json::json!({"status": "delivered"}));
    }
    to_binary(&serde_json::json!({"status": "not_found"}))
}

fn query_channel_info(deps: Deps, chain_id: Option<String>) -> StdResult<Binary> {
    let channels = get_all_a2a_channels(deps.storage);
    let filtered = if let Some(cid) = chain_id {
        channels.into_iter().filter(|c| c.remote_chain_id == cid).collect()
    } else { channels };
    to_binary(&filtered)
}

fn query_pending_messages(deps: Deps, limit: Option<u32>) -> StdResult<Binary> {
    let limit = limit.unwrap_or(20).min(100);
    let mut messages = Vec::new();
    for (_, value) in deps.storage.range(
        Some(PENDING_MSGS.as_bytes().to_vec()), None, cosmwasm_std::Order::Ascending,
    ) {
        if messages.len() >= limit as usize { break; }
        if let Ok(msg) = bincode2::deserialize::<PendingMessage>(&value) {
            messages.push(msg);
        }
    }
    to_binary(&messages)
}

fn query_stats(deps: Deps) -> StdResult<Binary> {
    if let Some(data) = deps.storage.get(A2A_CONFIG.as_bytes()) {
        let config: A2aGlobalConfig = bincode2::deserialize(&data)?;
        let pending_count = deps.storage
            .range(Some(PENDING_MSGS.as_bytes().to_vec()), None, cosmwasm_std::Order::Ascending)
            .count();
        return to_binary(&serde_json::json!({
            "total_messages": config.total_messages,
            "pending_messages": pending_count,
            "active_channels": config.active_channels,
            "chain_id": config.chain_id,
        }));
    }
    to_binary(&serde_json::json!({}))
}

// ---------- IBC 入口点 ----------

#[entry_point]
pub fn ibc_packet_receive(
    deps: DepsMut,
    env: Env,
    msg: IbcPacket,
) -> StdResult<IbcReceiveResponse> {
    let a2a_msg: A2aIbcPacket = bincode2::deserialize(&msg.data)?;
    let msg_time = parse_timestamp(&a2a_msg.timestamp)?;
    let elapsed = env.block.time.seconds() - msg_time;

    if elapsed > a2a_msg.ttl_seconds {
        let ack = A2AAcknowledgement {
            message_id: a2a_msg.message_id,
            status: A2ADeliveryStatus::Expired,
            error: Some("Message expired".to_string()),
        };
        return Ok(IbcReceiveResponse::new()
            .set_ack(bincode2::serialize(&ack)?)
            .add_attribute("action", "a2a_expired"));
    }

    let delivered_key = format!("{}{}", DELIVERED_MSGS, a2a_msg.message_id);
    deps.storage.set(delivered_key.as_bytes(), &msg.data);

    let ack = A2AAcknowledgement {
        message_id: a2a_msg.message_id.clone(),
        status: A2ADeliveryStatus::Delivered,
        error: None,
    };
    Ok(IbcReceiveResponse::new()
        .set_ack(bincode2::serialize(&ack)?)
        .add_attribute("action", "a2a_received")
        .add_attribute("message_id", &a2a_msg.message_id)
        .add_attribute("sender_did", &a2a_msg.sender_did))
}

#[entry_point]
pub fn ibc_packet_ack(
    deps: DepsMut,
    _env: Env,
    ack: IbcPacketAckMsg,
) -> StdResult<Response> {
    if let Ok(a2a_ack) = bincode2::deserialize::<A2AAcknowledgement>(&ack.acknowledgement.data) {
        let pending_key = format!("{}{}", PENDING_MSGS, a2a_ack.message_id);
        if a2a_ack.status == A2ADeliveryStatus::Delivered {
            if let Some(data) = deps.storage.get(pending_key.as_bytes()) {
                if let Ok(mut pending) = bincode2::deserialize::<PendingMessage>(&data) {
                    pending.status = PendingStatus::Completed;
                    deps.storage.set(pending_key.as_bytes(), &bincode2::serialize(&pending)?);
                }
            }
        }
        return Ok(Response::new()
            .add_attribute("action", "a2a_ack")
            .add_attribute("message_id", &a2a_ack.message_id)
            .add_attribute("status", format!("{:?}", a2a_ack.status)));
    }
    Ok(Response::new())
}

#[entry_point]
pub fn ibc_packet_timeout(
    deps: DepsMut,
    _env: Env,
    msg: IbcPacketTimeoutMsg,
) -> StdResult<Response> {
    if let Ok(a2a_msg) = bincode2::deserialize::<A2aIbcPacket>(&msg.packet.data) {
        let pending_key = format!("{}{}", PENDING_MSGS, a2a_msg.message_id);
        if let Some(data) = deps.storage.get(pending_key.as_bytes()) {
            if let Ok(mut pending) = bincode2::deserialize::<PendingMessage>(&data) {
                pending.status = PendingStatus::Failed;
                deps.storage.set(pending_key.as_bytes(), &bincode2::serialize(&pending)?);
            }
        }
        return Ok(Response::new()
            .add_attribute("action", "a2a_timeout")
            .add_attribute("message_id", &a2a_msg.message_id));
    }
    Ok(Response::new())
}

// ---------- 通道处理 ----------

#[entry_point]
pub fn ibc_channel_open(
    deps: DepsMut,
    _env: Env,
    msg: IbcChannelOpenMsg,
) -> StdResult<()> {
    let channel = msg.channel();
    if channel.endpoint.port_id != "a2a" {
        return Err(StdError::generic_err(
            format!("Expected port 'a2a', got '{}'", channel.endpoint.port_id)
        ));
    }
    Ok(())
}

#[entry_point]
pub fn ibc_channel_connect(
    deps: DepsMut,
    _env: Env,
    msg: IbcChannelConnectMsg,
) -> StdResult<Response> {
    let channel = msg.channel();
    let config = A2AChannelConfig {
        channel_id: channel.channel_id.clone(),
        remote_chain_id: channel.counterparty_endpoint.port_id.clone(),
        is_active: true,
        total_messages_sent: 0,
        total_messages_received: 0,
    };
    let key = format!("{}{}", A2A_CHANNELS, channel.channel_id);
    deps.storage.set(key.as_bytes(), &bincode2::serialize(&config)?);
    Ok(Response::new()
        .add_attribute("action", "a2a_channel_connect")
        .add_attribute("channel_id", &channel.channel_id))
}

#[entry_point]
pub fn ibc_channel_close(
    deps: DepsMut,
    _env: Env,
    msg: IbcChannelCloseMsg,
) -> StdResult<Response> {
    let channel = msg.channel();
    let key = format!("{}{}", A2A_CHANNELS, channel.channel_id);
    if let Some(data) = deps.storage.get(key.as_bytes()) {
        if let Ok(mut ch) = bincode2::deserialize::<A2AChannelConfig>(&data) {
            ch.is_active = false;
            deps.storage.set(key.as_bytes(), &bincode2::serialize(&ch)?);
        }
    }
    Ok(Response::new()
        .add_attribute("action", "a2a_channel_close")
        .add_attribute("channel_id", &channel.channel_id))
}

// ---------- 辅助 ----------

#[derive(Serialize, Deserialize, Clone, Debug)]
struct A2aGlobalConfig {
    admin: String,
    chain_id: String,
    total_messages: u64,
    active_channels: u64,
}

fn get_local_chain_id(store: &dyn Storage) -> String {
    store.get(A2A_CONFIG.as_bytes())
        .and_then(|d| bincode2::deserialize::<A2aGlobalConfig>(&d).ok())
        .map(|c| c.chain_id)
        .unwrap_or_default()
}

fn find_channel_for_target(store: &dyn Storage, target_chain: &str) -> Option<String> {
    let channels = get_all_a2a_channels(store);
    channels.into_iter()
        .find(|c| c.remote_chain_id == target_chain && c.is_active)
        .map(|c| c.channel_id)
}

fn get_all_a2a_channels(store: &dyn Storage) -> Vec<A2AChannelConfig> {
    let mut channels = Vec::new();
    for (_, value) in store.range(
        Some(A2A_CHANNELS.as_bytes().to_vec()), None, cosmwasm_std::Order::Ascending,
    ) {
        if let Ok(ch) = bincode2::deserialize::<A2AChannelConfig>(&value) {
            channels.push(ch);
        }
    }
    channels
}

fn update_message_count(store: &mut dyn Storage) {
    if let Some(data) = store.get(A2A_CONFIG.as_bytes()) {
        if let Ok(mut config) = bincode2::deserialize::<A2aGlobalConfig>(&data) {
            config.total_messages += 1;
            store.set(A2A_CONFIG.as_bytes(), &bincode2::serialize(&config).unwrap());
        }
    }
}

fn parse_timestamp(ts: &str) -> StdResult<u64> {
    if let Ok(secs) = ts.parse::<u64>() {
        return Ok(secs);
    }
    Ok(0)
}

5. 跨链支付结算

5.1 ICS-20 转账概览

IBC 的 ICS-20 标准定义了跨链同质化代币转账。AI Agent 可以使用 ICS-20 进行跨链支付结算。

// ============================================================
// 跨链支付结算 (Rust - CosmWasm)
// ============================================================

use cosmwasm_std::{
    entry_point, to_binary, Binary, Coin, Deps, DepsMut, Env, MessageInfo, Response,
    StdResult, StdError, IbcMsg, IbcPacket, IbcReceiveResponse, IbcPacketAckMsg,
    IbcPacketTimeoutMsg, IbcChannelOpenMsg, IbcChannelConnectMsg, IbcChannelCloseMsg,
    IbcOrder, IbcTimeout, Uint128, BankMsg,
};
use serde::{Deserialize, Serialize};

// ---------- 支付数据包 ----------

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct AgentPaymentOrder {
    pub order_id: String,
    pub payer_did: String,
    pub payee_did: String,
    pub payer_chain: String,
    pub payee_chain: String,
    pub amount: Coin,
    pub service_id: String,
    pub memo: String,
    pub timestamp: u64,
    pub timeout_height: Option<u64>,
    pub signature: Binary,
}

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct PaymentSettlementProof {
    pub order_id: String,
    pub ibc_denom: String,
    pub voucher_amount: Uint128,
    pub source_channel: String,
    pub destination_channel: String,
    pub packet_sequence: u64,
    pub ack_status: SettlementStatus,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
pub enum SettlementStatus {
    Pending,
    Confirmed,
    Failed,
    Refunded,
}

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct AgentPaymentChannel {
    pub channel_id: String,
    pub agent_a_did: String,
    pub agent_b_did: String,
    pub chain_a: String,
    pub chain_b: String,
    pub balance_a: Coin,
    pub balance_b: Coin,
    pub total_deposited_a: Coin,
    pub total_deposited_b: Coin,
    pub nonce: u64,
    pub status: ChannelStatus,
    pub created_at: u64,
    pub expires_at: u64,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
pub enum ChannelStatus {
    Open,
    Settling,
    Closed,
}

const PAYMENT_ORDERS: &str = "payment_orders";
const PAYMENT_CHANNELS: &str = "payment_channels";
const SETTLEMENT_PROOFS: &str = "settlement_proofs";

#[derive(Serialize, Deserialize)]
pub struct InstantiateMsg {
    pub admin: String,
    pub native_denom: String,
}

#[derive(Serialize, Deserialize)]
pub enum ExecuteMsg {
    ExecuteCrossChainPayment {
        target_chain: String,
        target_agent_did: String,
        amount: Coin,
        service_id: String,
        memo: String,
    },
    CreatePaymentChannel {
        counterparty_did: String,
        counterparty_chain: String,
        initial_deposit: Coin,
        duration_seconds: u64,
    },
    ClosePaymentChannel { channel_id: String },
    SettlePayment { order_id: String },
    RefundPayment { order_id: String },
}

#[entry_point]
pub fn execute(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    msg: ExecuteMsg,
) -> StdResult<Response> {
    match msg {
        ExecuteMsg::ExecuteCrossChainPayment { target_chain, target_agent_did, amount, service_id, memo } => {
            execute_cross_chain_payment(deps, env, info, target_chain, target_agent_did, amount, service_id, memo)
        }
        ExecuteMsg::CreatePaymentChannel { counterparty_did, counterparty_chain, initial_deposit, duration_seconds } => {
            execute_create_channel(deps, env, info, counterparty_did, counterparty_chain, initial_deposit, duration_seconds)
        }
        ExecuteMsg::ClosePaymentChannel { channel_id } => {
            execute_close_channel(deps, env, info, channel_id)
        }
        ExecuteMsg::SettlePayment { order_id } => execute_settle(deps, env, info, order_id),
        ExecuteMsg::RefundPayment { order_id } => execute_refund(deps, env, info, order_id),
    }
}

/// 执行跨链支付 — 使用 ICS-20 转账
pub fn execute_cross_chain_payment(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    target_chain: String,
    target_agent_did: String,
    amount: Coin,
    service_id: String,
    memo: String,
) -> StdResult<Response> {
    let store = deps.storage;
    let sender_balance = info.funds.iter()
        .find(|c| c.denom == amount.denom)
        .ok_or_else(|| StdError::generic_err("Insufficient funds"))?;
    if sender_balance.amount < amount.amount {
        return Err(StdError::generic_err("Insufficient funds"));
    }

    let order_id = format!("pay-{}-{}", env.block.height, info.sender);
    let order = AgentPaymentOrder {
        order_id: order_id.clone(),
        payer_did: info.sender.to_string(),
        payee_did: target_agent_did.clone(),
        payer_chain: "msg-chain-1".to_string(),
        payee_chain: target_chain.clone(),
        amount: amount.clone(),
        service_id,
        memo,
        timestamp: env.block.time.seconds(),
        timeout_height: Some(env.block.height + 100),
        signature: Binary::default(),
    };

    let order_key = format!("{}{}", PAYMENT_ORDERS, order_id);
    store.set(order_key.as_bytes(), &bincode2::serialize(&order)?);

    // ICS-20 转账
    let ibc_transfer = IbcMsg::Transfer {
        channel_id: get_transfer_channel(&target_chain),
        to_address: target_agent_did.clone(),
        amount: amount.clone(),
        timeout: IbcTimeout::with_timestamp(env.block.time.plus_seconds(600)),
    };

    // 支付通知
    let payment_notification = AgentPaymentNotification {
        order_id: order_id.clone(),
        payer_did: info.sender.to_string(),
        amount: amount.clone(),
        source_chain: "msg-chain-1".to_string(),
        memo: memo.clone(),
    };
    let notify_channel = get_a2a_channel(&target_chain);
    let notify_packet = IbcMsg::SendPacket {
        channel_id: notify_channel,
        data: Binary::from(bincode2::serialize(&payment_notification)?),
        timeout: IbcTimeout::with_timestamp(env.block.time.plus_seconds(600)),
    };

    Ok(Response::new()
        .add_message(ibc_transfer)
        .add_message(notify_packet)
        .add_attribute("action", "cross_chain_payment")
        .add_attribute("order_id", &order_id)
        .add_attribute("payer", &info.sender)
        .add_attribute("payee", &target_agent_did)
        .add_attribute("amount", &amount.to_string())
        .add_attribute("target_chain", &target_chain))
}

fn get_transfer_channel(target_chain: &str) -> String {
    match target_chain {
        "chain-a" => "channel-0".to_string(),
        "chain-b" => "channel-1".to_string(),
        _ => panic!("unknown transfer channel"),
    }
}

fn get_a2a_channel(target_chain: &str) -> String {
    match target_chain {
        "chain-a" => "channel-2".to_string(),
        "chain-b" => "channel-4".to_string(),
        _ => panic!("unknown a2a channel"),
    }
}

#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct AgentPaymentNotification {
    pub order_id: String,
    pub payer_did: String,
    pub amount: Coin,
    pub source_chain: String,
    pub memo: String,
}

fn execute_create_channel(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    counterparty_did: String,
    counterparty_chain: String,
    initial_deposit: Coin,
    duration_seconds: u64,
) -> StdResult<Response> {
    let store = deps.storage;
    let channel_id = format!("paych-{}-{}", info.sender, counterparty_did);
    let channel = AgentPaymentChannel {
        channel_id: channel_id.clone(),
        agent_a_did: info.sender.to_string(),
        agent_b_did: counterparty_did,
        chain_a: "msg-chain-1".to_string(),
        chain_b: counterparty_chain,
        balance_a: initial_deposit.clone(),
        balance_b: Coin::new(0, &initial_deposit.denom),
        total_deposited_a: initial_deposit.clone(),
        total_deposited_b: Coin::new(0, &initial_deposit.denom),
        nonce: 0,
        status: ChannelStatus::Open,
        created_at: env.block.time.seconds(),
        expires_at: env.block.time.seconds() + duration_seconds,
    };
    let key = format!("{}{}", PAYMENT_CHANNELS, channel_id);
    store.set(key.as_bytes(), &bincode2::serialize(&channel)?);
    Ok(Response::new()
        .add_attribute("action", "create_payment_channel")
        .add_attribute("channel_id", &channel_id))
}

fn execute_close_channel(
    deps: DepsMut,
    _env: Env,
    info: MessageInfo,
    channel_id: String,
) -> StdResult<Response> {
    let store = deps.storage;
    let key = format!("{}{}", PAYMENT_CHANNELS, channel_id);
    let data = store.get(key.as_bytes())
        .ok_or_else(|| StdError::generic_err("Channel not found"))?;
    let mut channel: AgentPaymentChannel = bincode2::deserialize(&data)?;
    if channel.agent_a_did != info.sender.to_string()
        && channel.agent_b_did != info.sender.to_string() {
        return Err(StdError::generic_err("Unauthorized"));
    }
    channel.status = ChannelStatus::Closed;
    store.set(key.as_bytes(), &bincode2::serialize(&channel)?);

    let refund_msg = BankMsg::Send {
        to_address: channel.agent_a_did.clone(),
        amount: vec![channel.balance_a.clone()],
    };
    Ok(Response::new()
        .add_message(refund_msg)
        .add_attribute("action", "close_payment_channel")
        .add_attribute("channel_id", &channel_id))
}

fn execute_settle(
    deps: DepsMut,
    _env: Env,
    _info: MessageInfo,
    order_id: String,
) -> StdResult<Response> {
    let store = deps.storage;
    let key = format!("{}{}", PAYMENT_ORDERS, order_id);
    let data = store.get(key.as_bytes())
        .ok_or_else(|| StdError::generic_err("Order not found"))?;
    let order: AgentPaymentOrder = bincode2::deserialize(&data)?;

    let proof = PaymentSettlementProof {
        order_id: order_id.clone(),
        ibc_denom: format!("ibc/{}", order.amount.denom),
        voucher_amount: order.amount.amount,
        source_channel: "channel-0".to_string(),
        destination_channel: "channel-0".to_string(),
        packet_sequence: 0,
        ack_status: SettlementStatus::Confirmed,
    };
    let proof_key = format!("{}{}", SETTLEMENT_PROOFS, order_id);
    store.set(proof_key.as_bytes(), &bincode2::serialize(&proof)?);

    Ok(Response::new()
        .add_attribute("action", "settle_payment")
        .add_attribute("order_id", &order_id))
}

fn execute_refund(
    deps: DepsMut,
    _env: Env,
    _info: MessageInfo,
    order_id: String,
) -> StdResult<Response> {
    let store = deps.storage;
    let key = format!("{}{}", PAYMENT_ORDERS, order_id);
    let data = store.get(key.as_bytes())
        .ok_or_else(|| StdError::generic_err("Order not found"))?;
    let order: AgentPaymentOrder = bincode2::deserialize(&data)?;

    let refund_msg = BankMsg::Send {
        to_address: order.payer_did,
        amount: vec![order.amount],
    };

    store.remove(key.as_bytes());

    Ok(Response::new()
        .add_message(refund_msg)
        .add_attribute("action", "refund_payment")
        .add_attribute("order_id", &order_id))
}

// ---------- 查询 ----------

#[derive(Serialize, Deserialize)]
pub enum QueryMsg {
    GetPaymentOrder { order_id: String },
    GetPaymentChannel { channel_id: String },
    GetSettlementProof { order_id: String },
    ListPaymentChannels { owner_did: Option<String> },
}

#[entry_point]
pub fn query(deps: Deps, _env: Env, msg: QueryMsg) -> StdResult<Binary> {
    match msg {
        QueryMsg::GetPaymentOrder { order_id } => {
            let key = format!("{}{}", PAYMENT_ORDERS, order_id);
            let data = deps.storage.get(key.as_bytes())
                .ok_or_else(|| StdError::generic_err("Order not found"))?;
            to_binary(&bincode2::deserialize::<AgentPaymentOrder>(&data)?)
        }
        QueryMsg::GetPaymentChannel { channel_id } => {
            let key = format!("{}{}", PAYMENT_CHANNELS, channel_id);
            let data = deps.storage.get(key.as_bytes())
                .ok_or_else(|| StdError::generic_err("Channel not found"))?;
            to_binary(&bincode2::deserialize::<AgentPaymentChannel>(&data)?)
        }
        QueryMsg::GetSettlementProof { order_id } => {
            let key = format!("{}{}", SETTLEMENT_PROOFS, order_id);
            let data = deps.storage.get(key.as_bytes())
                .ok_or_else(|| StdError::generic_err("Proof not found"))?;
            to_binary(&bincode2::deserialize::<PaymentSettlementProof>(&data)?)
        }
        QueryMsg::ListPaymentChannels { owner_did } => {
            let mut channels = Vec::new();
            for (_, value) in deps.storage.range(
                Some(PAYMENT_CHANNELS.as_bytes().to_vec()), None, cosmwasm_std::Order::Ascending,
            ) {
                if let Ok(ch) = bincode2::deserialize::<AgentPaymentChannel>(&value) {
                    if let Some(ref owner) = owner_did {
                        if ch.agent_a_did != *owner && ch.agent_b_did != *owner {
                            continue;
                        }
                    }
                    channels.push(ch);
                }
            }
            to_binary(&channels)
        }
    }
}

// ---------- 收款验证 ----------

#[entry_point]
pub fn ibc_packet_receive(
    deps: DepsMut,
    _env: Env,
    msg: IbcPacket,
) -> StdResult<IbcReceiveResponse> {
    // 处理 ICS-20 转账的 IBC 数据包
    // 验证是否为 Agent 支付
    if let Ok(notif) = bincode2::deserialize::<AgentPaymentNotification>(&msg.data) {
        let proof_key = format!("{}proof:{}", SETTLEMENT_PROOFS, notif.order_id);
        let proof = PaymentSettlementProof {
            order_id: notif.order_id.clone(),
            ibc_denom: String::new(),
            voucher_amount: notif.amount.amount,
            source_channel: msg.src.channel_id.clone(),
            destination_channel: msg.dest.channel_id.clone(),
            packet_sequence: msg.sequence,
            ack_status: SettlementStatus::Confirmed,
        };
        deps.storage.set(proof_key.as_bytes(), &bincode2::serialize(&proof)?);

        return Ok(IbcReceiveResponse::new()
            .set_ack(Binary::from(b"{\"status\":\"received\"}"))
            .add_attribute("action", "payment_received")
            .add_attribute("order_id", &notif.order_id)
            .add_attribute("amount", &notif.amount.to_string()));
    }
    Ok(IbcReceiveResponse::new()
        .set_ack(Binary::from(b"{\"status\":\"unknown\"}")))
}

6. 跨链Agent编排

6.1 跨链工作流引擎

AI Agent 编排涉及在多条链上协调多个 Agent 的执行。以下实现一个简单的跨链编排引擎。

# ============================================================
# 跨链 Agent 编排引擎 (Python)
# ============================================================

from dataclasses import dataclass, field
from enum import Enum
from typing import Any, Optional
import asyncio
import json
import uuid


class WorkflowStatus(Enum):
    PENDING = "pending"
    RUNNING = "running"
    COMPLETED = "completed"
    FAILED = "failed"
    TIMEOUT = "timeout"


class TaskType(Enum):
    LOCAL_ACTION = "local_action"
    CROSS_CHAIN_A2A = "cross_chain_a2a"
    CROSS_CHAIN_PAYMENT = "cross_chain_payment"
    SUB_WORKFLOW = "sub_workflow"


@dataclass
class WorkflowTask:
    task_id: str = field(default_factory=lambda: str(uuid.uuid4()))
    task_type: TaskType = TaskType.LOCAL_ACTION
    target_chain: str = ""
    target_agent: str = ""
    action: str = ""
    params: dict[str, Any] = field(default_factory=dict)
    depends_on: list[str] = field(default_factory=list)
    timeout_seconds: int = 120
    max_retries: int = 3
    status: WorkflowStatus = WorkflowStatus.PENDING
    result: Optional[dict[str, Any]] = None
    error: Optional[str] = None


@dataclass
class WorkflowDefinition:
    workflow_id: str = field(default_factory=lambda: str(uuid.uuid4()))
    name: str = ""
    description: str = ""
    tasks: list[WorkflowTask] = field(default_factory=list)
    max_concurrency: int = 5
    timeout_seconds: int = 600
    status: WorkflowStatus = WorkflowStatus.PENDING


class CrossChainWorkflowEngine:
    """跨链工作流编排引擎"""

    def __init__(self, a2a_client: CrossChainA2A):
        self.a2a = a2a_client
        self.active_workflows: dict[str, WorkflowDefinition] = {}
        self.completed_workflows: dict[str, WorkflowDefinition] = {}

    async def execute_workflow(
        self,
        workflow: WorkflowDefinition,
    ) -> WorkflowDefinition:
        workflow.status = WorkflowStatus.RUNNING
        self.active_workflows[workflow.workflow_id] = workflow

        try:
            completed: set[str] = set()
            running: dict[str, asyncio.Task] = {}

            semaphore = asyncio.Semaphore(workflow.max_concurrency)

            while len(completed) < len(workflow.tasks):
                # 找到可以执行的 task
                ready_tasks = [
                    t for t in workflow.tasks
                    if t.task_id not in completed
                    and t.task_id not in running
                    and all(dep in completed for dep in t.depends_on)
                ]

                for task in ready_tasks:
                    async with semaphore:
                        task_ref = asyncio.create_task(
                            self._execute_task(task, workflow)
                        )
                        running[task.task_id] = task_ref

                if not running:
                    break

                done, pending = await asyncio.wait(
                    running.values(),
                    timeout=10,
                    return_when=asyncio.FIRST_COMPLETED,
                )

                for done_task in done:
                    task = done_task.result()
                    completed.add(task.task_id)
                    running.pop(task.task_id, None)

            workflow.status = WorkflowStatus.COMPLETED
            return workflow

        except Exception as e:
            workflow.status = WorkflowStatus.FAILED
            return workflow
        finally:
            self.completed_workflows[workflow.workflow_id] = workflow
            self.active_workflows.pop(workflow.workflow_id, None)

    async def _execute_task(
        self,
        task: WorkflowTask,
        workflow: WorkflowDefinition,
    ) -> WorkflowTask:
        task.status = WorkflowStatus.RUNNING
        attempt = 0

        while attempt < task.max_retries:
            try:
                if task.task_type == TaskType.CROSS_CHAIN_A2A:
                    task.result = await self.a2a.send_message(
                        target_chain=task.target_chain,
                        target_agent_did=task.target_agent,
                        message_payload=task.params,
                        await_response=True,
                        response_timeout=task.timeout_seconds,
                    )
                elif task.task_type == TaskType.LOCAL_ACTION:
                    task.result = await self._execute_local_action(task)
                elif task.task_type == TaskType.CROSS_CHAIN_PAYMENT:
                    task.result = await self._execute_cross_chain_payment(task)

                task.status = WorkflowStatus.COMPLETED
                return task

            except Exception as e:
                attempt += 1
                task.error = str(e)
                if attempt >= task.max_retries:
                    task.status = WorkflowStatus.FAILED
                    return task
                await asyncio.sleep(2 ** attempt)

        return task

    async def _execute_local_action(self, task: WorkflowTask) -> dict:
        return {"status": "ok", "action": task.action, "params": task.params}

    async def _execute_cross_chain_payment(self, task: WorkflowTask) -> dict:
        return {"status": "ok", "payment_id": task.task_id}

6.2 跨链状态协调

// ============================================================
// 跨链状态协调 (Rust)
// ============================================================

use std::collections::HashMap;
use serde::{Deserialize, Serialize};

/// 跨链状态快照
#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct CrossChainState {
    pub workflow_id: String,
    pub chain_id: String,
    pub state_type: StateType,
    pub data: HashMap<String, String>,
    pub block_height: u64,
    pub timestamp: u64,
    pub signature: Vec<u8>,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
pub enum StateType {
    TaskResult,
    WorkflowStatus,
    LockAcquired,
    LockReleased,
}

/// 跨链分布式锁
#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct CrossChainLock {
    pub lock_id: String,
    pub holder_did: String,
    pub chain_id: String,
    pub resource_id: String,
    pub acquired_at: u64,
    pub expires_at: u64,
    pub status: LockStatus,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
pub enum LockStatus {
    Acquired,
    Released,
    Expired,
}

/// 跨链状态同步器
pub struct CrossChainStateSynchronizer {
    pub local_chain: String,
    pub ibc_client: IbcClient,
    pub state_cache: HashMap<String, CrossChainState>,
}

impl CrossChainStateSynchronizer {
    pub fn new(local_chain: String, ibc_client: IbcClient) -> Self {
        Self {
            local_chain,
            ibc_client,
            state_cache: HashMap::new(),
        }
    }

    /// 发送状态到目标链
    pub async fn sync_state(
        &self,
        target_chain: &str,
        state: CrossChainState,
    ) -> Result<(), String> {
        let channel = self.get_channel(target_chain)?;
        let packet = IbcPacket {
            source_port: "state-sync".to_string(),
            source_channel: channel.clone(),
            data: bincode2::serialize(&state).map_err(|e| e.to_string())?,
            timeout_height: None,
            timeout_timestamp: Some(
                std::time::SystemTime::now()
                    .duration_since(std::time::UNIX_EPOCH)
                    .unwrap()
                    .as_secs() + 300
            ),
        };
        self.ibc_client.send_packet(packet).await
    }

    /// 获取通道
    fn get_channel(&self, target_chain: &str) -> Result<String, String> {
        match target_chain {
            "chain-a" => Ok("channel-5".to_string()),
            "chain-b" => Ok("channel-6".to_string()),
            _ => Err(format!("Unknown chain: {}", target_chain)),
        }
    }

    /// 获取跨链分布式锁
    pub async fn acquire_lock(
        &self,
        resource_id: &str,
        holder_did: &str,
        ttl_seconds: u64,
    ) -> Result<CrossChainLock, String> {
        let lock = CrossChainLock {
            lock_id: format!("lock:{}", resource_id),
            holder_did: holder_did.to_string(),
            chain_id: self.local_chain.clone(),
            resource_id: resource_id.to_string(),
            acquired_at: std::time::SystemTime::now()
                .duration_since(std::time::UNIX_EPOCH)
                .unwrap()
                .as_secs(),
            expires_at: std::time::SystemTime::now()
                .duration_since(std::time::UNIX_EPOCH)
                .unwrap()
                .as_secs() + ttl_seconds,
            status: LockStatus::Acquired,
        };
        Ok(lock)
    }
}

struct IbcClient;

impl IbcClient {
    async fn send_packet(&self, _packet: IbcPacket) -> Result<(), String> {
        Ok(())
    }
}

#[derive(Clone, Debug)]
struct IbcPacket {
    source_port: String,
    source_channel: String,
    data: Vec<u8>,
    timeout_height: Option<u64>,
    timeout_timestamp: Option<u64>,
}

6.3 示例:跨链 DEX 套利 Agent

// ============================================================
// 跨链 DEX 套利 Agent (TypeScript)
// ============================================================

interface ArbitrageOpportunity {
    buyChain: string;
    sellChain: string;
    token: string;
    buyPrice: number;
    sellPrice: number;
    profit: number;
    estimatedGas: number;
}

class CrossChainArbitrageAgent {
    private a2aClient: CrossChainA2A;
    private workflowEngine: CrossChainWorkflowEngine;
    private minProfit: number;

    constructor(a2a: CrossChainA2A, workflow: CrossChainWorkflowEngine) {
        this.a2aClient = a2a;
        this.workflowEngine = workflow;
        this.minProfit = 10; // 最小利润(USD)
    }

    async executeArbitrage(): Promise<void> {
        // 1. 聚合多链价格
        const prices = await this.aggregatePrices();

        // 2. 检测套利机会
        const opportunities = this.findOpportunities(prices);

        for (const opp of opportunities) {
            if (opp.profit < this.minProfit) continue;
            await this.executeArbitrageTrade(opp);
        }
    }

    private async aggregatePrices(): Promise<Map<string, Map<string, number>>> {
        const prices = new Map<string, Map<string, number>>();
        const chains = ['msg-chain-1', 'chain-a', 'chain-b'];

        for (const chain of chains) {
            const response = await this.a2aClient.sendMessage(
                chain,
                `did:msg:agent:${chain}-price-oracle`,
                { action: 'query_price', tokens: ['ATOM', 'OSMO', 'USDC'] },
                A2AMessageType.REQUEST,
                60,
                true,
                30
            );
            if (response?.prices) {
                prices.set(chain, new Map(Object.entries(response.prices)));
            }
        }
        return prices;
    }

    private findOpportunities(
        prices: Map<string, Map<string, number>>
    ): ArbitrageOpportunity[] {
        const opportunities: ArbitrageOpportunity[] = [];
        const tokens = ['ATOM', 'OSMO', 'USDC'];
        const chains = Array.from(prices.keys());

        for (const token of tokens) {
            for (let i = 0; i < chains.length; i++) {
                for (let j = 0; j < chains.length; j++) {
                    if (i === j) continue;

                    const buyPrice = prices.get(chains[i])?.get(token) ?? 0;
                    const sellPrice = prices.get(chains[j])?.get(token) ?? 0;

                    if (buyPrice > 0 && sellPrice > buyPrice) {
                        opportunities.push({
                            buyChain: chains[i],
                            sellChain: chains[j],
                            token,
                            buyPrice,
                            sellPrice,
                            profit: (sellPrice - buyPrice) * 100, // 假设100 token
                            estimatedGas: 2,
                        });
                    }
                }
            }
        }
        return opportunities.sort((a, b) => b.profit - a.profit);
    }

    private async executeArbitrageTrade(opp: ArbitrageOpportunity): Promise<void> {
        const workflow = new WorkflowDefinition();
        workflow.name = `Arbitrage: ${opp.token} ${opp.buyChain} -> ${opp.sellChain}`;

        // Task 1: 在买入链上准备流动性
        workflow.tasks.push(new WorkflowTask(
            TaskType.CROSS_CHAIN_A2A,
            opp.buyChain,
            `did:msg:agent:${opp.buyChain}-dex`,
            'prepare_swap',
            { token: opp.token, amount: 100, maxSlippage: 0.01 }
        ));

        // Task 2: 在卖出链上准备卖出
        workflow.tasks.push(new WorkflowTask(
            TaskType.CROSS_CHAIN_A2A,
            opp.sellChain,
            `did:msg:agent:${opp.sellChain}-dex`,
            'prepare_swap',
            { token: opp.token, amount: 100, maxSlippage: 0.01 }
        ));

        // Task 3: 执行买入
        workflow.tasks.push(new WorkflowTask(
            TaskType.CROSS_CHAIN_PAYMENT,
            opp.buyChain,
            `did:msg:agent:${opp.buyChain}-dex`,
            'execute_swap',
            { token: opp.token, amount: 100 }
        ));

        // Task 4: 跨链转移资产
        workflow.tasks.push(new WorkflowTask(
            TaskType.CROSS_CHAIN_PAYMENT,
            opp.sellChain,
            `did:msg:agent:${opp.sellChain}-dex`,
            'ibc_transfer',
            { token: opp.token, amount: 100, targetChain: opp.sellChain }
        ));

        // Task 5: 执行卖出
        workflow.tasks.push(new WorkflowTask(
            TaskType.CROSS_CHAIN_A2A,
            opp.sellChain,
            `did:msg:agent:${opp.sellChain}-dex`,
            'execute_swap',
            { token: opp.token, amount: 100 }
        ));

        await this.workflowEngine.executeWorkflow(workflow);
    }
}

7. 安全考虑

7.1 IBC 数据包超时处理

IBC 数据包超时是跨链通信中的关键安全场景。Agent 必须正确处理超时以确保资金和消息安全。

// ============================================================
// IBC 超时处理 (TypeScript)
// ============================================================

class IbcTimeoutHandler {
    private pendingPackets: Map<string, {
        packet: IbcPacket;
        sentAt: number;
        timeoutMs: number;
        onTimeout: () => Promise<void>;
    }> = new Map();

    addPacket(
        packet: IbcPacket,
        timeoutSeconds: number,
        onTimeout: () => Promise<void>
    ): void {
        const key = `${packet.sourceChannel}:${packet.sequence}`;
        this.pendingPackets.set(key, {
            packet,
            sentAt: Date.now(),
            timeoutMs: timeoutSeconds * 1000,
            onTimeout,
        });
    }

    async checkTimeouts(): Promise<void> {
        const now = Date.now();
        const expired: string[] = [];

        for (const [key, entry] of this.pendingPackets) {
            if (now - entry.sentAt > entry.timeoutMs) {
                try {
                    await entry.onTimeout();
                } catch (e) {
                    console.error(`Timeout handler failed: ${e}`);
                }
                expired.push(key);
            }
        }

        for (const key of expired) {
            this.pendingPackets.delete(key);
        }
    }
}

7.2 跨链重放保护

// ============================================================
// 跨链重放保护 (Rust)
// ============================================================

use std::collections::HashSet;
use sha2::{Sha256, Digest};

/// 重放保护器
pub struct ReplayProtector {
    /// 已处理的 IBC 数据包序列
    processed_packets: HashSet<String>,
    max_cache_size: usize,
}

impl ReplayProtector {
    pub fn new(max_cache_size: usize) -> Self {
        Self {
            processed_packets: HashSet::new(),
            max_cache_size,
        }
    }

    /// 检查数据包是否已被处理
    pub fn is_replay(&mut self, source_chain: &str, sequence: u64, data: &[u8]) -> bool {
        let mut hasher = Sha256::new();
        hasher.update(source_chain.as_bytes());
        hasher.update(&sequence.to_be_bytes());
        hasher.update(data);
        let key = format!("{:x}", hasher.finalize());

        if self.processed_packets.contains(&key) {
            return true; // 重放检测
        }

        self.processed_packets.insert(key);

        // 限制缓存大小
        if self.processed_packets.len() > self.max_cache_size {
            self.processed_packets.clear();
        }

        false
    }
}

7.3 Relayer 信任假设

使用 IBC 时需注意以下信任假设:

假设 风险 缓解措施
Relayer 诚实转发 Relayer 可能选择性审查消息 运行多个独立 Relayer;设置超时回退
轻客户端安全 轻客户端可能被欺骗 使用多级验证;设置信任过期
链验证人诚实 验证人可能作恶 选择高安全性链;使用 ICS 安全模型
IBC 通道安全 通道可能被攻击 限制通道权限;审计通道配置

7.4 安全最佳实践

#!/bin/bash
# ============================================================
# 跨链 Agent 安全配置
# ============================================================

# 1. IBC 超时配置
IBC_TIMEOUT_SECONDS=120
IBC_MAX_RETRIES=3

# 2. 密钥管理
AGENT_KEY_DIR="/etc/agent/keys"
mkdir -p $AGENT_KEY_DIR
chmod 700 $AGENT_KEY_DIR

# 为每条链生成独立密钥
for chain in "msg-chain-1" "chain-a" "chain-b"; do
    msgd keys add "agent-$chain" \
        --keyring-backend file \
        --keyring-dir $AGENT_KEY_DIR
done

# 3. 通道白名单
CHANNEL_WHITELIST=(
    "channel-0:chain-a:agent-registry"
    "channel-1:chain-b:agent-registry"
    "channel-2:chain-a:a2a"
    "channel-3:chain-b:a2a"
    "channel-4:chain-a:transfer"
    "channel-5:chain-b:transfer"
)

for entry in "${CHANNEL_WHITELIST[@]}"; do
    IFS=':' read -r ch chain port <<< "$entry"
    echo "Whitelisting $ch -> $chain:$port"
done

# 4. 费率限制
RATE_LIMIT_PER_SECOND=10
RATE_LIMIT_PER_HOUR=1000

# 5. IBC 数据包大小限制
MAX_IBC_PACKET_SIZE=1048576  # 1MB

附录

A. 跨链端口与通道映射

端口 用途 通道类型 数据包格式
did-registry DID 同步 Ordered IbcDidSyncPacket
agent-registry Agent 注册 Ordered IbcAgentRegistration
a2a Agent 通信 Unordered A2aIbcPacket
transfer ICS-20 转账 Unordered FungibleTokenPacketData
state-sync 状态同步 Ordered CrossChainState

B. MSG Chain 常用命令

# 查询 Agent 注册信息
msgd query wasm contract-state smart $REGISTRY \
    '{"get_agent":{"agent_id":"agent-123"}}' \
    --node https://rpc.msg-chain-1.example.com

# 查询跨链 Agent
msgd query wasm contract-state smart $REGISTRY \
    '{"query_cross_chain_agents":{"chain_id":"chain-a","capability":"trade_execution"}}' \
    --node https://rpc.msg-chain-1.example.com

# 查询 IBC 通道
msgd query ibc channel connections \
    --node https://rpc.msg-chain-1.example.com

# 查询 IBC 数据包
msgd query ibc channel packet \
    transfer channel-0 \
    --node https://rpc.msg-chain-1.example.com

# 发送跨链转账
msgd tx ibc-transfer transfer \
    transfer channel-0 \
    cosmos1recipient... \
    1000000umsg \
    --from agent-key \
    --gas auto \
    --gas-prices 1000000000attoMSG

C. 数据包序列化格式

所有 IBC 数据包使用以下编码格式:


本指南覆盖了在 MSG Chain (msg-chain-1) 上部署 AI Agent 并实现 IBC 跨链通信的完整技术栈。所有代码示例遵循 Rust + TypeScript + Python 三种语言的实现模式,支持 Bech32 前缀 msg。