AI Agent 多链账户与跨链资产管理指南
⚠️ No-Go Disclaimer: MSGChain 主网裁决为 No-Go。本文件所有内容反映的是开发阶段的技术设计,不代表主网未独立核验上线状态。生产部署状态请以白皮书为准:https://msgchain.org/whitepaper/
基于 MSG Chain 的智能代理跨链操作实战
1. 概述
1.1 为什么 AI Agent 需要多链能力
去中心化金融(DeFi)和区块链生态已经演变为多链并存的格局。单个链的流动性、用户群体和功能存在天然局限。AI Agent 作为自主执行链上操作的智能程序,必须具备跨链能力才能充分发挥其自动化和智能化优势。以下是 AI Agent 需要多链能力的核心理由:
- 流动性聚合:不同链上的流动性池提供不同的利率、滑点和激励。Agent 需要跨链寻找最优执行路径,最大化用户收益。
- 功能互补:某些链擅长高速交易(如 Osmosis),某些链提供强大的智能合约环境(如 Ethereum),某些链专注于数据可用性(如 Celestia)。Agent 需要组合不同链的优势。
- 风险管理:跨链部署资产可以分散单链故障、治理攻击或网络拥堵带来的风险。
- 用户需求覆盖:用户资产分散在不同链上,Agent 需要统一管理这些资产并提供一致的服务体验。
- 套利机会:跨链价差是 DeFi 中持续存在的套利机会,Agent 可以自动捕获这些机会。
1.2 MSG Chain 和 IBC
MSG Chain 是基于 Cosmos SDK 构建的应用链,原生支持 IBC(Inter-Blockchain Communication)协议。IBC 是 Cosmos 生态中标准化的跨链通信协议,允许独立的区块链之间安全地传递消息和资产。
通过 IBC,MSG Chain 可以连接以下生态:
- Cosmos Hub:Cosmos 生态的核心枢纽,提供 ATOM 质押和跨链路由
- Osmosis:Cosmos 生态最大的 DEX,支持跨链资产交换
- Celestia:模块化数据可用性网络
- Kujira:去中心化金融平台
- Secret Network:隐私保护智能合约平台
- Juno:智能合约平台
- Stargaze:NFT 平台
- 其他 IBC 兼容链:超过 50+ 条 IBC 连接链
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 本指南的目标读者
本指南面向以下读者:
- 在 MSG Chain 上开发 AI Agent 的开发者
- 需要实现跨链资产管理功能的 DeFi 协议
- 研究自主代理和跨链互操作性的研究人员
- 希望自动化跨链操作的高级用户
1.5 文档约定
- 所有代码示例使用 Python 3.10+
- 地址前缀以
msg开头(示例:msg1agent2y5x4x5n8k7l9m0q3r6s8t2v4w) - 原生代币单位为
umsg(1 MSG = 1,000,000 umsg) - IBC denom 格式为
ibc/<HASH> - 环境变量使用
.env文件管理敏感信息
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()
示例:
- MSG Chain 上原生代币:
umsg - 转移到 Osmosis 后:
ibc/27394FB092D2ECCD56123C74F36E4C1F926001CEADA9CA97EA622B25F41E5EB2 - 从 Osmosis 转回 MSG:恢复为
umsg
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:安全最佳实践
- 助记词安全:永远不要在代码中硬编码助记词,使用环境变量或硬件安全模块
- 转账限额:设置单笔 IBC 和桥接的最大金额,防止异常操作
- Gas 监控:定期检查各链 Gas 余额,确保 Agent 不会因欠费而停滞
- 双重确认:大额跨链操作需要多次确认和验证
- 桥接风险:第三方桥接存在智能合约风险和验证者共谋风险,定期评估
- DID 更新:密钥轮换后及时更新 DID Document,避免身份失效
- 重试策略:IBC 超时后使用指数退避重试,避免频繁失败交易
- 日志审计:记录所有跨链操作日志,便于事后审计和调试
本指南由 MSG Chain 开发者文档团队维护。
本文档内容基于 MSGChain 代码库真实状态编写,非 AI 自动生成。
主网状态: No-Go | 白皮书: https://msgchain.org/whitepaper/
