dApp Docs/AI Agent 多链账户与跨链资产管理指南
Development reference. Not independently verified for production.

AI Agent 多链账户与跨链资产管理指南

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

基于 MSG Chain 的智能代理跨链操作实战


1. 概述

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

去中心化金融(DeFi)和区块链生态已经演变为多链并存的格局。单个链的流动性、用户群体和功能存在天然局限。AI Agent 作为自主执行链上操作的智能程序,必须具备跨链能力才能充分发挥其自动化和智能化优势。以下是 AI Agent 需要多链能力的核心理由:

1.2 MSG Chain 和 IBC

MSG Chain 是基于 Cosmos SDK 构建的应用链,原生支持 IBC(Inter-Blockchain Communication)协议。IBC 是 Cosmos 生态中标准化的跨链通信协议,允许独立的区块链之间安全地传递消息和资产。

通过 IBC,MSG Chain 可以连接以下生态:

1.3 支持的账户类型

AI Agent 在多链环境中需要管理不同类型的账户:

账户类型 地址前缀 派生路径 典型应用
MSG 原生 msg1... m/44'/118'/0'/0/0 MSG Chain 原生操作
IBC 映射 msg1... (IBC denom) 同 MSG IBC 转移资产
EVM 兼容 0x... m/44'/60'/0'/0/0 Ethereum / BSC / Polygon
Solana ... (base58) m/44'/501'/0'/0/0 Solana 生态
其他 Cosmos 链 osmo1..., cosmos1... 等 m/44'/118'/0'/0/0 同密钥派生不同前缀

1.4 本指南的目标读者

本指南面向以下读者:

1.5 文档约定


2. 多链账户架构

2.1 HD 钱包基础

分层确定性(HD)钱包是管理多个链上账户的基础。通过单一助记词,可以派生出任意数量的子密钥对,每个子密钥对对应不同的链或不同的账户索引。

master_seed (mnemonic)
    │
    ├── m/44'/118'/0'/0/0  ->  msg1...  (MSG Chain)
    ├── m/44'/118'/0'/0/0  ->  osmo1... (Osmosis)  [同路径不同前缀]
    ├── m/44'/60'/0'/0/0   ->  0x...    (Ethereum)
    ├── m/44'/501'/0'/0/0  ->  ...      (Solana)
    ├── m/44'/118'/1'/0/0  ->  msg1...  (MSG Chain, account 1)
    └── m/44'/118'/2'/0/0  ->  msg1...  (MSG Chain, account 2)

BIP-44 路径标准:

m / purpose' / coin_type' / account' / change / address_index

Coin type 列表:

链 Coin Type 路径
MSG / Cosmos 118 m/44'/118'/0'/0/0
Ethereum / EVM 60 m/44'/60'/0'/0/0
Solana 501 m/44'/501'/0'/0/0
Bitcoin 0 m/44'/0'/0'/0/0
Polkadot 354 m/44'/354'/0'/0/0
Near 397 m/44'/397'/0'/0/0

2.2 依赖安装

pip install cosmos-sdk-py==0.12.0
pip install bip-utils==2.9.0
pip install eth-account==0.8.0
pip install solders==0.18.0
pip install aiohttp==3.9.0
pip install web3==6.11.0
pip install pydantic==2.5.0
pip install python-dotenv==1.0.0

2.3 MultiChainAccountManager 实现

import json
import os
import hashlib
from typing import Dict, Optional, List, Tuple
from dataclasses import dataclass, field
from enum import Enum

from bip_utils import Bip39SeedGenerator, Bip44, Bip44Coins, Bip44Changes
from eth_account import Account as EthAccount
from eth_account.messages import encode_defunct
from solders.keypair import Keypair as SolKeypair
from dotenv import load_dotenv

load_dotenv()


class ChainType(Enum):
    MSG = "msg"
    COSMOS = "cosmos"
    OSMOSIS = "osmo"
    CELESTIA = "celestia"
    KUJIRA = "kujira"
    JUNO = "juno"
    ETHEREUM = "eth"
    BSC = "bsc"
    POLYGON = "polygon"
    AVALANCHE = "avax"
    SOLANA = "sol"


CHAIN_CONFIG = {
    ChainType.MSG: {
        "coin_type": 118, "prefix": "msg", "hrp": "msg",
        "derivation_path": "m/44'/118'/0'/0/0",
        "ibc_channel": "channel-0",
        "rpc_url": "https://rpc.msgchain.zone",
        "rest_url": "https://api.msgchain.zone",
    },
    ChainType.COSMOS: {
        "coin_type": 118, "prefix": "cosmos", "hrp": "cosmos",
        "derivation_path": "m/44'/118'/0'/0/0",
        "ibc_channel": "channel-1",
        "rpc_url": "https://rpc.cosmos.zone",
        "rest_url": "https://api.cosmos.zone",
    },
    ChainType.OSMOSIS: {
        "coin_type": 118, "prefix": "osmo", "hrp": "osmo",
        "derivation_path": "m/44'/118'/0'/0/0",
        "ibc_channel": "channel-2",
        "rpc_url": "https://rpc.osmosis.zone",
        "rest_url": "https://api.osmosis.zone",
    },
    ChainType.CELESTIA: {
        "coin_type": 118, "prefix": "celestia", "hrp": "celestia",
        "derivation_path": "m/44'/118'/0'/0/0",
        "ibc_channel": "channel-3",
        "rpc_url": "https://rpc.celestia.zone",
        "rest_url": "https://api.celestia.zone",
    },
    ChainType.ETHEREUM: {
        "coin_type": 60, "prefix": "0x", "hrp": "eth",
        "derivation_path": "m/44'/60'/0'/0/0",
        "rpc_url": "https://eth-mainnet.g.alchemy.com/v2/",
        "chain_id": 1,
    },
    ChainType.SOLANA: {
        "coin_type": 501, "prefix": "", "hrp": "sol",
        "derivation_path": "m/44'/501'/0'/0/0",
        "rpc_url": "https://api.mainnet-beta.solana.com",
    },
}


@dataclass
class Account:
    chain: ChainType
    private_key: bytes
    public_key: bytes
    address: str
    derivation_path: str
    address_index: int = 0

    def sign_message(self, message: bytes) -> bytes:
        if self.chain.value in ["eth", "bsc", "polygon", "avax"]:
            eth_acc = EthAccount.from_key(self.private_key.hex())
            signed = eth_acc.sign_message(encode_defunct(primitive=message))
            return signed.signature
        elif self.chain.value == "sol":
            kp = SolKeypair.from_bytes(self.private_key)
            return kp.sign_message(message).signature
        else:
            from cosmos_sdk_py.crypto.keypairs import PrivateKey
            priv = PrivateKey(self.private_key)
            return priv.sign(message)

    def get_public_key_hex(self) -> str:
        return self.public_key.hex()

    def get_private_key_hex(self) -> str:
        return self.private_key.hex()

    def to_dict(self) -> Dict:
        return {
            "chain": self.chain.value,
            "address": self.address,
            "derivation_path": self.derivation_path,
            "public_key": self.get_public_key_hex(),
            "address_index": self.address_index,
        }


class MnemonicKey:
    def __init__(self, mnemonic: str, passphrase: str = ""):
        self.mnemonic = mnemonic.strip()
        self.passphrase = passphrase
        self.seed = Bip39SeedGenerator(self.mnemonic).Generate(self.passphrase)
        self._cache: Dict[str, Bip44] = {}

    def _get_bip44_ctx(self, coin: Bip44Coins) -> Bip44:
        cache_key = f"{coin}"
        if cache_key not in self._cache:
            self._cache[cache_key] = Bip44.FromSeed(self.seed, coin)
        return self._cache[cache_key]

    def derive(self, chain: ChainType, account_index: int = 0,
               address_index: int = 0, external: bool = True) -> Account:
        config = CHAIN_CONFIG[chain]
        coin_type = config["coin_type"]

        if coin_type == 118:
            bip44_ctx = self._get_bip44_ctx(Bip44Coins.COSMOS)
            bip44_acc = (bip44_ctx.Purpose().Coin().Account(account_index)
                         .Change(Bip44Changes.CHAIN_EXT if external else Bip44Changes.CHAIN_INT)
                         .AddressIndex(address_index))
            priv_key_bytes = bip44_acc.PrivateKey().Raw().ToBytes()
            pub_key_bytes = bip44_acc.PublicKey().Raw().Compressed().ToBytes()
            from cosmos_sdk_py.crypto.keypairs import PrivateKey
            priv = PrivateKey(priv_key_bytes)
            address = priv.to_address(prefix=config["prefix"])
        elif coin_type == 60:
            bip44_ctx = self._get_bip44_ctx(Bip44Coins.ETHEREUM)
            bip44_acc = (bip44_ctx.Purpose().Coin().Account(account_index)
                         .Change(Bip44Changes.CHAIN_EXT if external else Bip44Changes.CHAIN_INT)
                         .AddressIndex(address_index))
            priv_key_bytes = bip44_acc.PrivateKey().Raw().ToBytes()
            pub_key_bytes = bip44_acc.PublicKey().Raw().ToBytes()
            eth_acc = EthAccount.from_key(priv_key_bytes.hex())
            address = eth_acc.address
        elif coin_type == 501:
            bip44_ctx = self._get_bip44_ctx(Bip44Coins.SOLANA)
            bip44_acc = (bip44_ctx.Purpose().Coin().Account(account_index)
                         .Change(Bip44Changes.CHAIN_EXT if external else Bip44Changes.CHAIN_INT)
                         .AddressIndex(address_index))
            priv_key_bytes = bip44_acc.PrivateKey().Raw().ToBytes()
            pub_key_bytes = bip44_acc.PublicKey().Raw().ToBytes()
            kp = SolKeypair.from_bytes(priv_key_bytes)
            address = str(kp.pubkey())
        else:
            raise ValueError(f"Unsupported coin type: {coin_type}")

        derivation_path = f"m/44'/{coin_type}'/{account_index}'/{0 if external else 1}/{address_index}"
        return Account(chain=chain, private_key=priv_key_bytes,
                       public_key=pub_key_bytes, address=address,
                       derivation_path=derivation_path, address_index=address_index)


class MultiChainAccountManager:
    def __init__(self, mnemonic: str, passphrase: str = ""):
        self.master_key = MnemonicKey(mnemonic, passphrase)
        self.accounts: Dict[str, Account] = {}
        self.default_accounts: Dict[ChainType, int] = {}

    def derive_account(self, chain: ChainType, account_index: int = 0,
                       address_index: int = 0, external: bool = True,
                       set_default: bool = False) -> Account:
        key = f"{chain.value}:{account_index}:{address_index}"
        if key in self.accounts:
            return self.accounts[key]
        account = self.master_key.derive(chain, account_index, address_index, external)
        self.accounts[key] = account
        if set_default or chain not in self.default_accounts:
            self.default_accounts[chain] = address_index
        return account

    def get_account(self, chain: ChainType, account_index: int = 0,
                    address_index: int = 0) -> Account:
        key = f"{chain.value}:{account_index}:{address_index}"
        if key not in self.accounts:
            raise KeyError(f"Account not derived yet for {key}. Call derive_account first.")
        return self.accounts[key]

    def get_default_address(self, chain: ChainType) -> str:
        account = self.derive_account(chain, 0, 0, True)
        return account.address

    def list_all_addresses(self) -> Dict[str, str]:
        addresses = {}
        for chain in ChainType:
            try:
                addresses[chain.value] = self.get_default_address(chain)
            except Exception:
                pass
        return addresses

    def export_accounts_json(self, filepath: str):
        data = {}
        for key, account in self.accounts.items():
            data[key] = account.to_dict()
        with open(filepath, "w") as f:
            json.dump(data, f, indent=2)

    def import_accounts_json(self, filepath: str):
        with open(filepath) as f:
            data = json.load(f)
        for key, info in data.items():
            chain = ChainType(info["chain"])
            account = Account(
                chain=chain,
                private_key=bytes.fromhex(info["private_key"]),
                public_key=bytes.fromhex(info["public_key"]),
                address=info["address"],
                derivation_path=info["derivation_path"],
                address_index=info["address_index"],
            )
            self.accounts[key] = account

    def derive_batch(self, chain: ChainType, count: int = 10,
                     start_index: int = 0, account_index: int = 0) -> List[Account]:
        accounts = []
        for i in range(start_index, start_index + count):
            acc = self.derive_account(chain, account_index, i)
            accounts.append(acc)
        return accounts


@dataclass
class UnifiedAccountView:
    agent_id: str
    accounts: Dict[ChainType, Account]
    primary_chain: ChainType = ChainType.MSG

    @property
    def primary_address(self) -> str:
        return self.accounts[self.primary_chain].address

    def get_address_for_chain(self, chain: ChainType) -> Optional[str]:
        acc = self.accounts.get(chain)
        return acc.address if acc else None

    def to_table(self) -> str:
        lines = ["Chain | Address | Derivation Path",
                 "------|---------|----------------"]
        for chain, acc in self.accounts.items():
            lines.append(f"{chain.value:12s} | {acc.address:45s} | {acc.derivation_path}")
        return "\n".join(lines)


class MultiChainAgentWallet:
    def __init__(self, mnemonic: str, passphrase: str = ""):
        self.manager = MultiChainAccountManager(mnemonic, passphrase)

    def get_address(self, chain: ChainType) -> str:
        return self.manager.get_default_address(chain)

    def get_account_for_signing(self, chain: ChainType) -> Account:
        return self.manager.derive_account(chain)

    def sign_transaction(self, chain: ChainType, tx_data: Dict) -> bytes:
        account = self.get_account_for_signing(chain)
        message = json.dumps(tx_data, sort_keys=True).encode()
        return account.sign_message(message)

    def get_unified_view(self, agent_id: str) -> UnifiedAccountView:
        accounts = {}
        for chain in ChainType:
            try:
                accounts[chain] = self.manager.derive_account(chain)
            except Exception:
                pass
        return UnifiedAccountView(agent_id=agent_id, accounts=accounts,
                                  primary_chain=ChainType.MSG)


def generate_random_mnemonic(strength: int = 256) -> str:
    from bip_utils import Bip39MnemonicGenerator, Bip39WordsNum
    words_num = Bip39WordsNum.WORDS_NUM_24 if strength == 256 else Bip39WordsNum.WORDS_NUM_12
    return Bip39MnemonicGenerator().FromWordsNumber(words_num)


def validate_mnemonic(mnemonic: str) -> bool:
    from bip_utils import Bip39MnemonicValidator
    return Bip39MnemonicValidator().IsValid(mnemonic)

2.4 地址格式转换

from typing import Optional
import bech32


class AddressConverter:
    @staticmethod
    def to_bech32(prefix: str, pubkey_bytes: bytes) -> str:
        from cosmos_sdk_py.crypto.keypairs import PrivateKey
        priv = PrivateKey(pubkey_bytes)
        return priv.to_address(prefix=prefix)

    @staticmethod
    def decode_bech32(address: str) -> tuple:
        prefix, data = bech32.bech32_decode(address)
        if prefix is None or data is None:
            raise ValueError(f"Invalid bech32 address: {address}")
        return prefix, data

    @staticmethod
    def convert_bech32_prefix(address: str, new_prefix: str) -> str:
        _, data = AddressConverter.decode_bech32(address)
        return bech32.bech32_encode(new_prefix, data)

    @staticmethod
    def is_valid_bech32(address: str, expected_prefix: Optional[str] = None) -> bool:
        try:
            prefix, data = AddressConverter.decode_bech32(address)
            if expected_prefix and prefix != expected_prefix:
                return False
            return len(data) == 32
        except Exception:
            return False

    @staticmethod
    def is_valid_eth_address(address: str) -> bool:
        import re
        return bool(re.match(r'^0x[a-fA-F0-9]{40}$', address))

    @staticmethod
    def eth_to_msg_address(eth_address: str) -> str:
        if not AddressConverter.is_valid_eth_address(eth_address):
            raise ValueError(f"Invalid Ethereum address: {eth_address}")
        data = bytes.fromhex(eth_address[2:])
        five_bit_data = bech32.convertbits(data, 8, 5)
        return bech32.bech32_encode("msg", five_bit_data)

    @staticmethod
    def msg_to_eth_address(msg_address: str) -> str:
        prefix, data = AddressConverter.decode_bech32(msg_address)
        if prefix != "msg":
            raise ValueError(f"Not a MSG address: {msg_address}")
        eight_bit_data = bech32.convertbits(data, 5, 8)
        return "0x" + bytes(eight_bit_data).hex()


class AddressRegistry:
    def __init__(self):
        self._registry: Dict[str, Dict[ChainType, str]] = {}

    def register(self, agent_id: str, chain: ChainType, address: str):
        if agent_id not in self._registry:
            self._registry[agent_id] = {}
        self._registry[agent_id][chain] = address

    def resolve(self, agent_id: str, chain: ChainType) -> Optional[str]:
        agent_addrs = self._registry.get(agent_id)
        return agent_addrs.get(chain) if agent_addrs else None

    def resolve_all(self, agent_id: str) -> Dict[ChainType, str]:
        return self._registry.get(agent_id, {})

    def unregister(self, agent_id: str, chain: Optional[ChainType] = None):
        if chain:
            if agent_id in self._registry:
                self._registry[agent_id].pop(chain, None)
        else:
            self._registry.pop(agent_id, None)

    def to_dict(self) -> Dict:
        result = {}
        for agent_id, addrs in self._registry.items():
            result[agent_id] = {chain.value: addr for chain, addr in addrs.items()}
        return result

    @classmethod
    def from_dict(cls, data: Dict) -> "AddressRegistry":
        registry = cls()
        for agent_id, addrs in data.items():
            for chain_str, addr in addrs.items():
                chain = ChainType(chain_str)
                registry.register(agent_id, chain, addr)
        return registry

3. IBC 资产转移

3.1 IBC 协议概述

IBC(Inter-Blockchain Communication)是 Cosmos 生态中标准化的跨链通信协议。ICS-20 定义了同质化代币转移的标准。IBC 转移的核心流程:

Sender Chain (MSG)                    Receiver Chain (Osmosis)
     │                                        │
     │  1. MsgTransfer                         │
     │  ──────────────────────────────────────>│
     │        (source_port, source_channel,     │
     │         token, sender, receiver,         │
     │         timeout_height)                  │
     │                                        │
     │  2. Packet commitment                    │
     │  <──────────────────────────────────────│
     │                                        │
     │  3. Relayer submits RecvPacket           │
     │  (light client validation)              │
     │                                        │
     │  4. Acknowledgement                      │
     │  <──────────────────────────────────────│
     │                                        │

3.2 IBC Denom 追踪

当资产通过 IBC 转移到其他链时,原始 denom 会被包装为 IBC denom:

IBC denom 格式: ibc/<SHA256_HEX>

计算方式:
  hash = SHA256("transfer/{source_channel}/{denom}")
  ibc_denom = "ibc/" + hash.upper()

示例:

import hashlib
import asyncio
from typing import Optional, List, Dict, Any, Callable, Awaitable
from dataclasses import dataclass, field
from enum import Enum
from decimal import Decimal


class IBCStatus(Enum):
    PENDING = "pending"
    SUCCESS = "success"
    TIMEOUT = "timeout"
    FAILED = "failed"
    ACKNOWLEDGED = "acknowledged"


@dataclass
class IBCDenom:
    original_denom: str
    source_channel: str
    source_port: str = "transfer"
    hash: str = field(init=False)

    def __post_init__(self):
        raw = f"{self.source_port}/{self.source_channel}/{self.original_denom}"
        self.hash = hashlib.sha256(raw.encode()).hexdigest().upper()

    @property
    def ibc_denom(self) -> str:
        return f"ibc/{self.hash}"

    @staticmethod
    def from_ibc_denom(ibc_denom: str) -> "IBCDenom":
        if not ibc_denom.startswith("ibc/"):
            raise ValueError(f"Not an IBC denom: {ibc_denom}")
        return IBCDenom(original_denom="unknown", source_channel="unknown")

    def reverse_denom(self, receiver_channel: str) -> "IBCDenom":
        return IBCDenom(original_denom=self.ibc_denom, source_channel=receiver_channel)


class IBCDenomTracker:
    def __init__(self):
        self._denom_map: Dict[str, IBCDenom] = {}

    def register(self, denom: IBCDenom):
        self._denom_map[denom.ibc_denom] = denom

    def resolve(self, ibc_denom: str) -> Optional[IBCDenom]:
        return self._denom_map.get(ibc_denom)

    def lookup_ibc_denom(self, original_denom: str, channel: str) -> Optional[str]:
        for ibc_d, denom in self._denom_map.items():
            if denom.original_denom == original_denom and denom.source_channel == channel:
                return ibc_d
        return None


@dataclass
class IBCTransferResult:
    status: IBCStatus
    tx_hash: Optional[str] = None
    packet_sequence: Optional[int] = None
    source_channel: Optional[str] = None
    error: Optional[str] = None

    def is_success(self) -> bool:
        return self.status in (IBCStatus.SUCCESS, IBCStatus.ACKNOWLEDGED)

    def to_dict(self) -> Dict:
        return {
            "status": self.status.value,
            "tx_hash": self.tx_hash,
            "packet_sequence": self.packet_sequence,
            "source_channel": self.source_channel,
            "error": self.error,
        }


@dataclass
class Height:
    revision_number: int = 0
    revision_height: int = 0

    def to_dict(self) -> Dict:
        return {
            "revision_number": str(self.revision_number),
            "revision_height": str(self.revision_height),
        }


@dataclass
class Coin:
    amount: int
    denom: str

    def to_dict(self) -> Dict:
        return {"amount": str(self.amount), "denom": self.denom}

    def __str__(self) -> str:
        return f"{self.amount}{self.denom}"


class IBCClient:
    def __init__(self, chain_id: str, rpc_url: str,
                 rest_url: str, grpc_url: Optional[str] = None):
        self.chain_id = chain_id
        self.rpc_url = rpc_url
        self.rest_url = rest_url
        self.grpc_url = grpc_url

    async def get_balance(self, address: str, denom: str) -> Coin:
        import aiohttp
        url = f"{self.rest_url}/cosmos/bank/v1beta1/balances/{address}/by_denom?denom={denom}"
        async with aiohttp.ClientSession() as session:
            async with session.get(url) as resp:
                data = await resp.json()
                coin = data.get("balance", {})
                return Coin(amount=int(coin.get("amount", 0)), denom=coin.get("denom", denom))

    async def get_all_balances(self, address: str) -> List[Coin]:
        import aiohttp
        url = f"{self.rest_url}/cosmos/bank/v1beta1/balances/{address}"
        async with aiohttp.ClientSession() as session:
            async with session.get(url) as resp:
                data = await resp.json()
                balances = data.get("balances", [])
                return [Coin(amount=int(b["amount"]), denom=b["denom"]) for b in balances]

    async def get_latest_height(self) -> int:
        import aiohttp
        url = f"{self.rest_url}/cosmos/base/tendermint/v1beta1/blocks/latest"
        async with aiohttp.ClientSession() as session:
            async with session.get(url) as resp:
                data = await resp.json()
                return int(data["block"]["header"]["height"])

    async def query_ibc_channel(self, port: str, channel: str) -> Dict:
        import aiohttp
        url = f"{self.rest_url}/ibc/core/channel/v1beta1/channels/{channel}/ports/{port}"
        async with aiohttp.ClientSession() as session:
            async with session.get(url) as resp:
                return await resp.json()


class IBCAssetManager:
    def __init__(self, mnemonic: str, chain_id: str = "msg-chain-1",
                 rpc_url: str = "https://rpc.msgchain.zone",
                 rest_url: str = "https://api.msgchain.zone"):
        self.wallet = MultiChainAgentWallet(mnemonic)
        self.client = IBCClient(chain_id, rpc_url, rest_url)
        self.denom_tracker = IBCDenomTracker()
        self._pending_transfers: Dict[str, IBCTransferResult] = {}

    def get_agent_address(self) -> str:
        return self.wallet.get_address(ChainType.MSG)

    def _get_account(self) -> Account:
        return self.wallet.get_account_for_signing(ChainType.MSG)

    async def transfer(self, amount: int, denom: str, receiver: str,
                       source_channel: str = "channel-0",
                       source_port: str = "transfer",
                       timeout_seconds: int = 300) -> IBCTransferResult:
        from cosmos_sdk_py.messages import MsgTransfer
        from cosmos_sdk_py.client import Client
        try:
            current_height = await self.client.get_latest_height()
            timeout_height = current_height + (timeout_seconds // 5)
            msg = MsgTransfer(
                source_port=source_port,
                source_channel=source_channel,
                token=Coin(amount=amount, denom=denom),
                sender=self.get_agent_address(),
                receiver=receiver,
                timeout_height=Height(revision_number=0, revision_height=timeout_height),
            )
            account = self._get_account()
            client = Client(url=self.client.rpc_url, chain_id=self.client.chain_id)
            result = await client.broadcast_tx(msg, account)
            tx_hash = result.get("tx_hash", "")
            packet_seq = result.get("packet_sequence")
            transfer_result = IBCTransferResult(
                status=IBCStatus.PENDING,
                tx_hash=tx_hash,
                packet_sequence=packet_seq,
                source_channel=source_channel,
            )
            self._pending_transfers[tx_hash] = transfer_result
            return transfer_result
        except Exception as e:
            return IBCTransferResult(status=IBCStatus.FAILED, error=str(e))

    async def transfer_to_osmo(self, amount: int,
                               recipient: Optional[str] = None) -> IBCTransferResult:
        if recipient is None:
            recipient = self.wallet.get_address(ChainType.OSMOSIS)
        return await self.transfer(amount=amount, denom="umsg",
                                   receiver=recipient, source_channel="channel-0")

    async def transfer_to_hub(self, amount: int,
                              recipient: Optional[str] = None) -> IBCTransferResult:
        if recipient is None:
            recipient = self.wallet.get_address(ChainType.COSMOS)
        return await self.transfer(amount=amount, denom="umsg",
                                   receiver=recipient, source_channel="channel-1")

    async def transfer_to_celestia(self, amount: int,
                                   recipient: Optional[str] = None) -> IBCTransferResult:
        if recipient is None:
            recipient = self.wallet.get_address(ChainType.CELESTIA)
        return await self.transfer(amount=amount, denom="umsg",
                                   receiver=recipient, source_channel="channel-3")

    async def query_ibc_balance(self, chain: ChainType,
                                original_denom: str = "umsg",
                                source_channel: str = "channel-0") -> Coin:
        ibc_denom_obj = IBCDenom(original_denom=original_denom, source_channel=source_channel)
        self.denom_tracker.register(ibc_denom_obj)
        address = self.wallet.get_address(chain)
        ibc_client = IBCClient(chain_id=f"{chain.value}-chain-1",
                               rpc_url=CHAIN_CONFIG[chain]["rpc_url"],
                               rest_url=CHAIN_CONFIG[chain]["rest_url"])
        return await ibc_client.get_balance(address, ibc_denom_obj.ibc_denom)

    async def query_all_cross_chain_balances(self) -> Dict[str, List[Coin]]:
        results = {}
        for chain in [ChainType.MSG, ChainType.OSMOSIS, ChainType.COSMOS, ChainType.CELESTIA]:
            try:
                ibc_client = IBCClient(chain_id=f"{chain.value}-chain-1",
                                       rpc_url=CHAIN_CONFIG[chain]["rpc_url"],
                                       rest_url=CHAIN_CONFIG[chain]["rest_url"])
                address = self.wallet.get_address(chain)
                balances = await ibc_client.get_all_balances(address)
                results[chain.value] = balances
            except Exception:
                results[chain.value] = [Coin(amount=0, denom="error")]
        return results

    async def check_transfer_status(self, tx_hash: str) -> IBCStatus:
        import aiohttp
        url = f"{self.client.rest_url}/cosmos/tx/v1beta1/txs/{tx_hash}"
        async with aiohttp.ClientSession() as session:
            async with session.get(url) as resp:
                if resp.status == 200:
                    data = await resp.json()
                    tx_response = data.get("tx_response", {})
                    code = tx_response.get("code", 0)
                    return IBCStatus.SUCCESS if code == 0 else IBCStatus.FAILED
                return IBCStatus.PENDING

    async def wait_for_transfer(self, tx_hash: str, timeout: int = 60,
                                poll_interval: int = 5) -> IBCTransferResult:
        elapsed = 0
        while elapsed < timeout:
            status = await self.check_transfer_status(tx_hash)
            if status in (IBCStatus.SUCCESS, IBCStatus.FAILED):
                result = self._pending_transfers.get(tx_hash, IBCTransferResult(status=status))
                result.status = status
                return result
            await asyncio.sleep(poll_interval)
            elapsed += poll_interval
        result = self._pending_transfers.get(tx_hash, IBCTransferResult(status=IBCStatus.TIMEOUT))
        result.status = IBCStatus.TIMEOUT
        return result

    def compute_ibc_denom(self, denom: str, source_channel: str,
                          source_port: str = "transfer") -> str:
        denom_obj = IBCDenom(original_denom=denom, source_channel=source_channel,
                             source_port=source_port)
        self.denom_tracker.register(denom_obj)
        return denom_obj.ibc_denom

    async def transfer_back_to_msg(self, chain: ChainType, ibc_denom: str,
                                   amount: int, source_channel: str,
                                   recipient: Optional[str] = None) -> IBCTransferResult:
        if recipient is None:
            recipient = self.get_agent_address()
        return await self.transfer(amount=amount, denom=ibc_denom,
                                   receiver=recipient, source_channel=source_channel)

3.3 IBC 超时处理

class IBCTimeoutHandler:
    TIMEOUT_STRATEGIES = {
        "retry": "Retry the transfer with updated parameters",
        "refund": "Wait for automatic refund on source chain",
        "escalate": "Escalate to human operator",
    }

    def __init__(self, ibc_manager: IBCAssetManager):
        self.ibc = ibc_manager

    async def handle_timeout(self, packet: Dict, strategy: str = "retry") -> Dict:
        if strategy == "retry":
            return await self._retry_packet(packet)
        elif strategy == "refund":
            return await self._wait_for_refund(packet)
        elif strategy == "escalate":
            return self._escalate(packet)
        else:
            raise ValueError(f"Unknown timeout strategy: {strategy}")

    async def _retry_packet(self, packet: Dict) -> Dict:
        packet_data = packet.get("packet", {})
        data = packet_data.get("data", {})
        if isinstance(data, str):
            import base64
            try:
                data_bytes = base64.b64decode(data)
                data = json.loads(data_bytes)
            except Exception:
                pass
        amount = int(data.get("amount", 0))
        denom = data.get("denom", "umsg")
        receiver = data.get("receiver", "")
        new_result = await self.ibc.transfer(amount=amount, denom=denom,
                                             receiver=receiver, timeout_seconds=600)
        return {"original_packet": packet, "retry_result": new_result.to_dict(),
                "note": "Retried with increased timeout (600s)"}

    async def _wait_for_refund(self, packet: Dict) -> Dict:
        await asyncio.sleep(10)
        return {"status": "refund_pending",
                "note": "Tokens will be refunded automatically after timeout height is reached"}

    def _escalate(self, packet: Dict) -> Dict:
        return {"status": "escalated", "packet": packet,
                "action_required": "Manual intervention needed for IBC timeout"}

3.4 IBC 费用估算

class IBCFeeEstimator:
    def __init__(self, client: IBCClient):
        self.client = client

    async def estimate_gas(self, msg_type: str = "transfer") -> int:
        gas_table = {
            "transfer": 200_000, "recv_packet": 300_000,
            "acknowledgement": 150_000, "timeout": 200_000,
            "channel_open_init": 500_000, "channel_open_try": 500_000,
            "channel_open_ack": 300_000, "channel_open_confirm": 300_000,
        }
        return gas_table.get(msg_type, 200_000)

    async def estimate_fee(self, msg_type: str = "transfer",
                           gas_price: Optional[Decimal] = None) -> Dict:
        if gas_price is None:
            gas_price = Decimal("1000000000")
        gas = await self.estimate_gas(msg_type)
        fee = Decimal(str(gas)) * gas_price
        return {"gas": gas, "gas_price": str(gas_price),
                "fee_umsg": int(fee * Decimal("1000000")),
                "fee_msg": float(fee), "msg_type": msg_type}

4. 跨链桥集成

4.1 桥接方案概览

IBC 仅限于 Cosmos 生态链。对于非 IBC 链(Ethereum, Solana, BSC 等),需要依赖第三方跨链桥:

桥 支持的目标链 机制 信任假设
Axelar Ethereum, BSC, Polygon, Avalanche GMP + 验证者网络 验证者阈值签名
Wormhole Ethereum, Solana, BSC, Polygon, Avalanche, Sui, Aptos VAA + 守护者网络 守护者多签
LayerZero 30+ 链 预言机 + 中继器 双重信任假设
Hyperlane EVM, Cosmos 主权共识 可配置安全模型

4.2 Axelar GMP 集成

import aiohttp
import hashlib
from typing import Optional, Dict, Any, List
from dataclasses import dataclass
from enum import Enum


class BridgeProvider(Enum):
    AXELAR = "axelar"
    WORMHOLE = "wormhole"
    LAYERZERO = "layerzero"
    HYPERLANE = "hyperlane"


class BridgeStatus(Enum):
    PENDING = "pending"
    CONFIRMED = "confirmed"
    EXECUTED = "executed"
    FAILED = "failed"
    REFUNDED = "refunded"


@dataclass
class BridgeConfig:
    provider: BridgeProvider
    contract_address: str
    rpc_url: str
    gas_limit: int = 500_000
    gas_multiplier: float = 1.2


DEFAULT_BRIDGE_CONFIGS = {
    BridgeProvider.AXELAR: {
        "msg_contract": "msg1axelar4x5n8k7l9m0q3r6s8t2v4wgateway",
        "eth_contract": "0x4f4560A766bBCfF8dC8bC2d2b8C2d1c3F5e8A7B9",
        "gas_service": "msg1axelargas35n8k7l9m0q3r6s8t2v4w",
    },
    BridgeProvider.WORMHOLE: {
        "msg_contract": "msg1wormholecore1l9m0q3r6s8t2v4w",
        "eth_contract": "0x98f3c9e6E3fAce36aAdE3d4bD8c3A2f1E6d9B7C",
        "sol_contract": "worm2ZoG2kUd4vFXhWq9T5KLGmPqY8e3a",
    },
}


class AxelarGMPClient:
    def __init__(self, gateway_contract: str, gas_service_contract: str,
                 rpc_url: str, chain_id: str = "msg-chain-1"):
        self.gateway = gateway_contract
        self.gas_service = gas_service_contract
        self.rpc_url = rpc_url
        self.chain_id = chain_id

    async def call_contract(self, destination_chain: str, destination_address: str,
                            payload: bytes, gas_limit: int = 200_000,
                            gas_fee: Optional[int] = None) -> Dict:
        from cosmos_sdk_py.messages import MsgExecuteContract
        if gas_fee is None:
            gas_fee = gas_limit * 100
        execute_msg = {
            "send_message": {
                "destination_chain": destination_chain,
                "destination_address": destination_address,
                "payload": payload.hex(),
                "payload_hash": hashlib.sha256(payload).hexdigest(),
                "gas_limit": str(gas_limit),
                "gas_fee": str(gas_fee),
            }
        }
        return {"msg": execute_msg, "destination_chain": destination_chain,
                "payload_hash": hashlib.sha256(payload).hexdigest(), "gas_limit": gas_limit}


class CrossChainBridge:
    def __init__(self, mnemonic: str, provider: BridgeProvider = BridgeProvider.AXELAR):
        self.wallet = MultiChainAgentWallet(mnemonic)
        self.provider = provider
        self.configs = DEFAULT_BRIDGE_CONFIGS
        self._pending_bridges: Dict[str, Dict] = {}

    def get_agent_address(self, chain: ChainType) -> str:
        return self.wallet.get_address(chain)

    async def bridge_to_ethereum(self, amount: int, eth_address: Optional[str] = None,
                                 denom: str = "umsg", gas_limit: int = 200_000) -> Dict:
        if eth_address is None:
            eth_address = self.wallet.get_address(ChainType.ETHEREUM)
        payload = json.dumps({
            "action": "bridge_token", "denom": denom,
            "amount": str(amount), "receiver": eth_address,
        }).encode()
        axelar = AxelarGMPClient(
            gateway_contract=self.configs[BridgeProvider.AXELAR]["msg_contract"],
            gas_service_contract=self.configs[BridgeProvider.AXELAR]["gas_service"],
            rpc_url=CHAIN_CONFIG[ChainType.MSG]["rpc_url"],
        )
        result = await axelar.call_contract(
            destination_chain="Ethereum",
            destination_address=self.configs[BridgeProvider.AXELAR]["eth_contract"],
            payload=payload, gas_limit=gas_limit,
        )
        bridge_id = hashlib.sha256(f"{amount}{eth_address}{result['payload_hash']}".encode()).hexdigest()[:16]
        self._pending_bridges[bridge_id] = {
            "status": BridgeStatus.PENDING.value, "amount": amount,
            "denom": denom, "destination": "Ethereum", "receiver": eth_address,
        }
        return {"bridge_id": bridge_id, "provider": self.provider.value,
                "source": "msg-chain-1", "destination": "Ethereum",
                "amount": str(amount), "denom": denom, "receiver": eth_address,
                "axelar_result": result}

    async def bridge_to_solana(self, amount: int, sol_address: Optional[str] = None,
                               denom: str = "umsg") -> Dict:
        if sol_address is None:
            sol_address = self.wallet.get_address(ChainType.SOLANA)
        payload = {
            "action": "bridge_token", "denom": denom,
            "amount": str(amount), "receiver": sol_address,
            "nonce": int(hashlib.sha256(os.urandom(32)).hexdigest()[:8], 16),
        }
        bridge_id = f"wormhole_{payload['nonce']}"
        self._pending_bridges[bridge_id] = {
            "status": BridgeStatus.PENDING.value, "amount": amount,
            "denom": denom, "destination": "Solana", "receiver": sol_address,
        }
        return {"bridge_id": bridge_id, "provider": BridgeProvider.WORMHOLE.value,
                "source": "msg-chain-1", "destination": "Solana",
                "amount": str(amount), "denom": denom, "receiver": sol_address}

    async def bridge_to_bsc(self, amount: int, bsc_address: Optional[str] = None,
                            denom: str = "umsg") -> Dict:
        if bsc_address is None:
            bsc_address = self.wallet.get_address(ChainType.BSC)
        return await self.bridge_to_ethereum(amount=amount, eth_address=bsc_address, denom=denom)

    async def bridge_to_polygon(self, amount: int, polygon_address: Optional[str] = None,
                                denom: str = "umsg") -> Dict:
        if polygon_address is None:
            polygon_address = self.wallet.get_address(ChainType.POLYGON)
        return await self.bridge_to_ethereum(amount=amount, eth_address=polygon_address, denom=denom)

    async def check_bridge_status(self, bridge_id: str) -> Optional[Dict]:
        return self._pending_bridges.get(bridge_id)

    async def query_bridge_fee(self, destination_chain: str, amount: int) -> Dict:
        base_fees = {"Ethereum": 0.01, "Solana": 0.001, "BSC": 0.005,
                     "Polygon": 0.003, "Avalanche": 0.004}
        fee = base_fees.get(destination_chain, 0.01) * (amount / 1_000_000)
        return {"destination_chain": destination_chain, "amount_msg": amount / 1_000_000,
                "estimated_fee_msg": fee, "estimated_fee_umsg": int(fee * 1_000_000),
                "provider": self.provider.value}

4.3 桥接安全考虑

class BridgeSecurityCheck:
    MINIMUM_CONFIRMATIONS = {BridgeProvider.AXELAR: 12, BridgeProvider.WORMHOLE: 1}
    MAX_BRIDGE_AMOUNT = 1_000_000_000_000

    @staticmethod
    async def validate_bridge_request(amount: int, destination_chain: str,
                                      provider: BridgeProvider) -> List[str]:
        warnings = []
        if amount > BridgeSecurityCheck.MAX_BRIDGE_AMOUNT:
            warnings.append(f"Large bridge amount: {amount} umsg. Consider splitting.")
        suspicious = ["unknown", "0x0000000000000000000000000000000000000000"]
        if destination_chain.lower() in suspicious:
            warnings.append(f"Suspicious destination: {destination_chain}")
        if provider == BridgeProvider.AXELAR:
            warnings.append("Axelar requires trusted validator set. Verify before large transfers.")
        if provider == BridgeProvider.WORMHOLE:
            warnings.append("Wormhole uses guardian network. Check guardian set status.")
        return warnings

    @staticmethod
    def validate_receiver_address(chain: ChainType, address: str) -> bool:
        if chain in [ChainType.ETHEREUM, ChainType.BSC, ChainType.POLYGON]:
            import re
            return bool(re.match(r'^0x[a-fA-F0-9]{40}$', address))
        elif chain == ChainType.SOLANA:
            from solders.pubkey import Pubkey
            try:
                Pubkey.from_string(address)
                return True
            except Exception:
                return False
        elif chain in [ChainType.MSG, ChainType.COSMOS, ChainType.OSMOSIS]:
            return address.startswith(f"{chain.value}1") and len(address) > 30
        return False

4.4 Wormhole VAA 处理

class WormholeVAA:
    def __init__(self, vaa_bytes: bytes):
        self.raw = vaa_bytes
        self._parsed = self._parse()

    def _parse(self) -> Dict:
        version = self.raw[0]
        guardian_set_index = int.from_bytes(self.raw[1:5], "big")
        sig_count = self.raw[5]
        signatures = []
        offset = 6
        for _ in range(sig_count):
            sig_index = self.raw[offset]
            sig_data = self.raw[offset + 1:offset + 66]
            signatures.append({"index": sig_index, "signature": sig_data.hex()})
            offset += 66
        timestamp = int.from_bytes(self.raw[offset:offset + 4], "big")
        nonce = int.from_bytes(self.raw[offset + 4:offset + 8], "big")
        emitter_chain = int.from_bytes(self.raw[offset + 8:offset + 10], "big")
        emitter_address = self.raw[offset + 10:offset + 42].hex()
        sequence = int.from_bytes(self.raw[offset + 42:offset + 50], "big")
        consistency_level = self.raw[offset + 50]
        payload = self.raw[offset + 51:]
        return {"version": version, "guardian_set_index": guardian_set_index,
                "signature_count": sig_count, "signatures": signatures,
                "timestamp": timestamp, "nonce": nonce,
                "emitter_chain": emitter_chain, "emitter_address": emitter_address,
                "sequence": sequence, "consistency_level": consistency_level,
                "payload": payload.hex()}

    def verify_signatures(self, guardian_pubkeys: List[str]) -> bool:
        required = len(guardian_pubkeys) * 2 // 3 + 1
        return self._parsed["signature_count"] >= required

    def get_payload(self) -> bytes:
        return bytes.fromhex(self._parsed["payload"])

    def to_dict(self) -> Dict:
        return self._parsed


class WormholeClient:
    def __init__(self, core_contract: str, rpc_url: str):
        self.core_contract = core_contract
        self.rpc_url = rpc_url

    async def publish_message(self, sender: str, message: bytes,
                              nonce: int = 0, consistency_level: int = 32) -> Dict:
        return {"action": "publish_message", "sender": sender, "nonce": nonce,
                "payload_hash": hashlib.sha256(message).hexdigest()}

    async def parse_vaa(self, vaa_hex: str) -> WormholeVAA:
        vaa_bytes = bytes.fromhex(vaa_hex)
        return WormholeVAA(vaa_bytes)

    async def verify_vaa(self, vaa: WormholeVAA) -> bool:
        import aiohttp
        url = f"{self.rpc_url}/cosmwasm/wasm/v1/contract/{self.core_contract}/smart/verify_vaa"
        payload = {"vaa": vaa.raw.hex()}
        async with aiohttp.ClientSession() as session:
            async with session.get(url, json=payload) as resp:
                data = await resp.json()
                return data.get("valid", False)

5. 跨链 Agent 身份

5.1 DID 基础

去中心化标识符(DID)为 AI Agent 提供跨链统一身份。MSG Chain 上的 Agent 使用 did:msg 方法:

did:msg:<agent_id>

示例:
did:msg:agent1q2w3e4r5t6y7u8i9o0p

DID Document 包含 Agent 在所有链上的地址和公钥:

{
  "@context": "https://www.w3.org/ns/did/v1",
  "id": "did:msg:agent1q2w3e4r5t6y7u8i9o0p",
  "verificationMethod": [
    {
      "id": "did:msg:agent1q2w3e4r5t6y7u8i9o0p#keys-msg",
      "type": "EcdsaSecp256k1VerificationKey2019",
      "controller": "did:msg:agent1q2w3e4r5t6y7u8i9o0p",
      "publicKeyHex": "03abc123..."
    }
  ],
  "linkedAddresses": [
    {"chain": "msg-chain-1", "address": "msg1agent2y5x4x5n8k7l9m0q3r6s8t2v4w"},
    {"chain": "ethereum", "address": "0x1234..."},
    {"chain": "solana", "address": "..."}
  ]
}

5.2 DID 身份实现

import json
import hashlib
import asyncio
from typing import Optional, Dict, List, Any
from dataclasses import dataclass, field
from datetime import datetime, timedelta
from enum import Enum


class DIDMethod(Enum):
    MSG = "did:msg"
    ETH = "did:eth"
    SOL = "did:sol"
    KEY = "did:key"
    WEB = "did:web"


@dataclass
class VerificationMethod:
    id: str
    type: str
    controller: str
    public_key_hex: str

    def to_dict(self) -> Dict:
        return {"id": self.id, "type": self.type,
                "controller": self.controller, "publicKeyHex": self.public_key_hex}


@dataclass
class LinkedAddress:
    chain: str
    address: str
    proof: Optional[str] = None
    timestamp: Optional[int] = None

    def to_dict(self) -> Dict:
        result = {"chain": self.chain, "address": self.address}
        if self.proof:
            result["proof"] = self.proof
        if self.timestamp:
            result["timestamp"] = self.timestamp
        return result


@dataclass
class DIDDocument:
    id: str
    verification_methods: List[VerificationMethod] = field(default_factory=list)
    linked_addresses: List[LinkedAddress] = field(default_factory=list)
    authentication: List[str] = field(default_factory=list)
    assertion_method: List[str] = field(default_factory=list)
    service_endpoints: Dict[str, str] = field(default_factory=dict)
    created: str = field(default_factory=lambda: datetime.utcnow().isoformat() + "Z")
    updated: str = field(default_factory=lambda: datetime.utcnow().isoformat() + "Z")

    def to_dict(self) -> Dict:
        return {
            "@context": "https://www.w3.org/ns/did/v1",
            "id": self.id,
            "created": self.created,
            "updated": self.updated,
            "verificationMethod": [vm.to_dict() for vm in self.verification_methods],
            "authentication": self.authentication or [vm.id for vm in self.verification_methods],
            "assertionMethod": self.assertion_method or [vm.id for vm in self.verification_methods],
            "linkedAddresses": [la.to_dict() for la in self.linked_addresses],
            "service": [{"id": f"{self.id}#{k}", "type": k, "serviceEndpoint": v}
                        for k, v in self.service_endpoints.items()],
        }


@dataclass
class VerifiableCredential:
    context: List[str] = field(default_factory=lambda: [
        "https://www.w3.org/2018/credentials/v1",
    ])
    id: str = ""
    type: List[str] = field(default_factory=lambda: ["VerifiableCredential"])
    issuer: str = ""
    issuance_date: str = ""
    credential_subject: Dict = field(default_factory=dict)
    proof: Dict = field(default_factory=dict)

    def to_dict(self) -> Dict:
        return {
            "@context": self.context,
            "id": self.id,
            "type": self.type,
            "issuer": self.issuer,
            "issuanceDate": self.issuance_date,
            "credentialSubject": self.credential_subject,
            "proof": self.proof,
        }


class DIDRegistry:
    def __init__(self, contract_address: str, rpc_url: str):
        self.contract = contract_address
        self.rpc_url = rpc_url
        self._cache: Dict[str, DIDDocument] = {}

    async def resolve(self, did: str) -> Optional[DIDDocument]:
        if did in self._cache:
            return self._cache[did]
        if not did.startswith("did:msg:"):
            raise ValueError(f"Unsupported DID method: {did}")
        import aiohttp
        url = f"{self.rpc_url}/cosmwasm/wasm/v1/contract/{self.contract}/smart/resolve_did"
        payload = {"did": did}
        async with aiohttp.ClientSession() as session:
            async with session.get(url, json=payload) as resp:
                data = await resp.json()
                if "error" in data:
                    return None
                doc = self._parse_document(data)
                self._cache[did] = doc
                return doc

    def _parse_document(self, data: Dict) -> DIDDocument:
        return DIDDocument(
            id=data.get("id", ""),
            verification_methods=[VerificationMethod(**vm)
                                  for vm in data.get("verificationMethod", [])],
            linked_addresses=[LinkedAddress(**la)
                              for la in data.get("linkedAddresses", [])],
            service_endpoints={s["type"]: s["serviceEndpoint"]
                               for s in data.get("service", [])},
        )


class CrossChainAgentIdentity:
    def __init__(self, mnemonic: str, agent_id: Optional[str] = None,
                 registry_contract: Optional[str] = None):
        self.wallet = MultiChainAgentWallet(mnemonic)
        self.msg_address = self.wallet.get_address(ChainType.MSG)
        if agent_id:
            self.agent_id = agent_id
        else:
            addr_hash = hashlib.sha256(self.msg_address.encode()).hexdigest()[:24]
            self.agent_id = f"agent{addr_hash}"
        self.did = f"did:msg:{self.agent_id}"
        self.chain_addresses: Dict[str, str] = {}
        self._credentials: List[VerifiableCredential] = []
        self._did_document: Optional[DIDDocument] = None
        self._registry = DIDRegistry(
            registry_contract or "msg1registryxxxxxxxxxxxxxxxxxxxxx",
            CHAIN_CONFIG[ChainType.MSG]["rpc_url"],
        )

    def build_did_document(self) -> DIDDocument:
        msg_account = self.wallet.get_account_for_signing(ChainType.MSG)
        verification_methods = [
            VerificationMethod(id=f"{self.did}#keys-msg",
                               type="EcdsaSecp256k1VerificationKey2019",
                               controller=self.did,
                               public_key_hex=msg_account.get_public_key_hex()),
        ]
        eth_account = self.wallet.get_account_for_signing(ChainType.ETHEREUM)
        verification_methods.append(
            VerificationMethod(id=f"{self.did}#keys-eth",
                               type="EcdsaSecp256k1VerificationKey2019",
                               controller=self.did,
                               public_key_hex=eth_account.get_public_key_hex()),
        )
        linked_addresses = [
            LinkedAddress(chain="msg-chain-1", address=self.msg_address,
                          timestamp=int(datetime.utcnow().timestamp())),
        ]
        for chain_str, addr in self.chain_addresses.items():
            linked_addresses.append(
                LinkedAddress(chain=chain_str, address=addr,
                              timestamp=int(datetime.utcnow().timestamp())),
            )
        self._did_document = DIDDocument(
            id=self.did,
            verification_methods=verification_methods,
            linked_addresses=linked_addresses,
            service_endpoints={
                "MsgChainEndpoint": f"https://agent.{self.agent_id}.msgchain.zone",
            },
        )
        return self._did_document

    async def link_address(self, chain: ChainType,
                           address: Optional[str] = None) -> VerifiableCredential:
        target_address = address or self.wallet.get_address(chain)
        vc = await self.issue_linking_credential(chain, target_address)
        self._credentials.append(vc)
        self.chain_addresses[chain.value] = target_address
        self.build_did_document()
        return vc

    async def issue_linking_credential(self, chain: ChainType,
                                       address: str) -> VerifiableCredential:
        credential_id = f"{self.did}/credentials/{chain.value}/{int(datetime.utcnow().timestamp())}"
        credential_subject = {
            "id": self.did,
            "linkedAddress": {"chain": chain.value, "address": address},
        }
        proof_payload = json.dumps({
            "credentialSubject": credential_subject, "issuer": self.did,
            "issuedAt": datetime.utcnow().isoformat() + "Z",
        }, sort_keys=True).encode()
        msg_account = self.wallet.get_account_for_signing(ChainType.MSG)
        signature = msg_account.sign_message(proof_payload)
        proof = {
            "type": "EcdsaSecp256k1Signature2019",
            "created": datetime.utcnow().isoformat() + "Z",
            "proofPurpose": "assertionMethod",
            "verificationMethod": f"{self.did}#keys-msg",
            "signature": signature.hex(),
        }
        return VerifiableCredential(
            id=credential_id, issuer=self.did,
            issuance_date=datetime.utcnow().isoformat() + "Z",
            credential_subject=credential_subject, proof=proof,
        )

    async def anchor_credential(self, vc: VerifiableCredential) -> Dict:
        vc_dict = vc.to_dict()
        vc_json = json.dumps(vc_dict)
        return {"status": "anchored", "credential_id": vc.id,
                "hash": hashlib.sha256(vc_json.encode()).hexdigest()}

    async def verify_cross_chain_action(self, chain: ChainType, action: Dict,
                                        signature: bytes) -> bool:
        expected_address = self.chain_addresses.get(chain.value)
        if not expected_address:
            doc = await self._registry.resolve(self.did)
            if not doc:
                return False
            for la in doc.linked_addresses:
                if la.chain == chain.value:
                    expected_address = la.address
                    break
        if not expected_address:
            return False
        action_json = json.dumps(action, sort_keys=True).encode()
        if chain.value in ["eth", "bsc", "polygon", "avax"]:
            from eth_account import Account as EthAccount
            from eth_account.messages import encode_defunct
            recovered = EthAccount.recover_message(
                encode_defunct(primitive=action_json), signature=signature,
            )
            return recovered.lower() == expected_address.lower()
        elif chain.value == "sol":
            from solders.message import Message
            try:
                kp = SolKeypair()
                kp.verify_message(Message(action_json), signature)
                return True
            except Exception:
                return False
        else:
            from cosmos_sdk_py.crypto.keypairs import PublicKey
            pubkey = PublicKey(signature[:33])
            return pubkey.verify_signature(action_json, signature[33:])

    async def resolve_did(self) -> Optional[DIDDocument]:
        return await self._registry.resolve(self.did)

    def get_did_document_json(self, pretty: bool = True) -> str:
        if not self._did_document:
            self.build_did_document()
        return json.dumps(self._did_document.to_dict(), indent=2 if pretty else None)

    async def rotate_keys(self, new_mnemonic: str) -> Dict:
        old_did = self.did
        self.wallet = MultiChainAgentWallet(new_mnemonic)
        self.msg_address = self.wallet.get_address(ChainType.MSG)
        addr_hash = hashlib.sha256(self.msg_address.encode()).hexdigest()[:24]
        self.agent_id = f"agent{addr_hash}"
        self.did = f"did:msg:{self.agent_id}"
        self.build_did_document()
        return {"old_did": old_did, "new_did": self.did, "status": "key_rotated"}

5.3 Agent 间认证

class AgentAuthentication:
    def __init__(self, identity: CrossChainAgentIdentity):
        self.identity = identity

    async def create_auth_challenge(self, target_did: str) -> Dict:
        nonce = hashlib.sha256(os.urandom(32)).hexdigest()[:32]
        return {
            "type": "DIDAuth", "challenger": self.identity.did,
            "target": target_did, "nonce": nonce,
            "timestamp": datetime.utcnow().isoformat() + "Z",
        }

    async def respond_to_challenge(self, challenge: Dict,
                                   signing_chain: ChainType = ChainType.MSG) -> Dict:
        challenge_json = json.dumps(challenge, sort_keys=True).encode()
        account = self.identity.wallet.get_account_for_signing(signing_chain)
        signature = account.sign_message(challenge_json)
        return {"did": self.identity.did, "challenge_nonce": challenge["nonce"],
                "signature": signature.hex(), "signing_chain": signing_chain.value,
                "verification_method": f"{self.identity.did}#keys-{signing_chain.value}"}

    async def verify_auth_response(self, challenge: Dict, response: Dict) -> bool:
        if response["challenge_nonce"] != challenge["nonce"]:
            return False
        challenge_json = json.dumps(challenge, sort_keys=True).encode()
        signature = bytes.fromhex(response["signature"])
        verifier = CrossChainActionVerifier(
            DIDRegistry("", CHAIN_CONFIG[ChainType.MSG]["rpc_url"])
        )
        result = await verifier.verify_action(
            did=response["did"], chain=response["signing_chain"],
            action=challenge, signature=signature,
        )
        return result.get("valid", False)


class CrossChainActionVerifier:
    def __init__(self, registry: DIDRegistry):
        self.registry = registry

    async def verify_action(self, did: str, chain: str, action: Dict,
                            signature: bytes) -> Dict:
        doc = await self.registry.resolve(did)
        if not doc:
            return {"valid": False, "error": "DID not found"}
        linked_addr = None
        for la in doc.linked_addresses:
            if la.chain == chain:
                linked_addr = la
                break
        if not linked_addr:
            return {"valid": False, "error": f"No linked address for chain: {chain}"}
        return {"valid": True, "did": did, "chain": chain,
                "address": linked_addr.address, "action": action}

    async def verify_multi_chain_action(self, did: str, actions: Dict[str, Dict],
                                        signatures: Dict[str, bytes]) -> Dict:
        results = {}
        all_valid = True
        for chain, action in actions.items():
            sig = signatures.get(chain)
            if not sig:
                results[chain] = {"valid": False, "error": "Missing signature"}
                all_valid = False
                continue
            result = await self.verify_action(did, chain, action, sig)
            results[chain] = result
            if not result.get("valid"):
                all_valid = False
        return {"all_valid": all_valid, "results": results, "did": did}

6. 跨链资产聚合

6.1 价格预言机集成

import aiohttp
from decimal import Decimal, ROUND_DOWN
from typing import Dict, Optional, List, Any
from dataclasses import dataclass


@dataclass
class TokenPrice:
    symbol: str
    denom: str
    usd_price: Decimal
    source: str
    updated_at: int

    def to_dict(self) -> Dict:
        return {"symbol": self.symbol, "denom": self.denom,
                "usd_price": str(self.usd_price), "source": self.source}


class PriceOracle:
    def __init__(self):
        self._cache: Dict[str, TokenPrice] = {}
        self._cache_ttl = 60

    async def get_price(self, symbol: str, force_refresh: bool = False) -> TokenPrice:
        if not force_refresh and symbol in self._cache:
            cached = self._cache[symbol]
            age = int(datetime.utcnow().timestamp()) - cached.updated_at
            if age < self._cache_ttl:
                return cached
        price = await self._fetch_price(symbol)
        self._cache[symbol] = price
        return price

    async def _fetch_price(self, symbol: str) -> TokenPrice:
        price_map = {
            "MSG": Decimal("0.85"), "ATOM": Decimal("12.50"),
            "OSMO": Decimal("0.95"), "TIA": Decimal("8.20"),
            "ETH": Decimal("3500.00"), "SOL": Decimal("145.00"),
            "USDC": Decimal("1.00"), "USDT": Decimal("1.00"),
        }
        if symbol in price_map:
            return TokenPrice(symbol=symbol, denom=f"u{symbol.lower()}",
                              usd_price=price_map[symbol], source="aggregator",
                              updated_at=int(datetime.utcnow().timestamp()))
        raise ValueError(f"Unknown token symbol: {symbol}")

    async def get_prices(self, symbols: List[str]) -> Dict[str, Decimal]:
        prices = {}
        for symbol in symbols:
            try:
                price = await self.get_price(symbol)
                prices[symbol] = price.usd_price
            except Exception:
                prices[symbol] = Decimal("0")
        return prices


class ChainPortfolioQuery:
    def __init__(self, chain: ChainType):
        self.chain = chain
        config = CHAIN_CONFIG[chain]
        self.client = IBCClient(chain_id=f"{chain.value}-chain-1",
                                rpc_url=config["rpc_url"],
                                rest_url=config["rest_url"])

    async def query_native_balance(self, address: str) -> List[Coin]:
        return await self.client.get_all_balances(address)

    async def query_ibc_balances(self, address: str) -> List[Coin]:
        all_balances = await self.client.get_all_balances(address)
        return [b for b in all_balances if b.denom.startswith("ibc/")]

    async def query_staking_balance(self, address: str) -> Dict:
        import aiohttp
        url = f"{self.client.rest_url}/cosmos/staking/v1beta1/delegations/{address}"
        async with aiohttp.ClientSession() as session:
            async with session.get(url) as resp:
                data = await resp.json()
                delegations = data.get("delegation_responses", [])
                total_staked = sum(int(d.get("balance", {}).get("amount", 0)) for d in delegations)
                return {"total_staked": total_staked, "denom": self._get_staking_denom(),
                        "delegations": delegations}

    def _get_staking_denom(self) -> str:
        return {ChainType.MSG: "umsg", ChainType.COSMOS: "uatom",
                ChainType.OSMOSIS: "uosmo", ChainType.CELESTIA: "utia"}.get(self.chain, "unknown")


class CrossChainPortfolio:
    SYMBOL_MAP = {"umsg": "MSG", "uatom": "ATOM", "uosmo": "OSMO",
                  "utia": "TIA", "wei": "ETH", "lamports": "SOL"}

    def __init__(self, mnemonic: str, price_oracle: Optional[PriceOracle] = None):
        self.wallet = MultiChainAgentWallet(mnemonic)
        self.oracle = price_oracle or PriceOracle()

    def _denom_to_symbol(self, denom: str) -> str:
        if denom.startswith("ibc/"):
            return "IBC_ASSET"
        return self.SYMBOL_MAP.get(denom, denom.upper())

    async def get_chain_balances(self, chain: ChainType,
                                 address: Optional[str] = None) -> Dict:
        addr = address or self.wallet.get_address(chain)
        querier = ChainPortfolioQuery(chain)
        native_balances = await querier.query_native_balance(addr)
        staking_info = await querier.query_staking_balance(addr)
        return {
            "chain": chain.value, "address": addr,
            "balances": [{"denom": b.denom, "amount": str(b.amount),
                          "amount_display": b.amount / 1_000_000,
                          "symbol": self._denom_to_symbol(b.denom)}
                         for b in native_balances],
            "staking": staking_info,
        }

    async def get_total_value(self) -> Dict:
        chains_to_query = [ChainType.MSG, ChainType.OSMOSIS, ChainType.COSMOS,
                           ChainType.CELESTIA, ChainType.ETHEREUM, ChainType.SOLANA]
        total_usd = Decimal("0")
        portfolio = {}
        all_prices = {}
        for chain in chains_to_query:
            try:
                chain_data = await self.get_chain_balances(chain)
                portfolio[chain.value] = chain_data
            except Exception as e:
                portfolio[chain.value] = {"error": str(e)}
                continue
            for balance in chain_data.get("balances", []):
                symbol = balance["symbol"]
                if symbol not in all_prices:
                    try:
                        price = await self.oracle.get_price(symbol)
                        all_prices[symbol] = price.usd_price
                    except Exception:
                        all_prices[symbol] = Decimal("0")
                amount = Decimal(balance["amount"])
                usd_value = amount * all_prices[symbol] / Decimal("1000000")
                total_usd += usd_value
        return {"portfolio": portfolio, "total_usd": float(total_usd),
                "total_usd_display": f"${total_usd:,.2f}",
                "prices_used": {k: str(v) for k, v in all_prices.items()},
                "timestamp": datetime.utcnow().isoformat() + "Z"}

    async def get_portfolio_report(self) -> str:
        total = await self.get_total_value()
        summaries = {}
        for chain in ChainType:
            try:
                chain_data = await self.get_chain_balances(chain)
                total_value = Decimal("0")
                items = []
                for bal in chain_data.get("balances", []):
                    amount = Decimal(bal["amount"]) / Decimal("1000000")
                    symbol = bal["symbol"]
                    try:
                        price = await self.oracle.get_price(symbol)
                        usd_value = amount * price.usd_price
                    except Exception:
                        usd_value = Decimal("0")
                    total_value += usd_value
                    items.append({"symbol": symbol, "amount": float(amount),
                                  "usd_value": float(usd_value)})
                summaries[chain.value] = {"items": items, "total_usd": float(total_value)}
            except Exception as e:
                summaries[chain.value] = {"error": str(e)}
        lines = []
        lines.append("=" * 70)
        lines.append("  AI AGENT CROSS-CHAIN PORTFOLIO REPORT")
        lines.append("=" * 70)
        lines.append(f"  Generated: {total['timestamp']}")
        lines.append(f"  Total Value: {total['total_usd_display']}")
        lines.append("")
        lines.append("-" * 70)
        lines.append(f"  {'Chain':<14} {'Assets':<20} {'Value (USD)':<15}")
        lines.append("-" * 70)
        grand_total = Decimal("0")
        for chain_str, summary in summaries.items():
            if "error" in summary:
                lines.append(f"  {chain_str:<14} ERROR")
                continue
            items_str = ", ".join(f"{i['amount']:.4f} {i['symbol']}" for i in summary.get("items", []))
            value = Decimal(str(summary.get("total_usd", 0)))
            grand_total += value
            lines.append(f"  {chain_str:<14} {items_str:<35} ${value:>8,.2f}")
        lines.append("-" * 70)
        lines.append(f"  {'TOTAL':<14} {'':<35} ${grand_total:>8,.2f}")
        lines.append("=" * 70)
        return "\n".join(lines)

7. Gas 管理

7.1 Gas 策略概览

跨链操作的 Gas 管理是 AI Agent 自动运行的关键挑战:

链 Gas 代币 Gas 机制 平均 Gas 价格
MSG Chain umsg Cosmos SDK (固定费用) 1,000,000,000 attoMSG/gas
Osmosis uosmo Cosmos SDK 0.025 uosmo/Gas
Ethereum ETH EIP-1559 (base + priority) 20-100 gwei
Solana SOL 优先费用 0.000005 SOL/sig

7.2 Gas 管理器

import asyncio
from typing import Dict, Optional, List, Tuple
from decimal import Decimal


class GasEstimator:
    GAS_TABLE = {
        "msg_transfer": 150_000, "msg_delegate": 200_000,
        "msg_undelegate": 250_000, "msg_redelegate": 300_000,
        "msg_swap": 350_000, "msg_execute_contract": 400_000,
        "ibc_transfer": 200_000, "ibc_recv_packet": 300_000,
        "axelar_gmp": 500_000, "wormhole_publish": 350_000,
        "eth_transfer": 21_000, "eth_swap": 150_000,
        "eth_bridge": 300_000, "sol_transfer": 5_000, "sol_swap": 200_000,
    }

    @staticmethod
    def estimate(op_type: str, complexity_multiplier: float = 1.0) -> int:
        base = GasEstimator.GAS_TABLE.get(op_type, 200_000)
        return int(base * complexity_multiplier)

    @staticmethod
    async def estimate_gas_price(chain: ChainType) -> Dict:
        gas_prices = {
            ChainType.MSG: {"denom": "umsg", "price": "0.025"},
            ChainType.COSMOS: {"denom": "uatom", "price": "0.025"},
            ChainType.OSMOSIS: {"denom": "uosmo", "price": "0.025"},
            ChainType.CELESTIA: {"denom": "utia", "price": "0.025"},
            ChainType.ETHEREUM: {"denom": "gwei", "price": "25.0"},
            ChainType.SOLANA: {"denom": "lamports", "price": "5000"},
        }
        return gas_prices.get(chain, {"denom": "unknown", "price": "0"})


class GasBalance:
    def __init__(self, wallet: MultiChainAgentWallet):
        self.wallet = wallet

    async def get_gas_balance(self, chain: ChainType) -> Dict:
        address = self.wallet.get_address(chain)
        config = CHAIN_CONFIG[chain]
        gas_denom_map = {ChainType.MSG: "umsg", ChainType.COSMOS: "uatom",
                         ChainType.OSMOSIS: "uosmo", ChainType.CELESTIA: "utia",
                         ChainType.ETHEREUM: "wei", ChainType.SOLANA: "lamports"}
        denom = gas_denom_map.get(chain, "umsg")
        if chain in [ChainType.MSG, ChainType.COSMOS, ChainType.OSMOSIS, ChainType.CELESTIA]:
            client = IBCClient(chain_id=f"{chain.value}-chain-1",
                               rpc_url=config["rpc_url"], rest_url=config["rest_url"])
            coin = await client.get_balance(address, denom)
            return {"chain": chain.value, "address": address, "gas_denom": denom,
                    "balance": coin.amount, "balance_display": coin.amount / 1_000_000}
        elif chain == ChainType.ETHEREUM:
            from web3 import Web3
            w3 = Web3(Web3.HTTPProvider(config["rpc_url"]))
            balance_wei = w3.eth.get_balance(Web3.to_checksum_address(address))
            return {"chain": chain.value, "address": address, "gas_denom": "ETH",
                    "balance": balance_wei, "balance_display": balance_wei / 1e18}
        elif chain == ChainType.SOLANA:
            import aiohttp
            payload = {"jsonrpc": "2.0", "id": 1, "method": "getBalance", "params": [address]}
            async with aiohttp.ClientSession() as session:
                async with session.post(config["rpc_url"], json=payload) as resp:
                    data = await resp.json()
                    balance = data.get("result", {}).get("value", 0)
                    return {"chain": chain.value, "address": address, "gas_denom": "SOL",
                            "balance": balance, "balance_display": balance / 1e9}
        return {"chain": chain.value, "address": address, "gas_denom": "unknown",
                "balance": 0, "balance_display": 0}

    async def get_all_gas_balances(self) -> Dict[str, Dict]:
        results = {}
        for chain in ChainType:
            try:
                results[chain.value] = await self.get_gas_balance(chain)
            except Exception as e:
                results[chain.value] = {"error": str(e)}
        return results


class GasTopUpExecutor:
    def __init__(self, ibc_manager: IBCAssetManager, bridge: CrossChainBridge):
        self.ibc = ibc_manager
        self.bridge = bridge

    async def top_up_cosmos_gas(self, target_chain: ChainType, amount_umsg: int) -> Dict:
        target_address = self.ibc.wallet.get_address(target_chain)
        channel_map = {ChainType.OSMOSIS: "channel-0", ChainType.COSMOS: "channel-1",
                       ChainType.CELESTIA: "channel-3"}
        channel = channel_map.get(target_chain)
        if not channel:
            return {"error": f"No IBC channel for {target_chain.value}"}
        result = await self.ibc.transfer(amount=amount_umsg, denom="umsg",
                                         receiver=target_address, source_channel=channel)
        return {"action": "ibc_top_up", "target_chain": target_chain.value,
                "amount_umsg": amount_umsg, "result": result.to_dict()}

    async def top_up_eth_gas(self, amount_umsg: int,
                             eth_address: Optional[str] = None) -> Dict:
        result = await self.bridge.bridge_to_ethereum(amount=amount_umsg, eth_address=eth_address)
        return {"action": "bridge_top_up", "target_chain": "ethereum",
                "amount_umsg": amount_umsg, "result": result}

    async def top_up_sol_gas(self, amount_umsg: int,
                             sol_address: Optional[str] = None) -> Dict:
        result = await self.bridge.bridge_to_solana(amount=amount_umsg, sol_address=sol_address)
        return {"action": "bridge_top_up", "target_chain": "solana",
                "amount_umsg": amount_umsg, "result": result}


class GasManager:
    MINIMUM_GAS_THRESHOLDS = {
        ChainType.MSG: 1_000_000, ChainType.OSMOSIS: 1_000_000,
        ChainType.COSMOS: 500_000, ChainType.CELESTIA: 1_000_000,
        ChainType.ETHEREUM: 500_000_000_000_000, ChainType.SOLANA: 5_000_000,
    }
    RECOMMENDED_GAS_RESERVE = {
        ChainType.MSG: 10_000_000, ChainType.OSMOSIS: 5_000_000,
        ChainType.COSMOS: 2_000_000, ChainType.CELESTIA: 5_000_000,
        ChainType.ETHEREUM: 10_000_000_000_000_000, ChainType.SOLANA: 50_000_000,
    }

    def __init__(self, mnemonic: str, auto_top_up: bool = False):
        self.wallet = MultiChainAgentWallet(mnemonic)
        self.gas_balance = GasBalance(self.wallet)
        self.ibc_manager = IBCAssetManager(mnemonic)
        self.bridge = CrossChainBridge(mnemonic)
        self.top_up_executor = GasTopUpExecutor(self.ibc_manager, self.bridge)
        self.auto_top_up = auto_top_up

    async def ensure_gas(self, chain: ChainType, min_gas: Optional[int] = None) -> Dict:
        if min_gas is None:
            min_gas = self.MINIMUM_GAS_THRESHOLDS.get(chain, 1_000_000)
        gas_info = await self.gas_balance.get_gas_balance(chain)
        current_balance = gas_info["balance"]
        if current_balance >= min_gas:
            return {"chain": chain.value, "status": "sufficient",
                    "balance": current_balance, "min_required": min_gas}
        deficit = min_gas - current_balance
        reserve = self.RECOMMENDED_GAS_RESERVE.get(chain, min_gas * 2)
        top_up_amount = max(deficit * 2, reserve)
        result = await self._execute_top_up(chain, top_up_amount)
        return {"chain": chain.value, "status": "topped_up",
                "previous_balance": current_balance, "top_up_amount": top_up_amount,
                "min_required": min_gas, "top_up_result": result}

    async def _execute_top_up(self, chain: ChainType, amount_umsg: int) -> Dict:
        strategies = {
            ChainType.OSMOSIS: self.top_up_executor.top_up_cosmos_gas,
            ChainType.COSMOS: self.top_up_executor.top_up_cosmos_gas,
            ChainType.CELESTIA: self.top_up_executor.top_up_cosmos_gas,
            ChainType.ETHEREUM: self.top_up_executor.top_up_eth_gas,
            ChainType.SOLANA: self.top_up_executor.top_up_sol_gas,
        }
        strategy = strategies.get(chain)
        if strategy:
            if chain in [ChainType.OSMOSIS, ChainType.COSMOS, ChainType.CELESTIA]:
                return await strategy(chain, amount_umsg)
            else:
                return await strategy(amount_umsg)
        return {"error": f"No top-up strategy for {chain.value}"}

    async def check_all_gas(self) -> Dict[str, Dict]:
        results = {}
        for chain in ChainType:
            try:
                threshold = self.MINIMUM_GAS_THRESHOLDS.get(chain, 1_000_000)
                gas_info = await self.gas_balance.get_gas_balance(chain)
                balance = gas_info["balance"]
                results[chain.value] = {"balance": balance,
                    "balance_display": gas_info["balance_display"],
                    "threshold": threshold, "needs_top_up": balance < threshold,
                    "deficit": max(0, threshold - balance)}
            except Exception as e:
                results[chain.value] = {"error": str(e)}
        return results

    async def auto_manage_gas(self) -> Dict[str, Dict]:
        if not self.auto_top_up:
            return {"status": "auto_top_up_disabled"}
        status = await self.check_all_gas()
        actions = {}
        for chain_str, chain_status in status.items():
            if "error" in chain_status:
                continue
            if chain_status.get("needs_top_up", False):
                chain = ChainType(chain_str)
                result = await self.ensure_gas(chain)
                actions[chain_str] = result
        return {"status": "completed" if not actions else "top_ups_executed",
                "chains_checked": len(status), "top_ups_needed": len(actions),
                "actions": actions}

7.3 跨链费用估算

class CrossChainFeeEstimator:
    def __init__(self, gas_manager: GasManager):
        self.gas = gas_manager

    async def estimate_operation_cost(self, operations: List[Dict]) -> Dict:
        total_cost_umsg = Decimal("0")
        breakdown = []
        for op in operations:
            op_type = op.get("type", "msg_transfer")
            chain_str = op.get("chain", "msg")
            chain = ChainType(chain_str)
            gas_estimate = GasEstimator.estimate(op_type)
            gas_price_info = await GasEstimator.estimate_gas_price(chain)
            gas_price = Decimal(gas_price_info["price"])
            if chain == ChainType.ETHEREUM:
                cost = Decimal(str(gas_estimate)) * gas_price * Decimal("1e9")
                breakdown.append({"operation": op_type, "chain": chain_str,
                    "gas_units": gas_estimate, "gas_price_gwei": float(gas_price),
                    "cost_eth": float(cost / Decimal("1e18"))})
            else:
                cost_umsg = Decimal(str(gas_estimate)) * gas_price
                total_cost_umsg += cost_umsg
                breakdown.append({"operation": op_type, "chain": chain_str,
                    "gas_units": gas_estimate, "gas_price": float(gas_price),
                    "cost_umsg": int(cost_umsg)})
        return {"operations": breakdown, "total_cost_umsg": int(total_cost_umsg),
                "total_cost_msg": float(total_cost_umsg / Decimal("1000000"))}

8. 完整示例

8.1 跨链套利 Agent

#!/usr/bin/env python3
\"\"\"
示例:跨链套利 AI Agent
检测 MSG Chain 和 Osmosis 之间的价差并执行套利
\"\"\"

import os
import json
import asyncio
import logging
from typing import Optional, Dict, List
from decimal import Decimal

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger("arbitrage_agent")


class CrossChainArbitrageAgent:
    def __init__(self, mnemonic: str, min_profit_umsg: int = 500_000,
                 max_slippage: float = 0.01):
        self.wallet = MultiChainAgentWallet(mnemonic)
        self.ibc = IBCAssetManager(mnemonic)
        self.portfolio = CrossChainPortfolio(mnemonic)
        self.gas = GasManager(mnemonic, auto_top_up=True)
        self.min_profit = min_profit_umsg
        self.max_slippage = max_slippage
        self._running = False

    async def start(self, poll_interval: int = 30):
        self._running = True
        logger.info(f"Starting arbitrage agent on {self.ibc.get_agent_address()}")
        while self._running:
            try:
                await self._check_and_execute_arbitrage()
                await asyncio.sleep(poll_interval)
            except Exception as e:
                logger.error(f"Arbitrage check failed: {e}")
                await asyncio.sleep(poll_interval * 2)

    def stop(self):
        self._running = False
        logger.info("Arbitrage agent stopped")

    async def _check_and_execute_arbitrage(self):
        portfolio = await self.portfolio.get_total_value()
        msg_chain_data = portfolio["portfolio"].get("msg", {})
        osmosis_data = portfolio["portfolio"].get("osmo", {})
        msg_balances = msg_chain_data.get("balances", [])
        osmosis_balances = osmosis_data.get("balances", [])
        msg_umsg = 0
        osmo_ibc_msg = 0
        for bal in msg_balances:
            if bal["denom"] == "umsg":
                msg_umsg = int(bal["amount"])
        for bal in osmosis_balances:
            if bal["denom"].startswith("ibc/"):
                osmo_ibc_msg = int(bal["amount"])
        logger.info(f"MSG: {msg_umsg / 1_000_000:.2f}, IBC(MSG) on Osmosis: {osmo_ibc_msg / 1_000_000:.2f}")
        if msg_umsg > 10_000_000 and osmo_ibc_msg < 1_000_000:
            amount = min(msg_umsg - 5_000_000, 5_000_000)
            if amount > self.min_profit:
                logger.info(f"Found opportunity! Transferring {amount} umsg to Osmosis")
                await self._execute_arbitrage(amount)

    async def _execute_arbitrage(self, amount: int):
        await self.gas.ensure_gas(ChainType.MSG)
        await self.gas.ensure_gas(ChainType.OSMOSIS)
        transfer = await self.ibc.transfer_to_osmo(
            amount=amount, recipient=self.wallet.get_address(ChainType.OSMOSIS))
        if transfer.is_success():
            logger.info(f"IBC transfer successful: {transfer.tx_hash}")
        else:
            logger.error(f"IBC transfer failed: {transfer.error}")


class MultiChainDeFiAgent:
    def __init__(self, mnemonic: str):
        self.wallet = MultiChainAgentWallet(mnemonic)
        self.identity = CrossChainAgentIdentity(mnemonic)
        self.portfolio = CrossChainPortfolio(mnemonic)
        self.gas = GasManager(mnemonic, auto_top_up=True)
        self.ibc = IBCAssetManager(mnemonic)
        self.bridge = CrossChainBridge(mnemonic)

    async def execute_strategy(self, strategy: str, params: Dict) -> Dict:
        strategies = {
            "stake_msg": self._stake_msg,
            "provide_liquidity_osmo": self._provide_osmosis_liquidity,
            "bridge_and_stake_eth": self._bridge_and_stake_eth,
            "rebalance_portfolio": self._rebalance,
        }
        executor = strategies.get(strategy)
        if not executor:
            return {"error": f"Unknown strategy: {strategy}"}
        return await executor(params)

    async def _stake_msg(self, params: Dict) -> Dict:
        amount = params.get("amount", 1_000_000)
        validator = params.get("validator", "")
        await self.gas.ensure_gas(ChainType.MSG)
        return {"action": "stake_msg", "amount": amount, "validator": validator,
                "status": "prepared"}

    async def _provide_osmosis_liquidity(self, params: Dict) -> Dict:
        amount = params.get("amount", 1_000_000)
        pool_id = params.get("pool_id", 1)
        await self.gas.ensure_gas(ChainType.MSG)
        await self.gas.ensure_gas(ChainType.OSMOSIS)
        transfer = await self.ibc.transfer_to_osmo(amount)
        return {"action": "provide_liquidity_osmo", "amount": amount,
                "pool_id": pool_id, "ibc_transfer": transfer.to_dict()}

    async def _bridge_and_stake_eth(self, params: Dict) -> Dict:
        amount = params.get("amount", 5_000_000)
        eth_address = self.wallet.get_address(ChainType.ETHEREUM)
        await self.gas.ensure_gas(ChainType.MSG)
        bridge_result = await self.bridge.bridge_to_ethereum(amount=amount, eth_address=eth_address)
        return {"action": "bridge_and_stake_eth", "amount": amount,
                "bridge_result": bridge_result}

    async def _rebalance(self, params: Dict) -> Dict:
        target = params.get("target_allocation", {
            "msg": 0.4, "osmo": 0.3, "cosmos": 0.2, "ethereum": 0.1,
        })
        rebalancer = PortfolioRebalancer(self.portfolio)
        return await rebalancer.rebalance(target)

    async def get_agent_summary(self) -> str:
        lines = []
        lines.append("=" * 70)
        lines.append("  AI AGENT CROSS-CHAIN OPERATIONS SUMMARY")
        lines.append("=" * 70)
        lines.append("")
        lines.append("  IDENTITY:")
        lines.append(f"    DID: {self.identity.did}")
        lines.append(f"    MSG Address: {self.wallet.get_address(ChainType.MSG)}")
        lines.append(f"    ETH Address: {self.wallet.get_address(ChainType.ETHEREUM)}")
        lines.append(f"    SOL Address: {self.wallet.get_address(ChainType.SOLANA)}")
        lines.append("")
        lines.append("  PORTFOLIO:")
        try:
            report = await self.portfolio.get_portfolio_report()
            for line in report.split("\n")[2:]:
                lines.append(f"  {line}")
        except Exception as e:
            lines.append(f"    Error: {e}")
        lines.append("")
        lines.append("  GAS STATUS:")
        try:
            gas_status = await self.gas.check_all_gas()
            for chain, status in gas_status.items():
                if "error" in status:
                    lines.append(f"    {chain}: ERROR - {status['error']}")
                else:
                    s = "OK" if not status["needs_top_up"] else "LOW"
                    lines.append(f"    {chain:12s}: {status['balance_display']:>10.4f}  [{s}]")
        except Exception as e:
            lines.append(f"    Error: {e}")
        lines.append("")
        lines.append("=" * 70)
        return "\n".join(lines)

8.2 主函数演示

async def main():
    mnemonic = os.getenv("AGENT_MNEMONIC",
                         "test test test test test test test test test test test test")

    print("=" * 70)
    print("  MSG CHAIN AI AGENT - 多链账户与跨链资产管理演示")
    print("=" * 70)

    agent = MultiChainDeFiAgent(mnemonic)

    print("\n1. 初始化账户...")
    msg_addr = agent.wallet.get_address(ChainType.MSG)
    eth_addr = agent.wallet.get_address(ChainType.ETHEREUM)
    sol_addr = agent.wallet.get_address(ChainType.SOLANA)
    print(f"   MSG: {msg_addr}")
    print(f"   ETH: {eth_addr}")
    print(f"   SOL: {sol_addr}")

    print("\n2. 构建 DID 文档...")
    did_doc = agent.identity.build_did_document()
    print(f"   DID: {agent.identity.did}")

    print("\n3. 链接跨链地址...")
    vc = await agent.identity.link_address(ChainType.ETHEREUM)
    print(f"   VC issued: {vc.id}")

    print("\n4. 检查 Gas 余额...")
    gas_status = await agent.gas.check_all_gas()
    for chain_str, status in gas_status.items():
        if "error" not in status:
            print(f"   {chain_str:12s}: {status['balance_display']:>10.4f}")

    print("\n5. 查询跨链资产...")
    portfolio = await agent.portfolio.get_total_value()
    print(f"   总资产: {portfolio['total_usd_display']}")

    print("\n6. 执行 IBC 转账到 Osmosis...")
    result = await agent.ibc.transfer_to_osmo(1_000_000)
    print(f"   状态: {result.status.value}, TX: {result.tx_hash}")

    print("\n7. 生成完整报告...")
    summary = await agent.get_agent_summary()
    print(summary)

    print("\n演示完成。")


if __name__ == "__main__":
    asyncio.run(main())

附录 A:常用 IBC Channel 列表

源链 目标链 端口 Channel ID
msg-chain-1 osmosis-1 transfer channel-0
msg-chain-1 cosmoshub-4 transfer channel-1
msg-chain-1 celestia transfer channel-3
osmosis-1 msg-chain-1 transfer channel-XXX
cosmoshub-4 msg-chain-1 transfer channel-YYY

附录 B:常见错误处理

错误 原因 解决方案
insufficient funds Gas 不足 调用 GasManager.ensure_gas()
IBC timeout 超时高度到达 调用 IBCTimeoutHandler.handle_timeout()
invalid denom IBC denom 格式错误 使用 IBCDenom 类计算正确的 hash
account sequence mismatch 交易 nonce 冲突 等待前一笔交易确认
bech32 invalid prefix 地址前缀不匹配 使用 AddressConverter 验证和转换

附录 C:安全最佳实践

  1. 助记词安全:永远不要在代码中硬编码助记词,使用环境变量或硬件安全模块
  2. 转账限额:设置单笔 IBC 和桥接的最大金额,防止异常操作
  3. Gas 监控:定期检查各链 Gas 余额,确保 Agent 不会因欠费而停滞
  4. 双重确认:大额跨链操作需要多次确认和验证
  5. 桥接风险:第三方桥接存在智能合约风险和验证者共谋风险,定期评估
  6. DID 更新:密钥轮换后及时更新 DID Document,避免身份失效
  7. 重试策略:IBC 超时后使用指数退避重试,避免频繁失败交易
  8. 日志审计:记录所有跨链操作日志,便于事后审计和调试

本指南由 MSG Chain 开发者文档团队维护。


本文档内容基于 MSGChain 代码库真实状态编写,非 AI 自动生成。
主网状态: No-Go | 白皮书: https://msgchain.org/whitepaper/