MSG Chain Rust SDK 开发指南
数据来源:MSG Chain 代码库核实
主网状态: No-Go — 当前 MSGChain 主网裁决为 No-Go,以下内容反映代码实际状态,不代表生产可用。
版本: 1.0.0
适用链: MSG Chain (msg-chain-1)
共识: Round-Robin + DAR
签名: Dilithium-5(后量子密码学)
智能合约: CosmWasm 1.5
定位: 类似 ethers-rs 的 Rust 全栈 SDK
目录
1. 概述
1.1 为什么需要 Rust SDK
MSG Chain 基于 Cosmos SDK + CosmWasm 构建,核心特性:
- 后量子签名: Dilithium-5,传统 ECDSA/Ed25519 SDK 无法直接适配
- AI 原生智能合约: 5 个专用 AI 合约(模型注册、推理证明、数据资产、Agent 注册、推理市场)
- Agent API: 18 个 REST + WebSocket 端点
- CosmWasm 生态: 46 个合约包
Rust SDK 的价值远超合约开发:
| 场景 | 传统方式 | Rust SDK |
|---|---|---|
| 链上数据索引 | 手动 RPC 调用 | 类型安全的查询客户端 |
| AI Agent 后端 | Python + 非结构化 HTTP | 强类型 Agent API 客户端 |
| 批量交易 | shell 脚本 + msg-cli | 异步并发 TxBuilder |
| 后量子签名 | 无原生支持 | Dilithium-5 实现 |
| 实时事件 | WebSocket 手写 | 类型安全的事件订阅 |
| 测试 | 无 | Mock 服务器 + 集成测试 |
1.2 适用人群
- MSG Chain 合约开发者
- AI Agent 后端工程师
- 需 Rust 替代 Python/Go 做链交互的团队
- 后量子密码学探索者
1.3 依赖关系图
msg-chain-sdk/
├── core/ # 核心客户端、错误处理、配置
├── types/ # 地址、交易、区块等类型
├── crypto/ # Dilithium-5, Ed25519, BIP39/44
├── clients/
│ ├── bank/ # 银行查询
│ ├── staking/ # 质押查询
│ ├── contract/ # 合约查询与执行
│ ├── gov/ # 治理
│ └── agent/ # Agent API (REST + WebSocket)
├── tx/ # 交易构建、签名、广播
└── testing/ # Mock 服务器与测试工具
1.4 链参数速查
| 参数 | 值 |
|---|---|
| Chain ID | msg-chain-1 |
| Bech32 前缀 | msg |
| CoinType | 118 |
| MSG 精度 | 18 位小数 |
| Gas 价格 | 1,000,000,000 attoMSG/gas |
| 区块时间 | ~5s |
| 共识 | Round-Robin + DAR |
| 签名算法 | Dilithium-5 |
2. 安装配置
2.1 系统要求
- Rust 1.75+(稳定版)
wasm32-unknown-unknown目标- protoc(可选,gRPC)
rustup target add wasm32-unknown-unknown
2.2 创建项目
cargo new my-msg-bot
cd my-msg-bot
2.3 完整 Cargo.toml
[package]
name = "msg-chain-sdk"
version = "1.0.0"
edition = "2021"
description = "MSG Chain Rust SDK - Post-quantum blockchain client"
license = "MIT"
[lib]
name = "msg_chain_sdk"
path = "src/lib.rs"
[dependencies]
tokio = { version = "1.36", features = ["full", "time"] }
futures-util = "0.3"
reqwest = { version = "0.12", features = ["json", "rustls-tls"] }
tokio-tungstenite = "0.21"
url = "2.5"
serde = { version = "1.0", features = ["derive"] }
serde_json = "1.0"
serde_with = "3.7"
bech32 = "0.11"
ripemd = "0.1"
sha2 = "0.10"
sha3 = "0.10"
hex = "0.4"
ed25519-dalek = { version = "2.1", features = ["serde"] }
ed25519 = "2.2"
pqcrypto-dilithium = "0.4"
pqcrypto-traits = "0.3"
bip39 = { version = "2.0", features = ["rand"] }
hdwallet = "1.0"
chrono = { version = "0.4", features = ["serde"] }
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
thiserror = "1.0"
anyhow = "1.0"
async-trait = "0.1"
base64 = "0.22"
uuid = { version = "1.7", features = ["v4"] }
rust_decimal = "1.34"
num-bigint = "0.4"
[dev-dependencies]
mockito = "1.4"
wiremock = "0.6"
tokio-test = "0.4"
criterion = { version = "0.5", features = ["async_tokio"] }
[features]
default = ["reqwest"]
native-tls = ["reqwest/native-tls"]
dilithium = ["pqcrypto-dilithium"]
2.4 环境变量配置
# .env
MSG_CHAIN_RPC_URL=https://rpc.msgchain.org
MSG_CHAIN_REST_URL=https://rest.msgchain.org
MSG_CHAIN_WS_URL=wss://ws.msgchain.org
MSG_CHAIN_CHAIN_ID=msg-chain-1
MSG_CHAIN_GAS_PRICE=1000000000
MSG_AGENT_API_KEY=your_api_key_here
MSG_AGENT_SECRET=your_agent_secret
MSG_WALLET_MNEMONIC="abandon abandon abandon ..."
2.5 配置模块
// src/config.rs
use std::env;
use std::str::FromStr;
use url::Url;
/// MSG Chain 网络配置
#[derive(Debug, Clone, serde::Deserialize, serde::Serialize)]
pub struct ChainConfig {
pub rpc_url: Url,
pub rest_url: Url,
pub ws_url: Url,
pub chain_id: String,
pub bech32_prefix: String,
pub coin_type: u32,
pub gas_price: Decimal,
pub gas_multiplier: f64,
pub block_time_seconds: u64,
}
impl Default for ChainConfig {
fn default() -> Self {
Self {
rpc_url: Url::from_str("https://rpc.msgchain.org").unwrap(),
rest_url: Url::from_str("https://rest.msgchain.org").unwrap(),
ws_url: Url::from_str("wss://ws.msgchain.org").unwrap(),
chain_id: "msg-chain-1".into(),
bech32_prefix: "msg".into(),
coin_type: 118,
gas_price: Decimal::from_str("1000000000").unwrap(),
gas_multiplier: 1.5,
block_time_seconds: 5,
}
}
}
impl ChainConfig {
pub fn from_env() -> Result<Self, ConfigError> {
Ok(Self {
rpc_url: env::var("MSG_CHAIN_RPC_URL")
.unwrap_or_else(|_| "https://rpc.msgchain.org".into())
.parse()
.map_err(|e| ConfigError::InvalidUrl(e))?,
rest_url: env::var("MSG_CHAIN_REST_URL")
.unwrap_or_else(|_| "https://rest.msgchain.org".into())
.parse()
.map_err(|e| ConfigError::InvalidUrl(e))?,
ws_url: env::var("MSG_CHAIN_WS_URL")
.unwrap_or_else(|_| "wss://ws.msgchain.org".into())
.parse()
.map_err(|e| ConfigError::InvalidUrl(e))?,
chain_id: env::var("MSG_CHAIN_CHAIN_ID")
.unwrap_or_else(|_| "msg-chain-1".into()),
bech32_prefix: String::from("msg"),
coin_type: 118,
gas_price: Decimal::from_str(
&env::var("MSG_CHAIN_GAS_PRICE")
.unwrap_or_else(|_| "1000000000".into())
).map_err(|e| ConfigError::ParseDecimal(e))?,
gas_multiplier: 1.5,
block_time_seconds: 5,
})
}
}
/// Gas 等级
#[derive(Debug, Clone, Copy, PartialEq)]
pub enum GasLevel {
Low,
Average,
High,
}
impl GasLevel {
pub fn price(&self) -> Decimal {
match self {
GasLevel::Low => Decimal::from_str("1000000000").unwrap(),
GasLevel::Average => Decimal::from_str("1000000000").unwrap(),
GasLevel::High => Decimal::from_str("1000000000").unwrap(),
}
}
}
/// 配置错误
#[derive(Debug, thiserror::Error)]
pub enum ConfigError {
#[error("无效 URL: {0}")]
InvalidUrl(#[from] url::ParseError),
#[error("环境变量缺失: {0}")]
MissingEnv(String),
#[error("解析十进制数失败: {0}")]
ParseDecimal(#[from] rust_decimal::Error),
#[error("IO 错误: {0}")]
Io(#[from] std::io::Error),
}
3. 核心客户端
3.1 全局错误类型
// src/error.rs
use std::time::Duration;
/// SDK 全局错误枚举
#[derive(Debug, thiserror::Error)]
pub enum SdkError {
// 网络层
#[error("HTTP 请求失败: {0}")]
HttpError(#[from] reqwest::Error),
#[error("WebSocket 错误: {0}")]
WsError(String),
#[error("连接超时 (timeout={0:?})")]
Timeout(Duration),
#[error("请求被限流 (retry_after={0:?})")]
RateLimited(Option<Duration>),
// 链上错误
#[error("RPC 错误: code={code}, message={message}")]
RpcError { code: i32, message: String },
#[error("交易失败: tx_hash={tx_hash}, code={code}, log={log}")]
TxFailed { tx_hash: String, code: u32, log: String },
#[error("交易未找到: {0}")]
TxNotFound(String),
#[error("Gas 估算失败: {0}")]
GasEstimationFailed(String),
// 序列化
#[error("JSON 序列化/反序列化错误: {0}")]
JsonError(#[from] serde_json::Error),
#[error("Base64 解码错误: {0}")]
Base64Error(#[from] base64::DecodeError),
#[error("Hex 解码错误: {0}")]
HexError(#[from] hex::FromHexError),
// 地址与密钥
#[error("Bech32 编码/解码错误: {0}")]
Bech32Error(String),
#[error("BIP39 助记词错误: {0}")]
MnemonicError(String),
#[error("签名错误: {0}")]
SigningError(String),
#[error("HD 派生路径错误: {0}")]
DerivationError(String),
// Agent API
#[error("Agent API 错误: {status}, body={body}")]
AgentApiError { status: u16, body: String },
#[error("API Key 未设置")]
ApiKeyMissing,
#[error("Stub 端点: {0}")]
StubEndpoint(String),
// 配置
#[error("配置错误: {0}")]
ConfigError(#[from] ConfigError),
#[error("无效输入: {0}")]
InvalidInput(String),
// 未知
#[error("未知错误: {0}")]
Unknown(String),
}
impl SdkError {
pub fn is_rate_limited(&self) -> bool { matches!(self, SdkError::RateLimited(_)) }
pub fn is_timeout(&self) -> bool { matches!(self, SdkError::Timeout(_)) }
pub fn is_stub(&self) -> bool { matches!(self, SdkError::StubEndpoint(_)) }
}
pub type Result<T> = std::result::Result<T, SdkError>;
3.2 Agent API 鉴权
// src/auth.rs
/// Agent API 鉴权配置
#[derive(Debug, Clone, serde::Deserialize, serde::Serialize)]
pub struct AgentAuth {
pub api_key: String,
pub api_secret: Option<String>,
}
impl AgentAuth {
pub fn from_env() -> Result<Self, SdkError> {
let api_key = std::env::var("MSG_AGENT_API_KEY")
.map_err(|_| SdkError::ApiKeyMissing)?;
let api_secret = std::env::var("MSG_AGENT_SECRET").ok();
Ok(Self { api_key, api_secret })
}
pub fn auth_header(&self) -> String {
format!("Bearer {}", self.api_key)
}
pub fn sign_header(&self, method: &str, path: &str, body: &str) -> Result<String, SdkError> {
let secret = self.api_secret.as_ref()
.ok_or_else(|| SdkError::SigningError("API secret not configured".into()))?;
let payload = format!("{}:{}:{}:{}", method, path, body, chrono::Utc::now().timestamp());
use sha2::Digest;
let mut mac = sha2::Sha256::new();
mac.update(secret.as_bytes());
mac.update(payload.as_bytes());
Ok(hex::encode(mac.finalize()))
}
}
3.3 令牌桶限流器
// src/rate_limiter.rs
use std::sync::Arc;
use tokio::sync::Mutex;
use std::time::Instant;
/// 令牌桶限流器
/// MSG Chain Agent API 限制: 100 req/s, burst 200
#[derive(Debug, Clone)]
pub struct RateLimiter {
inner: Arc<Mutex<Inner>>,
}
#[derive(Debug)]
struct Inner {
tokens: f64,
last_refill: Instant,
max_tokens: f64,
refill_rate: f64,
}
impl RateLimiter {
pub fn new(rate: f64, burst: f64) -> Self {
Self {
inner: Arc::new(Mutex::new(Inner {
tokens: burst,
last_refill: Instant::now(),
max_tokens: burst,
refill_rate: rate,
})),
}
}
pub fn agent_default() -> Self {
Self::new(100.0, 200.0)
}
pub async fn acquire(&self) -> Result<(), SdkError> {
loop {
let mut inner = self.inner.lock().await;
let now = Instant::now();
let elapsed = now.duration_since(inner.last_refill).as_secs_f64();
inner.tokens = (inner.tokens + elapsed * inner.refill_rate).min(inner.max_tokens);
inner.last_refill = now;
if inner.tokens >= 1.0 {
inner.tokens -= 1.0;
return Ok(());
}
let wait = Duration::from_secs_f64((1.0 - inner.tokens) / inner.refill_rate);
drop(inner);
tokio::time::sleep(wait).await;
}
}
}
3.4 完整核心客户端
// src/client.rs
use reqwest::header::{HeaderMap, HeaderValue, AUTHORIZATION, CONTENT_TYPE};
use reqwest::Client as HttpClient;
use std::sync::Arc;
use std::time::Duration;
use tokio::sync::RwLock;
use url::Url;
/// MSG Chain 核心客户端
///
/// 类似 ethers-rs 的 Middleware 模式,所有子客户端共享
/// 同一个 HTTP 连接池和限流器。
#[derive(Debug, Clone)]
pub struct MSGChainClient {
pub config: Arc<ChainConfig>,
http: HttpClient,
agent_auth: Option<Arc<AgentAuth>>,
rate_limiter: RateLimiter,
chain_id_cache: Arc<RwLock<Option<String>>>,
}
impl MSGChainClient {
/// 创建新客户端
pub fn new(config: ChainConfig) -> Result<Self, SdkError> {
let http = HttpClient::builder()
.timeout(Duration::from_secs(30))
.pool_max_idle_per_host(32)
.user_agent("msg-chain-rust-sdk/1.0.0")
.default_headers({
let mut h = HeaderMap::new();
h.insert(CONTENT_TYPE, HeaderValue::from_static("application/json"));
h
})
.build()?;
Ok(Self {
config: Arc::new(config),
http,
agent_auth: None,
rate_limiter: RateLimiter::agent_default(),
chain_id_cache: Arc::new(RwLock::new(None)),
})
}
/// 从环境变量创建
pub fn from_env() -> Result<Self, SdkError> {
let config = ChainConfig::from_env()?;
let mut client = Self::new(config)?;
if let Ok(auth) = AgentAuth::from_env() {
client = client.with_agent_auth(auth);
}
Ok(client)
}
/// 设置 Agent API 鉴权
pub fn with_agent_auth(mut self, auth: AgentAuth) -> Self {
self.agent_auth = Some(Arc::new(auth));
self
}
pub fn rest_base(&self) -> &Url { &self.config.rest_url }
pub fn rpc_url(&self) -> &Url { &self.config.rpc_url }
pub fn ws_url(&self) -> &Url { &self.config.ws_url }
// ========== 底层 HTTP 方法 ==========
pub async fn get<T: serde::de::DeserializeOwned>(
&self,
path: &str,
query: Option<&[(&str, &str)]>,
) -> Result<T> {
self.rate_limiter.acquire().await?;
let url = self.build_url(path, query)?;
let mut req = self.http.get(&url);
if let Some(auth) = &self.agent_auth {
req = req.header(AUTHORIZATION, auth.auth_header());
}
let resp = req.send().await?;
self.check_response_status(&resp).await?;
Ok(resp.json().await?)
}
pub async fn post<T: serde::de::DeserializeOwned, B: serde::Serialize + ?Sized>(
&self,
path: &str,
body: &B,
) -> Result<T> {
self.rate_limiter.acquire().await?;
let url = self.build_url(path, None)?;
let body_json = serde_json::to_string(body)?;
let mut req = self.http.post(&url)
.header(CONTENT_TYPE, "application/json")
.body(body_json.clone());
if let Some(auth) = &self.agent_auth {
req = req.header(AUTHORIZATION, auth.auth_header());
if let Some(secret) = &auth.api_secret {
let sig = auth.sign_header("POST", path, &body_json)?;
req = req.header("X-MSG-Signature", sig);
}
}
let resp = req.send().await?;
self.check_response_status(&resp).await?;
Ok(resp.json().await?)
}
/// 带 Stub 标记的 POST
pub async fn post_stub<T: serde::de::DeserializeOwned, B: serde::Serialize + ?Sized>(
&self,
path: &str,
body: &B,
) -> Result<T> {
self.rate_limiter.acquire().await?;
let url = self.build_url(path, None)?;
let body_json = serde_json::to_string(body)?;
let mut req = self.http.post(&url)
.header(CONTENT_TYPE, "application/json")
.header("X-MSG-Stub", "true")
.body(body_json.clone());
if let Some(auth) = &self.agent_auth {
req = req.header(AUTHORIZATION, auth.auth_header());
}
let resp = req.send().await?;
self.check_response_status(&resp).await?;
Ok(resp.json().await?)
}
fn build_url(&self, path: &str, query: Option<&[(&str, &str)]>) -> Result<Url, SdkError> {
let base = self.rest_base().to_string().trim_end_matches('/').to_string();
let p = path.trim_start_matches('/');
let mut url = Url::parse(&format!("{}/{}", base, p))
.map_err(|e| SdkError::InvalidInput(format!("URL 构建失败: {}", e)))?;
if let Some(params) = query {
url.query_pairs_mut().extend_pairs(params.iter());
}
Ok(url)
}
async fn check_response_status(&self, resp: &reqwest::Response) -> Result<()> {
let status = resp.status();
if status.is_success() { return Ok(()); }
let body = resp.text().await.unwrap_or_default();
match status.as_u16() {
429 => {
let retry_after = resp.headers()
.get("Retry-After")
.and_then(|v| v.to_str().ok())
.and_then(|v| v.parse::<u64>().ok())
.map(Duration::from_secs);
return Err(SdkError::RateLimited(retry_after));
}
408 | 504 => return Err(SdkError::Timeout(Duration::from_secs(30))),
401 | 403 => return Err(SdkError::AgentApiError {
status: status.as_u16(), body,
}),
_ => {
if status.is_server_error() {
return Err(SdkError::RpcError {
code: status.as_u16() as i32, message: body,
});
}
return Err(SdkError::AgentApiError {
status: status.as_u16(), body,
});
}
}
}
/// 健康检查
pub async fn health_check(&self) -> Result<bool> {
match self.get::<serde_json::Value>(
"/cosmos/base/tendermint/v1beta1/node_info", None
).await {
Ok(_) => Ok(true),
Err(e) => {
tracing::warn!("Health check failed: {}", e);
Ok(false)
}
}
}
/// 获取最新区块高度
pub async fn latest_block_height(&self) -> Result<u64> {
#[derive(serde::Deserialize)]
struct BlockResp { block: BlockInfo }
#[derive(serde::Deserialize)]
struct BlockInfo { header: BlockHeader }
#[derive(serde::Deserialize)]
struct BlockHeader { height: String }
let resp: BlockResp = self.get(
"/cosmos/base/tendermint/v1beta1/blocks/latest", None
).await?;
Ok(resp.block.header.height.parse::<u64>()
.map_err(|e| SdkError::InvalidInput(format!("无效区块高度: {}", e)))?)
}
}
// ========== 便捷构造 ==========
impl MSGChainClient {
pub fn bank_client(&self) -> BankClient { BankClient::new(self.clone()) }
pub fn staking_client(&self) -> StakingClient { StakingClient::new(self.clone()) }
pub fn contract_query_client(&self) -> ContractQueryClient { ContractQueryClient::new(self.clone()) }
pub fn tx_builder(&self) -> TxBuilder { TxBuilder::new(self.clone()) }
pub fn agent_client(&self) -> AgentClient { AgentClient::new(self.clone()) }
pub fn governance_client(&self) -> GovernanceClient { GovernanceClient::new(self.clone()) }
}
3.5 核心类型定义
// src/types.rs
use std::str::FromStr;
/// 交易哈希
#[derive(Debug, Clone, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
pub struct TxHash(pub String);
impl fmt::Display for TxHash {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "{}", self.0)
}
}
impl FromStr for TxHash {
type Err = SdkError;
fn from_str(s: &str) -> Result<Self, SdkError> {
if s.len() != 64 || !s.chars().all(|c| c.is_ascii_hexdigit()) {
return Err(SdkError::InvalidInput("无效交易哈希".into()));
}
Ok(TxHash(s.to_lowercase()))
}
}
/// 地址
#[derive(Debug, Clone, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
pub struct Address(pub String);
impl fmt::Display for Address {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "{}", self.0)
}
}
/// 代币面值
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct Coin {
pub denom: String,
pub amount: String,
}
impl Coin {
pub fn new(denom: impl Into<String>, amount: impl Into<String>) -> Self {
Self { denom: denom.into(), amount: amount.into() }
}
pub fn msg(amount: impl Into<String>) -> Self {
Self::new("umsg", amount)
}
pub fn is_zero(&self) -> bool { self.amount == "0" || self.amount == "0.0" }
}
/// 区块 ID
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct BlockId {
pub hash: String,
pub part_set_header: PartSetHeader,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct PartSetHeader {
pub total: u32,
pub hash: String,
}
/// 区块
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Block {
pub block_id: BlockId,
pub block: BlockData,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct BlockData {
pub header: BlockHeader,
pub data: Option<TxData>,
pub evidence: EvidenceList,
pub last_commit: Option<Commit>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct BlockHeader {
pub version: VersionInfo,
pub chain_id: String,
pub height: String,
pub time: String,
pub last_block_id: Option<BlockId>,
pub last_commit_hash: String,
pub data_hash: String,
pub validators_hash: String,
pub next_validators_hash: String,
pub consensus_hash: String,
pub app_hash: String,
pub last_results_hash: String,
pub evidence_hash: String,
pub proposer_address: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct VersionInfo {
pub block: String,
pub app: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct TxData {
pub txs: Option<Vec<String>>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct EvidenceList {
pub evidence: Vec<Evidence>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Evidence {
pub height: String,
pub time: String,
pub address: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Commit {
pub height: String,
pub round: i32,
pub block_id: BlockId,
pub signatures: Vec<CommitSig>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct CommitSig {
pub block_id_flag: i32,
pub validator_address: Option<String>,
pub timestamp: Option<String>,
pub signature: Option<String>,
}
/// 交易
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Transaction {
pub hash: TxHash,
pub height: String,
pub index: u32,
pub tx_result: TxResult,
pub tx: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct TxResult {
pub code: u32,
pub data: Option<String>,
pub log: String,
pub info: String,
pub gas_wanted: String,
pub gas_used: String,
pub events: Vec<Event>,
}
/// 事件
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Event {
#[serde(rename = "type")]
pub event_type: String,
pub attributes: Vec<EventAttribute>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct EventAttribute {
pub key: String,
pub value: String,
pub index: bool,
}
/// 金额(精度处理)
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Amount(pub String);
impl Amount {
/// 从 umsg 转换为 MSG(18 位小数)
pub fn from_umsg(umsg: impl Into<String>) -> Self {
let raw = umsg.into();
let len = raw.len();
if len <= 18 {
let padded = format!("{:0>18}", raw);
let int_part = &padded[..padded.len() - 18];
let dec_part = &padded[padded.len() - 18..];
let trimmed = dec_part.trim_end_matches('0');
let dec = if trimmed.is_empty() { "0" } else { trimmed };
if int_part.is_empty() || int_part == "0" {
Amount(format!("0.{}", dec))
} else {
Amount(format!("{}.{}", int_part, dec))
}
} else {
let split = len - 18;
Amount(format!("{}.{}", &raw[..split], &raw[split..]))
}
}
/// 转换为 umsg(乘以 10^18)
pub fn to_umsg(&self) -> String {
if let Some(dot) = self.0.find('.') {
let int_part = &self.0[..dot];
let dec_part = &self.0[dot + 1..];
let dec_padded = format!("{:<018}", dec_part);
format!("{}{}", int_part, dec_padded)
} else {
format!("{}0000000000000000000", self.0)
}
}
}
/// 分页
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Pagination {
pub next_key: Option<String>,
pub total: String,
}
/// 分页请求参数
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct PaginationParams {
pub key: Option<String>,
pub offset: Option<u64>,
pub limit: Option<u64>,
pub count_total: Option<bool>,
pub reverse: Option<bool>,
}
4. 地址与密钥
4.1 Bech32 地址
// src/crypto/address.rs
use bech32::{Bech32, Hrp};
use ripemd::Ripemd160;
use sha2::{Digest, Sha256};
/// MSG Chain 地址 (Bech32)
///
/// 格式: `msg1<32字符hex>`
#[derive(Debug, Clone, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
#[serde(try_from = "String", into = "String")]
pub struct Bech32Address {
pub prefix: String,
pub data: Vec<u8>,
pub bech32: String,
}
impl Bech32Address {
/// 从公钥生成地址 (RIPEMD160(SHA256(pubkey)))
pub fn from_public_key_bytes(prefix: &str, pubkey: &[u8]) -> Result<Self, SdkError> {
let sha = Sha256::digest(pubkey);
let ripemd = Ripemd160::digest(&sha);
Self::from_bytes(prefix, &ripemd)
}
pub fn from_bytes(prefix: &str, data: &[u8]) -> Result<Self, SdkError> {
let hrp = Hrp::parse(prefix)
.map_err(|e| SdkError::Bech32Error(format!("无效 HRP: {}", e)))?;
let bech32 = bech32::encode::<Bech32>(hrp, data)
.map_err(|e| SdkError::Bech32Error(format!("Bech32 编码失败: {}", e)))?;
Ok(Self {
prefix: prefix.to_string(),
data: data.to_vec(),
bech32,
})
}
pub fn from_bech32(s: &str) -> Result<Self, SdkError> {
let (hrp, data) = bech32::decode(s)
.map_err(|e| SdkError::Bech32Error(format!("Bech32 解码失败: {}", e)))?;
Ok(Self {
prefix: hrp.to_string(),
data: data.to_vec(),
bech32: s.to_string(),
})
}
pub fn is_valid(s: &str) -> bool { bech32::decode(s).is_ok() }
pub fn msg_address(data: &[u8]) -> Result<Self, SdkError> {
Self::from_bytes("msg", data)
}
pub fn hex(&self) -> String { hex::encode(&self.data) }
}
impl fmt::Display for Bech32Address {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "{}", self.bech32)
}
}
impl FromStr for Bech32Address {
type Err = SdkError;
fn from_str(s: &str) -> Result<Self, SdkError> { Self::from_bech32(s) }
}
impl From<Bech32Address> for String {
fn from(addr: Bech32Address) -> String { addr.bech32 }
}
impl TryFrom<String> for Bech32Address {
type Error = SdkError;
fn try_from(s: String) -> Result<Self, SdkError> { Self::from_bech32(&s) }
}
4.2 密钥对与钱包
// src/crypto/keys.rs
use bip39::{Language, Mnemonic, MnemonicType, Seed};
use hdwallet::{DerivationPath, ExtendedPrivKey};
/// Dilithium-5 密钥对(后量子)
#[derive(Debug, Clone)]
pub struct DilithiumKeypair {
pub public_key: Vec<u8>,
pub secret_key: Vec<u8>,
}
impl DilithiumKeypair {
pub fn generate() -> Self {
let (pk, sk) = pqcrypto_dilithium::dilithium5::keypair();
Self {
public_key: pk.as_ref().to_vec(),
secret_key: sk.as_ref().to_vec(),
}
}
pub fn from_seed(seed: &[u8], path: &str) -> Result<Self, SdkError> {
let derivation_path = DerivationPath::from_str(path)
.map_err(|e| SdkError::DerivationError(e.to_string()))?;
let ext_key = ExtendedPrivKey::new(seed)
.map_err(|e| SdkError::DerivationError(e.to_string()))?;
let child = ext_key.derive(&derivation_path)
.map_err(|e| SdkError::DerivationError(e.to_string()))?;
let seed_for_dilithium = Sha256::digest(child.private_key().to_bytes());
let (pk, sk) = pqcrypto_dilithium::dilithium5::keypair_from_seed(&seed_for_dilithium);
Ok(Self {
public_key: pk.as_ref().to_vec(),
secret_key: sk.as_ref().to_vec(),
})
}
pub fn sign(&self, msg: &[u8]) -> Vec<u8> {
pqcrypto_dilithium::dilithium5::sign(msg, &self.secret_key.clone().into())
.as_ref().to_vec()
}
pub fn verify(&self, msg: &[u8], sig: &[u8]) -> bool {
pqcrypto_dilithium::dilithium5::verify(msg, sig, &self.public_key.clone().into()).is_ok()
}
pub fn address(&self) -> Result<Bech32Address, SdkError> {
Bech32Address::from_public_key_bytes("msg", &self.public_key)
}
}
/// Ed25519 密钥对(兼容模式)
#[derive(Debug, Clone)]
pub struct Ed25519Keypair {
pub keypair: ed25519_dalek::SigningKey,
}
impl Ed25519Keypair {
pub fn generate() -> Self {
let mut rng = rand::thread_rng();
Self { keypair: ed25519_dalek::SigningKey::generate(&mut rng) }
}
pub fn from_seed(seed: &[u8]) -> Result<Self, SdkError> {
let bytes: [u8; 32] = seed[..32].try_into()
.map_err(|_| SdkError::SigningError("种子需要至少32字节".into()))?;
Ok(Self { keypair: ed25519_dalek::SigningKey::from_bytes(&bytes) })
}
pub fn sign(&self, msg: &[u8]) -> ed25519_dalek::Signature {
use ed25519_dalek::Signer;
self.keypair.sign(msg)
}
pub fn verify(&self, msg: &[u8], sig: &ed25519_dalek::Signature) -> bool {
use ed25519_dalek::Verifier;
self.keypair.verify(msg, sig).is_ok()
}
pub fn public_key_bytes(&self) -> Vec<u8> {
self.keypair.verifying_key().as_bytes().to_vec()
}
pub fn address(&self) -> Result<Bech32Address, SdkError> {
Bech32Address::from_public_key_bytes("msg", &self.public_key_bytes())
}
}
/// 完整钱包
#[derive(Debug, Clone)]
pub struct Wallet {
pub mnemonic: String,
pub seed: Vec<u8>,
pub dilithium_key: DilithiumKeypair,
pub address: Bech32Address,
}
impl Wallet {
pub fn from_mnemonic(mnemonic_str: &str) -> Result<Self, SdkError> {
let mnemonic = Mnemonic::from_phrase(mnemonic_str, Language::English)
.map_err(|e| SdkError::MnemonicError(e.to_string()))?;
let seed = Seed::new(&mnemonic, "");
let path = "m/44'/118'/0'/0/0";
let dilithium_key = DilithiumKeypair::from_seed(seed.as_bytes(), path)?;
let address = dilithium_key.address()?;
Ok(Self {
mnemonic: mnemonic.phrase().to_string(),
seed: seed.as_bytes().to_vec(),
dilithium_key,
address,
})
}
pub fn generate() -> Result<Self, SdkError> {
let mnemonic = Mnemonic::new(MnemonicType::Words24, Language::English);
Self::from_mnemonic(mnemonic.phrase())
}
pub fn sign_tx(&self, tx_bytes: &[u8]) -> Vec<u8> {
self.dilithium_key.sign(tx_bytes)
}
}
/// HD 派生路径
pub mod derivation {
use hdwallet::DerivationPath;
use std::str::FromStr;
pub const MSG_DEFAULT_PATH: &str = "m/44'/118'/0'/0/0";
pub fn bip44_path(coin_type: u32, account: u32, change: u32, index: u32) -> String {
format!("m/44'/{coin_type}'/{account}'/{change}/{index}")
}
pub fn msg_path(account: u32, index: u32) -> String {
bip44_path(118, account, 0, index)
}
pub fn parse_path(path: &str) -> Result<DerivationPath, SdkError> {
DerivationPath::from_str(path)
.map_err(|e| SdkError::DerivationError(e.to_string()))
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_address_roundtrip() {
let addr = Bech32Address::from_bytes("msg", &[0u8; 20]).unwrap();
let decoded = Bech32Address::from_bech32(&addr.bech32).unwrap();
assert_eq!(addr, decoded);
}
#[test]
fn test_address_format() {
let addr = Bech32Address::from_bytes("msg", &[1u8; 20]).unwrap();
assert!(addr.bech32.starts_with("msg1"));
}
#[test]
fn test_wallet_generation() {
let wallet = Wallet::generate().unwrap();
assert!(wallet.address.bech32.starts_with("msg1"));
assert_eq!(wallet.address.prefix, "msg");
}
#[test]
fn test_dilithium_sign_verify() {
let kp = DilithiumKeypair::generate();
let msg = b"hello msg chain";
let sig = kp.sign(msg);
assert!(kp.verify(msg, &sig));
assert!(!kp.verify(b"tampered", &sig));
}
}
5. 银行查询
// src/clients/bank.rs
/// 银行查询客户端
#[derive(Debug, Clone)]
pub struct BankClient {
client: MSGChainClient,
}
impl BankClient {
pub fn new(client: MSGChainClient) -> Self { Self { client } }
/// 查询账户信息
pub async fn account(&self, address: &Bech32Address) -> Result<AccountInfo, SdkError> {
let path = format!("/cosmos/auth/v1beta1/accounts/{}", address.bech32);
self.client.get(&path, None).await
}
/// 查询单一代币余额
pub async fn balance(
&self, address: &Bech32Address, denom: &str,
) -> Result<Coin, SdkError> {
#[derive(serde::Deserialize)]
struct BalanceResp { balance: Coin }
let path = format!(
"/cosmos/bank/v1beta1/balances/{}/by_denom", address.bech32
);
let resp: BalanceResp = self.client
.get(&path, Some(&[("denom", denom)])).await?;
Ok(resp.balance)
}
/// 查询所有余额
pub async fn balances(&self, address: &Bech32Address) -> Result<Vec<Coin>, SdkError> {
#[derive(serde::Deserialize)]
struct BalancesResp { balances: Vec<Coin>, pagination: Option<Pagination> }
let path = format!("/cosmos/bank/v1beta1/balances/{}", address.bech32);
let resp: BalancesResp = self.client.get(&path, None).await?;
Ok(resp.balances)
}
/// 查询总供应量
pub async fn total_supply(&self) -> Result<Vec<Coin>, SdkError> {
#[derive(serde::Deserialize)]
struct SupplyResp { supply: Vec<Coin> }
let resp: SupplyResp = self.client.get("/cosmos/bank/v1beta1/supply", None).await?;
Ok(resp.supply)
}
/// 查询 MSG 余额
pub async fn msg_balance(&self, address: &Bech32Address) -> Result<Amount, SdkError> {
let coin = self.balance(address, "umsg").await?;
Ok(Amount::from_umsg(coin.amount))
}
/// 查询历史高度的余额
pub async fn balance_at_height(
&self, address: &Bech32Address, denom: &str, height: u64,
) -> Result<Coin, SdkError> {
let path = format!(
"/cosmos/bank/v1beta1/balances/{}/by_denom", address.bech32
);
let resp = self.client.get::<serde_json::Value>(
&path, Some(&[("denom", denom), ("height", &height.to_string())]),
).await?;
Ok(serde_json::from_value::<Coin>(resp["balance"].clone())?)
}
/// 可花费余额
pub async fn spendable_balances(
&self, address: &Bech32Address,
) -> Result<Vec<Coin>, SdkError> {
#[derive(serde::Deserialize)]
struct SpendableResp { balances: Vec<Coin> }
let path = format!(
"/cosmos/bank/v1beta1/spendable_balances/{}", address.bech32
);
let resp: SpendableResp = self.client.get(&path, None).await?;
Ok(resp.balances)
}
}
/// 账户信息
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct AccountInfo { pub account: AccountData }
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "@type")]
pub enum AccountData {
#[serde(rename = "/cosmos.auth.v1beta1.BaseAccount")]
Base(BaseAccount),
#[serde(rename = "/cosmos.auth.v1beta1.ModuleAccount")]
Module(ModuleAccount),
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct BaseAccount {
pub address: String,
pub pub_key: Option<PublicKey>,
pub account_number: String,
pub sequence: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct ModuleAccount {
pub base_account: BaseAccount,
pub name: String,
pub permissions: Vec<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct PublicKey {
#[serde(rename = "@type")]
pub key_type: String,
pub key: String,
}
6. 质押查询
// src/clients/staking.rs
/// 质押查询客户端
#[derive(Debug, Clone)]
pub struct StakingClient {
client: MSGChainClient,
}
impl StakingClient {
pub fn new(client: MSGChainClient) -> Self { Self { client } }
pub async fn validators(&self, status: Option<&str>) -> Result<Vec<Validator>, SdkError> {
#[derive(serde::Deserialize)]
struct ValidatorsResp { validators: Vec<Validator>, pagination: Pagination }
let mut query = Vec::new();
if let Some(s) = status { query.push(("status", s)); }
let resp: ValidatorsResp = self.client
.get("/cosmos/staking/v1beta1/validators", Some(&query)).await?;
Ok(resp.validators)
}
pub async fn validator(&self, address: &str) -> Result<Validator, SdkError> {
#[derive(serde::Deserialize)]
struct ValidatorResp { validator: Validator }
let path = format!("/cosmos/staking/v1beta1/validators/{}", address);
let resp: ValidatorResp = self.client.get(&path, None).await?;
Ok(resp.validator)
}
pub async fn delegations(
&self, delegator: &Bech32Address,
) -> Result<Vec<Delegation>, SdkError> {
#[derive(serde::Deserialize)]
struct DelegationsResp { delegation_responses: Vec<Delegation> }
let path = format!("/cosmos/staking/v1beta1/delegations/{}", delegator.bech32);
let resp: DelegationsResp = self.client.get(&path, None).await?;
Ok(resp.delegation_responses)
}
pub async fn unbonding_delegations(
&self, delegator: &Bech32Address,
) -> Result<Vec<UnbondingDelegation>, SdkError> {
#[derive(serde::Deserialize)]
struct UnbondingResp { unbonding_responses: Vec<UnbondingDelegation> }
let path = format!(
"/cosmos/staking/v1beta1/delegators/{}/unbonding_delegations", delegator.bech32
);
let resp: UnbondingResp = self.client.get(&path, None).await?;
Ok(resp.unbonding_responses)
}
pub async fn rewards(
&self, delegator: &Bech32Address,
) -> Result<Rewards, SdkError> {
let path = format!(
"/cosmos/distribution/v1beta1/delegators/{}/rewards", delegator.bech32
);
let resp = self.client.get(&path, None).await?;
Ok(resp)
}
pub async fn delegation_total_rewards(
&self, delegator: &Bech32Address,
) -> Result<Vec<Coin>, SdkError> {
let rewards = self.rewards(delegator).await?;
Ok(rewards.total)
}
pub async fn staking_pool(&self) -> Result<StakingPool, SdkError> {
#[derive(serde::Deserialize)]
struct PoolResp { pool: StakingPool }
let resp: PoolResp = self.client
.get("/cosmos/staking/v1beta1/pool", None).await?;
Ok(resp.pool)
}
}
/// 验证者
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Validator {
pub operator_address: String,
pub consensus_pubkey: Option<PublicKey>,
pub jailed: bool,
pub status: BondStatus,
pub tokens: String,
pub delegator_shares: String,
pub description: ValidatorDescription,
pub unbonding_height: String,
pub unbonding_time: String,
pub commission: CommissionInfo,
pub min_self_delegation: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct ValidatorDescription {
pub moniker: String,
pub identity: Option<String>,
pub website: Option<String>,
pub security_contact: Option<String>,
pub details: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct CommissionInfo {
pub commission_rates: CommissionRates,
pub update_time: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct CommissionRates {
pub rate: String,
pub max_rate: String,
pub max_change_rate: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub enum BondStatus {
Unspecified, Unbonded, Unbonding, Bonded,
}
/// 委托
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Delegation {
pub delegation: DelegationInfo,
pub balance: Coin,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct DelegationInfo {
pub delegator_address: String,
pub validator_address: String,
pub shares: String,
}
/// 未委托
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct UnbondingDelegation {
pub delegator_address: String,
pub validator_address: String,
pub entries: Vec<UnbondingEntry>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct UnbondingEntry {
pub creation_height: String,
pub completion_time: String,
pub initial_balance: String,
pub balance: String,
}
/// 收益
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Rewards {
pub rewards: Vec<DelegationRewards>,
pub total: Vec<Coin>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct DelegationRewards {
pub validator_address: String,
pub reward: Vec<Coin>,
}
/// 质押池
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct StakingPool {
pub not_bonded_tokens: String,
pub bonded_tokens: String,
}
7. 合约查询
7.1 CosmWasm 基础查询
// src/clients/contract/query.rs
/// 合约查询客户端
///
/// 支持 CosmWasm 智能合约的各类查询,包括 5 个 AI 原生合约。
#[derive(Debug, Clone)]
pub struct ContractQueryClient {
client: MSGChainClient,
}
impl ContractQueryClient {
pub fn new(client: MSGChainClient) -> Self { Self { client } }
/// 通用合约查询
pub async fn query_contract<T: DeserializeOwned>(
&self,
contract: &Bech32Address,
query_msg: &impl Serialize,
) -> Result<T, SdkError> {
let msg_val = serde_json::to_value(query_msg)?;
let path = format!("/cosmwasm/wasm/v1/contract/{}/smart", contract.bech32);
let req = serde_json::json!({ "msg": msg_val });
#[derive(serde::Deserialize)]
struct QueryResp { data: serde_json::Value }
let resp: QueryResp = self.client.post(&path, &req).await?;
Ok(serde_json::from_value(resp.data)?)
}
/// 合约信息
pub async fn contract_info(
&self, contract: &Bech32Address,
) -> Result<ContractInfo, SdkError> {
#[derive(serde::Deserialize)]
struct ContractInfoResp { contract_info: ContractInfo }
let path = format!("/cosmwasm/wasm/v1/contract/{}", contract.bech32);
let resp: ContractInfoResp = self.client.get(&path, None).await?;
Ok(resp.contract_info)
}
/// 合约状态
pub async fn contract_state(
&self, contract: &Bech32Address,
) -> Result<Vec<ModelEntry>, SdkError> {
#[derive(serde::Deserialize)]
struct StateResp { models: Vec<ModelEntry>, pagination: Pagination }
let path = format!("/cosmwasm/wasm/v1/contract/{}/state", contract.bech32);
let resp: StateResp = self.client.get(&path, None).await?;
Ok(resp.models)
}
/// 查询所有合约
pub async fn contracts(&self) -> Result<Vec<String>, SdkError> {
#[derive(serde::Deserialize)]
struct ContractsResp { contract_infos: Vec<ContractInfoRaw>, pagination: Pagination }
#[derive(serde::Deserialize)]
struct ContractInfoRaw { address: String }
let resp: ContractsResp = self.client.get("/cosmwasm/wasm/v1/code", None).await?;
Ok(resp.contract_infos.into_iter().map(|c| c.address).collect())
}
/// 代码信息
pub async fn code_info(&self, code_id: u64) -> Result<CodeInfo, SdkError> {
#[derive(serde::Deserialize)]
struct CodeResp { code_info: CodeInfo }
let path = format!("/cosmwasm/wasm/v1/code/{}", code_id);
let resp: CodeResp = self.client.get(&path, None).await?;
Ok(resp.code_info)
}
}
/// 合约信息
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct ContractInfo {
pub code_id: String,
pub creator: String,
pub admin: Option<String>,
pub label: String,
pub created: Option<AbsoluteTxPosition>,
pub ibc_port_id: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct AbsoluteTxPosition {
pub block_height: String,
pub tx_index: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct ModelEntry {
pub key: String,
pub value: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct CodeInfo {
pub code_id: String,
pub creator: String,
pub data_hash: String,
pub instantiate_permission: Option<AccessConfig>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct AccessConfig {
pub permission: AccessType,
pub addresses: Vec<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub enum AccessType {
Unspecified, Nobody, Everybody, AnyOfAddresses,
}
7.2 AI 合约类型化查询
// src/clients/contract/ai_contracts.rs
/// AI 合约类型化查询接口
///
/// MSG Chain 5 个 AI 原生 CosmWasm 合约:
/// 1. ai_model_registry — AI 模型元数据注册
/// 2. ai_inference_proof — 推理证明验证
/// 3. ai_data_asset — 数据资产管理(NFT)
/// 4. ai_agent_registry — AI Agent 注册与信誉
/// 5. ai_inference_market — 推理请求与匹配市场
impl ContractQueryClient {
// ===== 1. AI 模型注册 =====
pub async fn ai_model_info(
&self, contract: &Bech32Address, model_id: &str,
) -> Result<AiModelInfo, SdkError> {
self.query_contract(contract, &serde_json::json!({
"model_info": { "model_id": model_id }
})).await
}
pub async fn ai_list_models(
&self, contract: &Bech32Address, start_after: Option<String>, limit: Option<u32>,
) -> Result<Vec<AiModelSummary>, SdkError> {
self.query_contract(contract, &serde_json::json!({
"list_models": { "start_after": start_after, "limit": limit }
})).await
}
pub async fn ai_model_versions(
&self, contract: &Bech32Address, model_id: &str,
) -> Result<Vec<AiModelVersion>, SdkError> {
self.query_contract(contract, &serde_json::json!({
"model_versions": { "model_id": model_id }
})).await
}
// ===== 2. 推理证明 =====
pub async fn ai_inference_proof(
&self, contract: &Bech32Address, proof_id: &str,
) -> Result<InferenceProof, SdkError> {
self.query_contract(contract, &serde_json::json!({
"proof": { "proof_id": proof_id }
})).await
}
pub async fn ai_verify_proof(
&self, contract: &Bech32Address, proof_id: &str, expected_output_hash: &str,
) -> Result<ProofVerificationResult, SdkError> {
self.query_contract(contract, &serde_json::json!({
"verify_proof": { "proof_id": proof_id, "expected_output_hash": expected_output_hash }
})).await
}
// ===== 3. 数据资产 (NFT) =====
pub async fn ai_data_asset(
&self, contract: &Bech32Address, token_id: &str,
) -> Result<DataAsset, SdkError> {
self.query_contract(contract, &serde_json::json!({
"data_asset": { "token_id": token_id }
})).await
}
pub async fn ai_asset_owner(
&self, contract: &Bech32Address, token_id: &str,
) -> Result<AssetOwner, SdkError> {
self.query_contract(contract, &serde_json::json!({
"owner_of": { "token_id": token_id }
})).await
}
// ===== 4. AI Agent 注册 =====
pub async fn ai_agent_info(
&self, contract: &Bech32Address, agent_id: &str,
) -> Result<AgentInfo, SdkError> {
self.query_contract(contract, &serde_json::json!({
"agent_info": { "agent_id": agent_id }
})).await
}
pub async fn ai_agent_reputation(
&self, contract: &Bech32Address, agent_id: &str,
) -> Result<AgentReputation, SdkError> {
self.query_contract(contract, &serde_json::json!({
"reputation": { "agent_id": agent_id }
})).await
}
pub async fn ai_agents_by_capability(
&self, contract: &Bech32Address, capability: &str, limit: Option<u32>,
) -> Result<Vec<AgentSummary>, SdkError> {
self.query_contract(contract, &serde_json::json!({
"agents_by_capability": { "capability": capability, "limit": limit }
})).await
}
// ===== 5. 推理市场 =====
pub async fn ai_inference_request(
&self, contract: &Bech32Address, request_id: &str,
) -> Result<InferenceRequest, SdkError> {
self.query_contract(contract, &serde_json::json!({
"inference_request": { "request_id": request_id }
})).await
}
pub async fn ai_pending_requests(
&self, contract: &Bech32Address, start_after: Option<String>, limit: Option<u32>,
) -> Result<Vec<InferenceRequest>, SdkError> {
self.query_contract(contract, &serde_json::json!({
"pending_requests": { "start_after": start_after, "limit": limit }
})).await
}
pub async fn ai_market_stats(
&self, contract: &Bech32Address,
) -> Result<MarketStats, SdkError> {
self.query_contract(contract, &serde_json::json!({
"market_stats": {}
})).await
}
}
// ===== AI 合约类型定义 =====
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct AiModelInfo {
pub model_id: String,
pub name: String,
pub version: String,
pub model_type: String,
pub owner: String,
pub description: Option<String>,
pub ipfs_hash: Option<String>,
pub created_at: String,
pub status: ModelStatus,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum ModelStatus { Active, Deprecated, Banned }
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct AiModelSummary {
pub model_id: String,
pub name: String,
pub version: String,
pub owner: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct AiModelVersion {
pub version: String,
pub ipfs_hash: String,
pub created_at: String,
pub checksum: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct InferenceProof {
pub proof_id: String,
pub request_id: String,
pub model_id: String,
pub input_hash: String,
pub output_hash: String,
pub prover: String,
pub proof_data: String,
pub verified: bool,
pub created_at: String,
pub proof_type: ProofType,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum ProofType { ZkSnark, OpTee, SecureEnclave }
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct ProofVerificationResult {
pub verified: bool,
pub verification_time: String,
pub details: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct DataAsset {
pub token_id: String,
pub owner: String,
pub name: String,
pub description: Option<String>,
pub data_type: String,
pub data_hash: String,
pub data_uri: Option<String>,
pub license: Option<String>,
pub created_at: String,
pub royalty_percentage: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct AssetOwner {
pub owner: String,
pub approvals: Vec<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct AgentInfo {
pub agent_id: String,
pub name: String,
pub owner: String,
pub description: Option<String>,
pub model_id: String,
pub capabilities: Vec<String>,
pub endpoint: Option<String>,
pub status: AgentStatus,
pub created_at: String,
pub metadata: Option<serde_json::Value>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum AgentStatus { Active, Paused, Banned }
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct AgentSummary {
pub agent_id: String,
pub name: String,
pub capabilities: Vec<String>,
pub status: AgentStatus,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct AgentReputation {
pub agent_id: String,
pub total_tasks: u64,
pub completed_tasks: u64,
pub success_rate: String,
pub avg_response_time: String,
pub total_rewards: String,
pub rating: String,
pub reviews: u64,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct InferenceRequest {
pub request_id: String,
pub requester: String,
pub model_id: String,
pub input_hash: String,
pub required_proof: ProofType,
pub max_price: String,
pub status: RequestStatus,
pub assigned_agent: Option<String>,
pub created_at: String,
pub expires_at: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum RequestStatus {
Pending, Assigned, Completed, Cancelled, Expired,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct MarketStats {
pub total_requests: u64,
pub completed_requests: u64,
pub total_agents: u64,
pub total_models: u64,
pub total_volume: String,
pub avg_price: String,
}
8. 合约执行
8.1 交易构建与广播
// src/tx/builder.rs
/// 交易构建器
///
/// 类似 ethers-rs 的 TransactionBuilder:
/// - 自动 Gas 估算与调整
/// - 序列号管理
/// - Dilithium-5 签名
/// - 并发广播与等待确认
#[derive(Debug, Clone)]
pub struct TxBuilder {
client: MSGChainClient,
gas_price: Decimal,
gas_multiplier: f64,
timeout_height: u64,
}
impl TxBuilder {
pub fn new(client: MSGChainClient) -> Self {
Self {
client,
gas_price: Decimal::from_str("1000000000").unwrap(),
gas_multiplier: 1.5,
timeout_height: 100,
}
}
pub fn with_gas_price(mut self, price: Decimal) -> Self { self.gas_price = price; self }
pub fn with_gas_multiplier(mut self, mult: f64) -> Self { self.gas_multiplier = mult; self }
/// MsgSend 转账
pub async fn send_tokens(
&self, wallet: &Wallet, to: &Bech32Address,
amount: Vec<Coin>, memo: Option<&str>,
) -> Result<TxResponse, SdkError> {
let msg = serde_json::json!({
"@type": "/cosmos.bank.v1beta1.MsgSend",
"from_address": wallet.address.bech32,
"to_address": to.bech32,
"amount": amount,
});
self.build_and_broadcast(wallet, vec![msg], memo).await
}
/// 合约执行
pub async fn execute_contract(
&self, wallet: &Wallet, contract: &Bech32Address,
execute_msg: &impl Serialize, funds: Vec<Coin>, memo: Option<&str>,
) -> Result<TxResponse, SdkError> {
let msg_val = serde_json::to_value(execute_msg)?;
let msg = serde_json::json!({
"@type": "/cosmwasm.wasm.v1.MsgExecuteContract",
"sender": wallet.address.bech32,
"contract": contract.bech32,
"msg": serde_json::to_string(&msg_val).unwrap_or_default(),
"funds": funds,
});
self.build_and_broadcast(wallet, vec![msg], memo).await
}
/// 合约实例化
pub async fn instantiate_contract(
&self, wallet: &Wallet, code_id: u64,
instantiate_msg: &impl Serialize, label: &str,
admin: Option<&Bech32Address>, funds: Vec<Coin>,
) -> Result<TxResponse, SdkError> {
let msg_val = serde_json::to_value(instantiate_msg)?;
let admin_str = admin.map(|a| a.bech32.clone());
let msg = serde_json::json!({
"@type": "/cosmwasm.wasm.v1.MsgInstantiateContract",
"sender": wallet.address.bech32,
"admin": admin_str,
"code_id": code_id.to_string(),
"label": label,
"msg": serde_json::to_string(&msg_val).unwrap_or_default(),
"funds": funds,
});
self.build_and_broadcast(wallet, vec![msg], None).await
}
/// 质押委托
pub async fn delegate_tokens(
&self, wallet: &Wallet, validator_addr: &str, amount: Coin,
) -> Result<TxResponse, SdkError> {
let msg = serde_json::json!({
"@type": "/cosmos.staking.v1beta1.MsgDelegate",
"delegator_address": wallet.address.bech32,
"validator_address": validator_addr,
"amount": amount,
});
self.build_and_broadcast(wallet, vec![msg], None).await
}
/// 治理投票
pub async fn gov_vote(
&self, wallet: &Wallet, proposal_id: u64, option: VoteOption,
) -> Result<TxResponse, SdkError> {
let msg = serde_json::json!({
"@type": "/cosmos.gov.v1beta1.MsgVote",
"proposal_id": proposal_id.to_string(),
"voter": wallet.address.bech32,
"option": option as u32,
});
self.build_and_broadcast(wallet, vec![msg], None).await
}
/// 核心方法:构建 + 签名 + 广播
async fn build_and_broadcast(
&self, wallet: &Wallet, msgs: Vec<serde_json::Value>, memo: Option<&str>,
) -> Result<TxResponse, SdkError> {
// 1. 获取账户信息
let bank = self.client.bank_client();
let account_info = bank.account(&wallet.address).await?;
let (sequence, account_number) = match &account_info.account {
AccountData::Base(acc) => (acc.sequence.clone(), acc.account_number.clone()),
AccountData::Module(acc) => (
acc.base_account.sequence.clone(),
acc.base_account.account_number.clone(),
),
};
// 2. 估算 Gas
let gas_limit = self.estimate_gas(&msgs, memo).await?;
// 3. 构建交易
let tx = self.build_std_tx(&msgs, memo.unwrap_or(""), &gas_limit, &sequence);
// 4. 签名
let tx_bytes = serde_json::to_vec(&tx)?;
let signature = wallet.sign_tx(&tx_bytes);
// 5. 构建签名的交易并广播
let signed_tx = serde_json::json!({
"tx": tx,
"signatures": [{
"pub_key": {
"@type": "/msgchain.crypto.DilithiumPubKey",
"key": hex::encode(&wallet.dilithium_key.public_key),
},
"signature": hex::encode(&signature),
}],
});
let broadcast_req = serde_json::json!({
"tx_bytes": base64::encode(&serde_json::to_vec(&signed_tx)?),
"mode": "BROADCAST_MODE_SYNC",
});
#[derive(serde::Deserialize)]
struct BroadcastResp { tx_response: TxResponse }
let resp: BroadcastResp = self.client
.post("/cosmos/tx/v1beta1/txs", &broadcast_req).await?;
Ok(resp.tx_response)
}
/// 估算 Gas
async fn estimate_gas(
&self, msgs: &[serde_json::Value], memo: Option<&str>,
) -> Result<String, SdkError> {
let sim_req = serde_json::json!({
"tx": {
"body": {
"messages": msgs,
"memo": memo.unwrap_or(""),
"timeout_height": "0",
},
"auth_info": {
"signer_infos": [],
"fee": { "gas_limit": "200000", "amount": [] },
},
"signatures": [],
},
});
#[derive(serde::Deserialize)]
struct SimResp { gas_info: GasInfo }
#[derive(serde::Deserialize)]
struct GasInfo { gas_used: String }
let resp: SimResp = self.client
.post("/cosmos/tx/v1beta1/simulate", &sim_req).await?;
let gas_used: u64 = resp.gas_info.gas_used.parse()
.map_err(|e| SdkError::GasEstimationFailed(e.to_string()))?;
let gas_adjusted = (gas_used as f64 * self.gas_multiplier).ceil() as u64;
Ok(gas_adjusted.to_string())
}
fn build_std_tx(
&self, msgs: &[serde_json::Value], memo: &str,
gas_limit: &str, sequence: &str,
) -> serde_json::Value {
let fee = self.calculate_fee(gas_limit);
serde_json::json!({
"body": {
"messages": msgs,
"memo": memo,
"timeout_height": self.timeout_height.to_string(),
"extension_options": [],
"non_critical_extension_options": [],
},
"auth_info": {
"signer_infos": [{
"public_key": { "@type": "/msgchain.crypto.DilithiumPubKey" },
"mode_info": { "single": { "mode": "SIGN_MODE_DIRECT" } },
"sequence": sequence,
}],
"fee": { "amount": fee, "gas_limit": gas_limit, "payer": "", "granter": "" },
},
"signatures": [],
})
}
fn calculate_fee(&self, gas_limit: &str) -> Vec<Coin> {
let gas: u64 = gas_limit.parse().unwrap_or(200_000);
let fee_umsg = (self.gas_price * Decimal::from(gas)).round_dp(0).to_string();
vec![Coin::new("umsg", fee_umsg)]
}
/// 等待交易确认
pub async fn wait_for_tx(
&self, tx_hash: &TxHash, timeout: Duration,
) -> Result<TxResponse, SdkError> {
let start = std::time::Instant::now();
loop {
if start.elapsed() > timeout {
return Err(SdkError::Timeout(timeout));
}
match self.query_tx(tx_hash).await {
Ok(resp) => {
if resp.code != 0 {
return Err(SdkError::TxFailed {
tx_hash: tx_hash.0.clone(),
code: resp.code, log: resp.raw_log,
});
}
return Ok(resp);
}
Err(SdkError::TxNotFound(_)) => {
tokio::time::sleep(Duration::from_millis(500)).await;
continue;
}
Err(e) => return Err(e),
}
}
}
/// 查询交易
pub async fn query_tx(&self, tx_hash: &TxHash) -> Result<TxResponse, SdkError> {
#[derive(serde::Deserialize)]
struct TxResp { tx_response: TxResponse }
let path = format!("/cosmos/tx/v1beta1/txs/{}", tx_hash.0);
let resp: TxResp = self.client.get(&path, None).await?;
Ok(resp.tx_response)
}
}
/// 交易响应
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct TxResponse {
pub height: String,
pub txhash: String,
pub code: u32,
pub data: Option<String>,
pub raw_log: String,
pub logs: Vec<TxLog>,
pub info: String,
pub gas_wanted: String,
pub gas_used: String,
pub tx: Option<serde_json::Value>,
pub timestamp: Option<String>,
pub events: Vec<Event>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct TxLog {
pub msg_index: u32,
pub log: String,
pub events: Vec<Event>,
}
/// 投票选项
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum VoteOption {
Unspecified = 0,
Yes = 1,
Abstain = 2,
No = 3,
NoWithVeto = 4,
}
impl VoteOption {
pub fn as_str(&self) -> &'static str {
match self {
VoteOption::Yes => "VOTE_OPTION_YES",
VoteOption::Abstain => "VOTE_OPTION_ABSTAIN",
VoteOption::No => "VOTE_OPTION_NO",
VoteOption::NoWithVeto => "VOTE_OPTION_NO_WITH_VETO",
VoteOption::Unspecified => "VOTE_OPTION_UNSPECIFIED",
}
}
}
8.2 合约部署
// src/tx/deploy.rs
impl TxBuilder {
/// 上传合约代码
pub async fn store_code(
&self, wallet: &Wallet, wasm_bytes: &[u8],
instantiate_permission: Option<AccessConfig>,
) -> Result<StoreCodeResponse, SdkError> {
let wasm_b64 = base64::encode(wasm_bytes);
let msg = serde_json::json!({
"@type": "/cosmwasm.wasm.v1.MsgStoreCode",
"sender": wallet.address.bech32,
"wasm_byte_code": wasm_b64,
"instantiate_permission": instantiate_permission,
});
let resp = self.build_and_broadcast(wallet, vec![msg], None).await?;
let code_id = resp.logs.iter()
.flat_map(|log| &log.events)
.flat_map(|evt| &evt.attributes)
.find(|attr| attr.key == "code_id")
.map(|attr| attr.value.clone())
.ok_or_else(|| SdkError::Unknown("无法解析 code_id".into()))?;
Ok(StoreCodeResponse {
tx_hash: TxHash(resp.txhash),
code_id: code_id.parse::<u64>()
.map_err(|e| SdkError::InvalidInput(e.to_string()))?,
})
}
/// 完整部署流程:上传 + 实例化
pub async fn deploy_contract(
&self, wallet: &Wallet, wasm_bytes: &[u8],
instantiate_msg: &impl Serialize, label: &str,
admin: Option<&Bech32Address>,
) -> Result<DeployResponse, SdkError> {
// 上传
let store = self.store_code(wallet, wasm_bytes, None).await?;
tracing::info!("Code stored: code_id={}", store.code_id);
// 实例化
let tx_resp = self.instantiate_contract(
wallet, store.code_id, instantiate_msg, label, admin, vec![],
).await?;
// 解析合约地址
let contract_addr = tx_resp.logs.iter()
.flat_map(|log| &log.events)
.flat_map(|evt| &evt.attributes)
.find(|attr| attr.key == "_contract_address")
.map(|attr| attr.value.clone())
.ok_or_else(|| SdkError::Unknown("无法解析合约地址".into()))?;
Ok(DeployResponse {
tx_hash: TxHash(tx_resp.txhash),
code_id: store.code_id,
contract_address: Bech32Address::from_bech32(&contract_addr)?,
height: tx_resp.height.parse::<u64>().unwrap_or(0),
})
}
}
#[derive(Debug, Clone)]
pub struct StoreCodeResponse {
pub tx_hash: TxHash,
pub code_id: u64,
}
#[derive(Debug, Clone)]
pub struct DeployResponse {
pub tx_hash: TxHash,
pub code_id: u64,
pub contract_address: Bech32Address,
pub height: u64,
}
9. Agent API
9.1 完整客户端
// src/clients/agent/mod.rs
use futures_util::{SinkExt, StreamExt};
use tokio_tungstenite::{connect_async, tungstenite::Message};
/// Agent API 客户端
///
/// 18 个端点,分为三类:
/// - 公开 (9): account, balance, balances, tx, block, blocks/range,
/// events/history, events/subscribe, oracle/price
/// - 保护 (6): wallet, mpc/sign, payment/session
/// - Stub (3): defi/*, registry/*, bridge/*
///
/// 限流: 100 req/s, burst 200
#[derive(Debug, Clone)]
pub struct AgentClient {
client: MSGChainClient,
}
impl AgentClient {
pub fn new(client: MSGChainClient) -> Self { Self { client } }
fn agent_path(&self, endpoint: &str) -> String {
format!("/api/v1/agent/{}", endpoint)
}
// ========================================
// 公开端点
// ========================================
/// 1. 账户查询
/// GET /api/v1/agent/account?address={address}
pub async fn account(&self, address: &Bech32Address) -> Result<AgentAccountInfo, SdkError> {
self.client.get(&self.agent_path("account"),
Some(&[("address", &address.bech32)])).await
}
/// 2. 余额查询
/// GET /api/v1/agent/balance?address={address}&denom={denom}
pub async fn balance(
&self, address: &Bech32Address, denom: &str,
) -> Result<AgentBalance, SdkError> {
self.client.get(&self.agent_path("balance"),
Some(&[("address", &address.bech32), ("denom", denom)])).await
}
/// 3. 所有余额
/// GET /api/v1/agent/balances?address={address}
pub async fn balances(&self, address: &Bech32Address) -> Result<AgentBalances, SdkError> {
self.client.get(&self.agent_path("balances"),
Some(&[("address", &address.bech32)])).await
}
/// 4. 交易查询
/// GET /api/v1/agent/tx?hash={hash}
pub async fn tx(&self, hash: &TxHash) -> Result<AgentTxInfo, SdkError> {
self.client.get(&self.agent_path("tx"),
Some(&[("hash", &hash.0)])).await
}
/// 5. 区块查询
/// GET /api/v1/agent/block?height={height}
pub async fn block(&self, height: Option<u64>) -> Result<AgentBlockInfo, SdkError> {
let mut params = Vec::new();
if let Some(h) = height { params.push(("height", h.to_string())); }
self.client.get(&self.agent_path("block"), Some(¶ms)).await
}
/// 6. 区块范围
/// GET /api/v1/agent/blocks/range?start={start}&end={end}
pub async fn blocks_range(
&self, start: u64, end: u64,
) -> Result<Vec<AgentBlockInfo>, SdkError> {
self.client.get(&self.agent_path("blocks/range"),
Some(&[("start", &start.to_string()), ("end", &end.to_string())])).await
}
/// 7. 事件历史
/// GET /api/v1/agent/events/history?filter={filter}&from={from}&to={to}
pub async fn events_history(
&self, filter: Option<&str>, from_block: Option<u64>, to_block: Option<u64>,
) -> Result<Vec<AgentEvent>, SdkError> {
let mut params = Vec::new();
if let Some(f) = filter { params.push(("filter", f)); }
if let Some(fb) = from_block { params.push(("from", fb.to_string())); }
if let Some(tb) = to_block { params.push(("to", tb.to_string())); }
self.client.get(&self.agent_path("events/history"), Some(¶ms)).await
}
/// 8. 事件订阅 (WebSocket)
/// WS /api/v1/agent/events/subscribe?filter={filter}
pub async fn events_subscribe(
&self, filter: Option<&str>,
) -> Result<EventStream, SdkError> {
let ws_base = self.client.ws_url().to_string().trim_end_matches('/').to_string();
let mut path = format!("{}/api/v1/agent/events/subscribe", ws_base);
if let Some(f) = filter {
path.push_str(&format!("?filter={}",
url::form_urlencoded::byte_serialize(f.as_bytes())));
}
let (ws_stream, _) = connect_async(&path).await
.map_err(|e| SdkError::WsError(format!("WS 连接失败: {}", e)))?;
Ok(EventStream { inner: ws_stream })
}
/// 9. 预言机价格
/// GET /api/v1/agent/oracle/price?symbol={symbol}
pub async fn oracle_price(&self, symbol: &str) -> Result<OraclePrice, SdkError> {
self.client.get(&self.agent_path("oracle/price"),
Some(&[("symbol", symbol)])).await
}
// ========================================
// 保护端点 (需 API Key)
// ========================================
/// 10. 钱包创建/导入
/// POST /api/v1/agent/wallet
pub async fn wallet_create(
&self, request: &WalletCreateRequest,
) -> Result<WalletInfo, SdkError> {
self.client.post_stub(&self.agent_path("wallet"), request).await
}
/// 11. MPC 签名
/// POST /api/v1/agent/mpc/sign
pub async fn mpc_sign(
&self, request: &MpcSignRequest,
) -> Result<MpcSignResponse, SdkError> {
self.client.post_stub(&self.agent_path("mpc/sign"), request).await
}
/// 12. 支付会话
/// POST /api/v1/agent/payment/session
pub async fn payment_session_create(
&self, request: &PaymentSessionRequest,
) -> Result<PaymentSession, SdkError> {
self.client.post_stub(&self.agent_path("payment/session"), request).await
}
// ========================================
// Stub 端点 (X-MSG-Stub=true)
// ========================================
/// 13. DeFi
/// POST /api/v1/agent/defi/{protocol}/{action}
pub async fn defi_interact(
&self, protocol: &str, action: &str, request: &DeFiRequest,
) -> Result<serde_json::Value, SdkError> {
self.client.post_stub(
&self.agent_path(&format!("defi/{}/{}", protocol, action)), request
).await
}
/// 14. 注册表
/// POST /api/v1/agent/registry/{entity}
pub async fn registry_query(
&self, entity: &str, request: &RegistryRequest,
) -> Result<serde_json::Value, SdkError> {
self.client.post_stub(
&self.agent_path(&format!("registry/{}", entity)), request
).await
}
/// 15. 桥接
/// POST /api/v1/agent/bridge/{chain}/{action}
pub async fn bridge_operate(
&self, chain: &str, action: &str, request: &BridgeRequest,
) -> Result<serde_json::Value, SdkError> {
self.client.post_stub(
&self.agent_path(&format!("bridge/{}/{}", chain, action)), request
).await
}
}
// ========================================
// WebSocket 事件流
// ========================================
#[derive(Debug)]
pub struct EventStream {
inner: tokio_tungstenite::WebSocketStream<
tokio_tungstenite::MaybeTlsStream<tokio::net::TcpStream>>,
}
impl EventStream {
pub async fn next_event(&mut self) -> Result<Option<AgentEvent>, SdkError> {
loop {
match self.inner.next().await {
Some(Ok(Message::Text(text))) => {
return Ok(Some(serde_json::from_str(&text)?));
}
Some(Ok(Message::Ping(_))) => {
self.inner.send(Message::Pong(vec![])).await
.map_err(|e| SdkError::WsError(e.to_string()))?;
continue;
}
Some(Ok(Message::Close(_))) => return Ok(None),
Some(Err(e)) => return Err(SdkError::WsError(e.to_string())),
None => return Ok(None),
_ => continue,
}
}
}
pub async fn close(mut self) -> Result<(), SdkError> {
self.inner.close(None).await
.map_err(|e| SdkError::WsError(e.to_string()))
}
}
// ========================================
// Agent API 类型定义
// ========================================
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct AgentAccountInfo {
pub address: String,
pub balances: Vec<AgentBalanceItem>,
pub account_number: u64,
pub sequence: u64,
pub pub_key_type: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct AgentBalance {
pub address: String,
pub denom: String,
pub amount: String,
pub amount_human: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct AgentBalanceItem {
pub denom: String,
pub amount: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct AgentBalances {
pub address: String,
pub balances: Vec<AgentBalanceItem>,
pub total_usd: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct AgentTxInfo {
pub hash: String,
pub height: u64,
pub index: u32,
pub from: String,
pub to: Option<String>,
pub amount: Option<Vec<Coin>>,
pub gas_used: u64,
pub gas_wanted: u64,
pub fee: Vec<Coin>,
pub memo: String,
pub status: String,
pub timestamp: String,
pub events: Vec<AgentEvent>,
pub messages: Vec<serde_json::Value>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct AgentBlockInfo {
pub height: u64,
pub hash: String,
pub time: String,
pub proposer: String,
pub tx_count: u32,
pub total_gas: u64,
pub size: u64,
pub last_commit_hash: String,
pub data_hash: String,
pub validators_hash: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct AgentEvent {
#[serde(rename = "type")]
pub event_type: String,
pub attributes: Vec<AgentEventAttribute>,
pub height: Option<u64>,
pub tx_hash: Option<String>,
pub timestamp: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct AgentEventAttribute {
pub key: String,
pub value: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct OraclePrice {
pub symbol: String,
pub price: String,
pub price_usd: String,
pub timestamp: String,
pub source: String,
pub confidence: Option<f64>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct WalletCreateRequest {
pub mnemonic: Option<String>,
pub derivation_path: Option<String>,
pub key_type: WalletKeyType,
pub password: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum WalletKeyType { Dilithium5, Ed25519, Secp256k1 }
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct WalletInfo {
pub address: String,
pub public_key: String,
pub key_type: String,
pub mnemonic: Option<String>,
pub derivation_path: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct MpcSignRequest {
pub wallet_id: String,
pub tx_bytes: String,
pub sign_mode: SignMode,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum SignMode { Direct, AminoJson, Textual }
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct MpcSignResponse {
pub signature: String,
pub public_key: String,
pub signature_type: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct PaymentSessionRequest {
pub agent_id: String,
pub amount: String,
pub denom: String,
pub max_spend: Option<String>,
pub expiry_seconds: Option<u64>,
pub purpose: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct PaymentSession {
pub session_id: String,
pub agent_id: String,
pub budget: String,
pub denom: String,
pub spent: String,
pub remaining: String,
pub expires_at: String,
pub status: SessionStatus,
pub purpose: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum SessionStatus { Active, Paused, Exhausted, Expired, Cancelled }
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct DeFiRequest {
pub user: String,
pub amount: String,
pub denom: String,
pub protocol: String,
pub action: String,
pub params: Option<serde_json::Value>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct RegistryRequest {
pub entity_type: String,
pub entity_id: Option<String>,
pub query: Option<serde_json::Value>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct BridgeRequest {
pub sender: String,
pub receiver: String,
pub amount: String,
pub denom: String,
pub source_chain: String,
pub target_chain: String,
pub memo: Option<String>,
}
9.2 批量操作
// src/clients/agent/batch.rs
use futures_util::future::join_all;
impl AgentClient {
/// 批量查询余额
pub async fn batch_balances(
&self, addresses: &[Bech32Address],
) -> Result<Vec<AgentBalances>, SdkError> {
let futures: Vec<_> = addresses.iter().map(|a| self.balances(a)).collect();
join_all(futures).await.into_iter().collect::<Result<Vec<_>, _>>()
}
/// 批量查询交易
pub async fn batch_tx_info(
&self, hashes: &[TxHash],
) -> Result<Vec<AgentTxInfo>, SdkError> {
let futures: Vec<_> = hashes.iter().map(|h| self.tx(h)).collect();
join_all(futures).await.into_iter().collect::<Result<Vec<_>, _>>()
}
/// 批量查询价格
pub async fn batch_oracle_prices(
&self, symbols: &[&str],
) -> Result<Vec<OraclePrice>, SdkError> {
let futures: Vec<_> = symbols.iter().map(|s| self.oracle_price(s)).collect();
join_all(futures).await.into_iter().collect::<Result<Vec<_>, _>>()
}
}
9.3 Agent 错误码
// src/clients/agent/errors.rs
/// Agent API 错误码枚举
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum AgentErrorCode {
BadRequest = 40000,
Unauthorized = 40100,
Forbidden = 40300,
NotFound = 40400,
RateLimited = 42900,
InternalError = 50000,
WalletCreateFailed = 41001,
WalletNotFound = 41002,
InsufficientFunds = 41003,
MpcSignFailed = 42001,
MpcKeyNotFound = 42002,
MpcSessionExpired = 42003,
PaymentSessionExpired = 43001,
PaymentLimitExceeded = 43002,
PaymentUnauthorized = 43003,
OracleSymbolNotFound = 44001,
OraclePriceStale = 44002,
StubNotImplemented = 49001,
}
impl AgentErrorCode {
pub fn from_code(code: u32) -> Option<Self> {
use AgentErrorCode::*;
Some(match code {
40000 => BadRequest, 40100 => Unauthorized, 40300 => Forbidden,
40400 => NotFound, 42900 => RateLimited, 50000 => InternalError,
41001 => WalletCreateFailed, 41002 => WalletNotFound, 41003 => InsufficientFunds,
42001 => MpcSignFailed, 42002 => MpcKeyNotFound, 42003 => MpcSessionExpired,
43001 => PaymentSessionExpired, 43002 => PaymentLimitExceeded,
43003 => PaymentUnauthorized, 44001 => OracleSymbolNotFound,
44002 => OraclePriceStale, 49001 => StubNotImplemented,
_ => return None,
})
}
pub fn description(&self) -> &'static str {
use AgentErrorCode::*;
match self {
BadRequest => "请求参数错误",
Unauthorized => "API Key 无效或缺失",
Forbidden => "无权限访问此端点",
NotFound => "资源不存在",
RateLimited => "请求频率超限",
InternalError => "服务端内部错误",
WalletCreateFailed => "钱包创建失败",
WalletNotFound => "钱包未找到",
InsufficientFunds => "余额不足",
MpcSignFailed => "MPC 签名失败",
MpcKeyNotFound => "MPC 密钥未找到",
MpcSessionExpired => "MPC 会话过期",
PaymentSessionExpired => "支付会话过期",
PaymentLimitExceeded => "支付限制超限",
PaymentUnauthorized => "支付未授权",
OracleSymbolNotFound => "价格符号不存在",
OraclePriceStale => "价格数据过期",
StubNotImplemented => "Stub 端点未实现",
}
}
}
10. 治理
// src/clients/gov.rs
/// 治理客户端
#[derive(Debug, Clone)]
pub struct GovernanceClient {
client: MSGChainClient,
}
impl GovernanceClient {
pub fn new(client: MSGChainClient) -> Self { Self { client } }
/// 查询所有提案
pub async fn proposals(
&self, status: Option<ProposalStatus>,
) -> Result<Vec<Proposal>, SdkError> {
#[derive(serde::Deserialize)]
struct ProposalsResp { proposals: Vec<Proposal>, pagination: Pagination }
let mut params = Vec::new();
if let Some(s) = status {
params.push(("proposal_status", (s as u32).to_string()));
}
let resp: ProposalsResp = self.client
.get("/cosmos/gov/v1beta1/proposals", Some(¶ms)).await?;
Ok(resp.proposals)
}
/// 提案详情
pub async fn proposal(&self, proposal_id: u64) -> Result<Proposal, SdkError> {
#[derive(serde::Deserialize)]
struct ProposalResp { proposal: Proposal }
let path = format!("/cosmos/gov/v1beta1/proposals/{}", proposal_id);
let resp: ProposalResp = self.client.get(&path, None).await?;
Ok(resp.proposal)
}
/// 提案投票
pub async fn proposal_votes(
&self, proposal_id: u64,
) -> Result<Vec<Vote>, SdkError> {
#[derive(serde::Deserialize)]
struct VotesResp { votes: Vec<Vote> }
let path = format!("/cosmos/gov/v1beta1/proposals/{}/votes", proposal_id);
let resp: VotesResp = self.client.get(&path, None).await?;
Ok(resp.votes)
}
/// 提案质押
pub async fn proposal_deposits(
&self, proposal_id: u64,
) -> Result<Vec<Deposit>, SdkError> {
#[derive(serde::Deserialize)]
struct DepositsResp { deposits: Vec<Deposit> }
let path = format!("/cosmos/gov/v1beta1/proposals/{}/deposits", proposal_id);
let resp: DepositsResp = self.client.get(&path, None).await?;
Ok(resp.deposits)
}
/// 治理参数
pub async fn gov_params(
&self, param_type: GovParamType,
) -> Result<serde_json::Value, SdkError> {
let path = format!("/cosmos/gov/v1beta1/params/{}", param_type.as_str());
#[derive(serde::Deserialize)]
struct ParamResp { #[serde(flatten)] data: serde_json::Value }
let resp: ParamResp = self.client.get(&path, None).await?;
Ok(resp.data)
}
/// 投票统计
pub async fn tally(&self, proposal_id: u64) -> Result<TallyResult, SdkError> {
#[derive(serde::Deserialize)]
struct TallyResp { tally: TallyResult }
let path = format!("/cosmos/gov/v1beta1/proposals/{}/tally", proposal_id);
let resp: TallyResp = self.client.get(&path, None).await?;
Ok(resp.tally)
}
}
/// 提案状态
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub enum ProposalStatus {
Unspecified = 0, DepositPeriod = 1, VotingPeriod = 2,
Passed = 3, Rejected = 4, Failed = 5,
}
/// 提案
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Proposal {
pub proposal_id: String,
pub content: ProposalContent,
pub status: ProposalStatus,
pub final_tally_result: Option<TallyResult>,
pub submit_time: String,
pub deposit_end_time: String,
pub total_deposit: Vec<Coin>,
pub voting_start_time: Option<String>,
pub voting_end_time: Option<String>,
}
/// 提案内容枚举
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "@type")]
pub enum ProposalContent {
#[serde(rename = "/cosmos.gov.v1beta1.TextProposal")]
Text(TextProposal),
#[serde(rename = "/cosmos.gov.v1beta1.ParameterChangeProposal")]
ParameterChange(ParameterChangeProposal),
#[serde(rename = "/cosmwasm.wasm.v1.StoreCodeProposal")]
StoreCode(StoreCodeProposal),
#[serde(rename = "/cosmwasm.wasm.v1.InstantiateContractProposal")]
InstantiateContract(InstantiateContractProposal),
#[serde(rename = "/cosmwasm.wasm.v1.MigrateContractProposal")]
MigrateContract(MigrateContractProposal),
#[serde(rename = "/cosmos.upgrade.v1beta1.SoftwareUpgradeProposal")]
SoftwareUpgrade(SoftwareUpgradeProposal),
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct TextProposal { pub title: String, pub description: String }
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct ParameterChangeProposal {
pub title: String, pub description: String, pub changes: Vec<ParamChange>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct ParamChange {
pub subspace: String, pub key: String, pub value: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct StoreCodeProposal {
pub title: String, pub description: String, pub run_as: String,
pub wasm_byte_code: String, pub instantiate_permission: Option<AccessConfig>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct InstantiateContractProposal {
pub title: String, pub description: String, pub run_as: String,
pub admin: Option<String>, pub code_id: String, pub label: String,
pub msg: String, pub funds: Vec<Coin>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct MigrateContractProposal {
pub title: String, pub description: String,
pub contract: String, pub code_id: String, pub msg: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct SoftwareUpgradeProposal {
pub title: String, pub description: String, pub plan: UpgradePlan,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct UpgradePlan {
pub name: String, pub height: String, pub info: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct TallyResult {
pub yes: String, pub abstain: String, pub no: String, pub no_with_veto: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Vote {
pub proposal_id: String, pub voter: String,
pub option: VoteOption, pub options: Vec<WeightedVoteOption>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct WeightedVoteOption {
pub option: VoteOption, pub weight: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Deposit {
pub proposal_id: String, pub depositor: String, pub amount: Vec<Coin>,
}
#[derive(Debug, Clone, Copy)]
pub enum GovParamType { Deposit, Tallying, Voting }
impl GovParamType {
pub fn as_str(&self) -> &'static str {
match self { GovParamType::Deposit => "deposit", GovParamType::Tallying => "tallying", GovParamType::Voting => "voting" }
}
}
11. 完整示例
11.1 Agent 余额监控机器人
// examples/agent_bot.rs
//! Agent 余额监控机器人
//!
//! 功能:
//! - 每 10 秒查询 MSG 余额
//! - 低于阈值告警
//! - 同步监控区块和价格
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
tracing_subscriber::fmt::init();
let client = MSGChainClient::from_env()?;
let agent = client.agent_client();
let bank = client.bank_client();
let monitor_address = Bech32Address::from_bech32(
"msg1qypqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqz8n7z"
).unwrap_or_else(|_| {
let wallet = Wallet::from_mnemonic(
&std::env::var("MSG_WALLET_MNEMONIC")
.expect("需要 MSG_WALLET_MNEMONIC")
).unwrap();
wallet.address
});
info!("监控地址: {}", monitor_address);
let alert_threshold = Amount::from_umsg("1000000000000000000"); // 1 MSG
loop {
// 1. 余额
match bank.msg_balance(&monitor_address).await {
Ok(bal) => {
info!("余额: {} MSG", bal.0);
if bal.to_umsg() < alert_threshold.to_umsg() {
warn!("余额过低! {}", bal.0);
}
}
Err(e) => error!("查询余额失败: {}", e),
}
// 2. 最新区块
match agent.block(None).await {
Ok(block) => info!("区块: height={}, tx_count={}", block.height, block.tx_count),
Err(e) => error!("区块查询失败: {}", e),
}
// 3. 价格
match agent.oracle_price("MSG").await {
Ok(price) => info!("MSG 价格: ${}", price.price_usd),
Err(e) => error!("价格查询失败: {}", e),
}
tokio::time::sleep(Duration::from_secs(10)).await;
}
}
11.2 合约部署工具
// examples/contract_deployer.rs
//! 合约部署工具
//!
//! 1. 读取 Wasm 文件
//! 2. 上传 + 实例化
//! 3. 输出合约地址
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
tracing_subscriber::fmt::init();
let client = MSGChainClient::from_env()?;
let wallet = Wallet::from_mnemonic(&std::env::var("MSG_WALLET_MNEMONIC")?)?;
let tx_builder = client.tx_builder();
let query = client.contract_query_client();
info!("部署钱包: {}", wallet.address);
let wasm_path = Path::new(
"target/wasm32-unknown-unknown/release/ai_model_registry.wasm"
);
let wasm_bytes = std::fs::read(wasm_path).expect(
"请先编译: cd contracts/cosmwasm/all/ai_model_registry && \
cargo build --target wasm32-unknown-unknown --release"
);
info!("Wasm: {} ({} bytes)", wasm_path.display(), wasm_bytes.len());
let instantiate_msg = serde_json::json!({
"owner": wallet.address.bech32,
"name": "MSG AI Model Registry",
"version": "1.0.0",
});
let deploy = tx_builder.deploy_contract(
&wallet, &wasm_bytes, &instantiate_msg,
"AI Model Registry v1", Some(&wallet.address),
).await?;
info!("部署成功!");
info!(" tx_hash: {}", deploy.tx_hash);
info!(" code_id: {}", deploy.code_id);
info!(" 合约地址: {}", deploy.contract_address);
info!(" 高度: {}", deploy.height);
let info = query.contract_info(&deploy.contract_address).await?;
info!("合约信息: code_id={}, creator={}, label={}",
info.code_id, info.creator, info.label);
let tx = tx_builder.wait_for_tx(
&deploy.tx_hash, Duration::from_secs(30)
).await?;
info!("确认: height={}, gas_used={}", tx.height, tx.gas_used);
Ok(())
}
11.3 实时事件监控
// examples/balance_monitor.rs
//! 实时余额监控器
//!
//! - 通过 WebSocket 订阅转账事件
//! - 实时更新余额
//! - 大额告警
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
tracing_subscriber::fmt::init();
let client = MSGChainClient::from_env()?;
let agent = client.agent_client();
let bank = client.bank_client();
let watch_addresses: Vec<Bech32Address> = vec![
Bech32Address::from_bech32("msg1...")?,
Bech32Address::from_bech32("msg1...")?,
];
let mut balance_cache: HashMap<String, Amount> = HashMap::new();
// 初始加载
for addr in &watch_addresses {
if let Ok(bal) = bank.msg_balance(addr).await {
balance_cache.insert(addr.bech32.clone(), bal.clone());
info!("初始 {}: {} MSG", addr, bal.0);
}
}
info!("订阅转账事件...");
let mut stream = agent.events_subscribe(Some("transfer")).await?;
while let Some(event) = stream.next_event().await? {
let from = event.attributes.iter()
.find(|a| a.key == "from").map(|a| a.value.clone());
let to = event.attributes.iter()
.find(|a| a.key == "to").map(|a| a.value.clone());
let amount = event.attributes.iter()
.find(|a| a.key == "amount").map(|a| a.value.clone());
let involved = watch_addresses.iter().any(|addr| {
from.as_ref().map_or(false, |f| f == &addr.bech32)
|| to.as_ref().map_or(false, |t| t == &addr.bech32)
});
if !involved { continue; }
info!("转账: from={:?} to={:?} amount={:?} tx={:?}",
from, to, amount, event.tx_hash);
// 大额告警
if let Some(ref amt) = amount {
if let Some(umsg) = amt.strip_suffix("umsg") {
let val: f64 = umsg.parse().unwrap_or(0.0);
if val > 100.0 * 1e18 {
warn!("大额交易! {} umsg", umsg);
}
}
}
// 更新余额
for addr in &watch_addresses {
let a = &addr.bech32;
if from.as_ref().map_or(false, |f| f == a)
|| to.as_ref().map_or(false, |t| t == a)
{
if let Ok(new) = bank.msg_balance(addr).await {
let old = balance_cache.get(a).cloned()
.unwrap_or_else(|| Amount::from_umsg("0"));
balance_cache.insert(a.clone(), new.clone());
info!("余额更新 {}: {} -> {} MSG", addr, old.0, new.0);
}
}
}
}
Ok(())
}
11.4 批量 AI 合约查询
// examples/batch_ai_contracts.rs
//! 同时查询 5 个 AI 原生合约
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let client = MSGChainClient::from_env()?;
let query = client.contract_query_client();
let model_registry = Bech32Address::from_bech32("msg1model...")?;
let agent_registry = Bech32Address::from_bech32("msg1agent...")?;
let market = Bech32Address::from_bech32("msg1market...")?;
let data_assets = Bech32Address::from_bech32("msg1data...")?;
let proof = Bech32Address::from_bech32("msg1proof...")?;
let (models, agent_info, stats, _asset, _proof) = tokio::try_join!(
query.ai_list_models(&model_registry, None, Some(10)),
query.ai_agent_info(&agent_registry, "agent-001"),
query.ai_market_stats(&market),
query.ai_data_asset(&data_assets, "asset-001"),
query.ai_inference_proof(&proof, "proof-001"),
)?;
info!("模型数: {}", models.len());
info!("Agent: name={}, status={:?}", agent_info.name, agent_info.status);
info!("市场: requests={}, agents={}", stats.total_requests, stats.total_agents);
Ok(())
}
11.5 编译与运行
# 编译
cd contracts/cosmwasm/all/ai_model_registry && cargo build --target wasm32-unknown-unknown --release
# 运行示例
cargo run --example agent_bot
cargo run --example contract_deployer
cargo run --example balance_monitor
# 设置环境变量
export MSG_CHAIN_RPC_URL=https://rpc.msgchain.org
export MSG_CHAIN_REST_URL=https://rest.msgchain.org
export MSG_CHAIN_WS_URL=wss://ws.msgchain.org
export MSG_WALLET_MNEMONIC="your 24 word mnemonic here"
export MSG_AGENT_API_KEY=your_api_key
# 运行测试
cargo test
12. 测试
12.1 单元测试
// tests/unit_tests.rs
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_amount_conversion() {
let a = Amount::from_umsg("1000000000000000000");
assert_eq!(a.0, "1.0");
let a = Amount::from_umsg("0");
assert_eq!(a.0, "0.0");
let a = Amount::from_umsg("123456789012345678901234567890");
assert_eq!(a.to_umsg(), "123456789012345678901234567890");
}
#[test]
fn test_tx_hash_validation() {
let valid = "a1b2c3d4e5f6a7b8c9d0e1f2a3b4c5d6e7f8a9b0c1d2e3f4a5b6c7d8e9f0a1b2";
assert!(TxHash::from_str(valid).is_ok());
let invalid = "short";
assert!(TxHash::from_str(invalid).is_err());
}
#[test]
fn test_gas_level_pricing() {
assert_eq!(GasLevel::Low.price().to_string(), "1000000000");
assert_eq!(GasLevel::Average.price().to_string(), "1000000000");
assert_eq!(GasLevel::High.price().to_string(), "1000000000");
}
#[test]
fn test_vote_option_mapping() {
assert_eq!(VoteOption::Yes.as_str(), "VOTE_OPTION_YES");
assert_eq!(VoteOption::No.as_str(), "VOTE_OPTION_NO");
}
}
12.2 集成测试(Mock 服务器)
// tests/integration_test.rs
//! 使用 mockito mock 服务器进行集成测试
use msg_chain_sdk::prelude::*;
use mockito::{Matcher, Server};
#[tokio::test]
async fn test_bank_balance_mock() {
let mut server = Server::new_async().await;
let mock = server.mock("GET", "/cosmos/bank/v1beta1/balances/msg1test/")
.with_status(200)
.with_header("content-type", "application/json")
.with_body(r#"{
"balances": [{"denom": "umsg", "amount": "1000000000000000000"}],
"pagination": {"next_key": null, "total": "1"}
}"#)
.create_async().await;
let config = ChainConfig {
rest_url: url::Url::parse(&server.url()).unwrap(),
..Default::default()
};
let client = MSGChainClient::new(config).unwrap();
let bank = client.bank_client();
let addr = Bech32Address::from_bech32("msg1test").unwrap();
let balances = bank.balances(&addr).await.unwrap();
assert_eq!(balances.len(), 1);
assert_eq!(balances[0].denom, "umsg");
mock.assert_async().await;
}
#[tokio::test]
async fn test_agent_oracle_price_mock() {
let mut server = Server::new_async().await;
let mock = server.mock("GET", Matcher::Regex(r"^/api/v1/agent/oracle/price.*$"))
.with_status(200)
.with_body(r#"{
"symbol": "MSG",
"price": "0.52",
"price_usd": "0.52",
"timestamp": "2026-01-15T10:00:00Z",
"source": "coingecko",
"confidence": 0.95
}"#)
.create_async().await;
let config = ChainConfig {
rest_url: url::Url::parse(&server.url()).unwrap(),
..Default::default()
};
let client = MSGChainClient::new(config).unwrap();
let agent = client.agent_client();
let price = agent.oracle_price("MSG").await.unwrap();
assert_eq!(price.symbol, "MSG");
assert_eq!(price.price_usd, "0.52");
mock.assert_async().await;
}
#[tokio::test]
async fn test_stub_endpoint() {
let mut server = Server::new_async().await;
let mock = server.mock("POST", "/api/v1/agent/wallet")
.match_header("X-MSG-Stub", "true")
.with_status(200)
.with_body(r#"{
"address": "msg1stub...",
"public_key": "abc123",
"key_type": "Dilithium5",
"derivation_path": "m/44'/118'/0'/0/0"
}"#)
.create_async().await;
let req = WalletCreateRequest {
mnemonic: None,
derivation_path: None,
key_type: WalletKeyType::Dilithium5,
password: None,
};
let config = ChainConfig {
rest_url: url::Url::parse(&server.url()).unwrap(),
..Default::default()
};
let client = MSGChainClient::new(config).unwrap()
.with_agent_auth(AgentAuth {
api_key: "test-key".into(),
api_secret: None,
});
let agent = client.agent_client();
let wallet = agent.wallet_create(&req).await.unwrap();
assert_eq!(wallet.key_type, "Dilithium5");
mock.assert_async().await;
}
#[tokio::test]
async fn test_rate_limit_response() {
let mut server = Server::new_async().await;
let mock = server.mock("GET", "/cosmos/base/tendermint/v1beta1/node_info")
.with_status(429)
.with_header("Retry-After", "5")
.with_body("rate limited")
.create_async().await;
let config = ChainConfig {
rest_url: url::Url::parse(&server.url()).unwrap(),
..Default::default()
};
let client = MSGChainClient::new(config).unwrap();
let result = client.health_check().await;
assert!(result.is_err());
let err = result.unwrap_err();
assert!(err.is_rate_limited());
mock.assert_async().await;
}
#[tokio::test]
async fn test_tx_simulation_mock() {
let mut server = Server::new_async().await;
let mock = server.mock("POST", "/cosmos/tx/v1beta1/simulate")
.with_status(200)
.with_body(r#"{"gas_info": {"gas_used": "120000"}}"#)
.create_async().await;
let config = ChainConfig {
rest_url: url::Url::parse(&server.url()).unwrap(),
..Default::default()
};
let client = MSGChainClient::new(config).unwrap();
let tx_builder = client.tx_builder();
let msgs = vec![serde_json::json!({"@type": "/cosmos.bank.v1beta1.MsgSend"})];
let gas = tx_builder.estimate_gas(&msgs, None).await.unwrap();
// 120000 * 1.5 = 180000
assert_eq!(gas, "180000");
mock.assert_async().await;
}
#[tokio::test]
async fn test_events_subscribe_mock() {
// 此测试验证 WebSocket 连接逻辑使用正确的 URL
let client = MSGChainClient::from_env().unwrap_or_else(|_| {
MSGChainClient::new(ChainConfig {
ws_url: url::Url::parse("wss://ws.msgchain.org").unwrap(),
..Default::default()
}).unwrap()
});
let agent = client.agent_client();
// 验证 URL 构建不含语法错误
let ws_url = client.ws_url().to_string();
assert!(ws_url.starts_with("ws"));
}
12.3 性能基准测试
// benches/signing.rs
use criterion::{black_box, criterion_group, criterion_main, Criterion};
use msg_chain_sdk::crypto::keys::DilithiumKeypair;
fn bench_dilithium_sign(c: &mut Criterion) {
let kp = DilithiumKeypair::generate();
let msg = vec![0u8; 256];
c.bench_function("dilithium5_sign_256b", |b| {
b.iter(|| kp.sign(black_box(&msg)))
});
}
fn bench_dilithium_verify(c: &mut Criterion) {
let kp = DilithiumKeypair::generate();
let msg = vec![0u8; 256];
let sig = kp.sign(&msg);
c.bench_function("dilithium5_verify_256b", |b| {
b.iter(|| kp.verify(black_box(&msg), black_box(&sig)))
});
}
fn bench_address_generation(c: &mut Criterion) {
let kp = DilithiumKeypair::generate();
c.bench_function("bech32_address_from_pubkey", |b| {
b.iter(|| Bech32Address::from_public_key_bytes("msg", black_box(&kp.public_key)))
});
}
criterion_group!(benches, bench_dilithium_sign, bench_dilithium_verify, bench_address_generation);
criterion_main!(benches);
12.4 测试指导原则
// tests/README.md 内容概要
// MSG Chain Rust SDK 测试策略:
//
// 1. 单元测试 — 纯逻辑测试,不依赖网络
// - Amount 精度转换
// - TxHash 验证
// - 地址编解码
// - Dilithium 签名/验签
//
// 2. 集成测试 — 使用 mockito/wiremock
// - 模拟 REST 端点响应
// - 验证请求路径、查询参数、请求头(包括 X-MSG-Stub)
// - 验证限流重试逻辑
// - 验证错误处理分支
//
// 3. 端到端测试 — 需要连接真实或本地节点
// - 运行前设置 MSG 链本地开发环境
// - 测试真实交易广播
//
// 运行:
// cargo test # 所有测试
// cargo test --test integration_test # 仅集成测试
// cargo test unit_tests # 仅单元测试
// cargo bench # 基准测试
附录
A. 模块导出 (lib.rs)
// src/lib.rs
pub mod config;
pub mod error;
pub mod rate_limiter;
pub mod auth;
pub mod client;
pub mod types;
pub mod crypto {
pub mod address;
pub mod keys;
}
pub mod clients {
pub mod bank;
pub mod staking;
pub mod contract {
pub mod query;
pub mod ai_contracts;
}
pub mod gov;
pub mod agent {
pub mod r#mod;
pub mod batch;
pub mod errors;
}
}
pub mod tx {
pub mod builder;
pub mod deploy;
}
/// Prelude: 推荐的通导入方式
pub mod prelude {
pub use crate::client::MSGChainClient;
pub use crate::config::{ChainConfig, ConfigError, GasLevel};
pub use crate::error::{Result, SdkError};
pub use crate::types::*;
pub use crate::crypto::address::Bech32Address;
pub use crate::crypto::keys::{DilithiumKeypair, Ed25519Keypair, Wallet};
pub use crate::clients::bank::*;
pub use crate::clients::staking::*;
pub use crate::clients::contract::query::*;
pub use crate::clients::contract::ai_contracts::*;
pub use crate::clients::gov::*;
pub use crate::clients::agent::r#mod::*;
pub use crate::clients::agent::errors::*;
pub use crate::tx::builder::*;
pub use crate::tx::deploy::*;
pub use crate::auth::AgentAuth;
pub use crate::rate_limiter::RateLimiter;
}
B. CosmWasm 合约标准依赖
[dependencies]
cosmwasm-std = "1.5"
cw-storage-plus = "1.2"
cw2 = "1.1"
编译合约:
cd contracts/cosmwasm/all/{contract_name}
cargo build --target wasm32-unknown-unknown --release
cargo test
C. 常见问题
Q: 为什么使用 Dilithium-5 而不是 Ed25519?
A: MSG Chain 选择后量子密码学 Dilithium-5 以抵御量子计算攻击。Ed25519 以兼容模式提供。
Q: 如何估算 Gas?
A: TxBuilder 自动调用 /cosmos/tx/v1beta1/simulate 并乘以 gas_multiplier(默认 1.5x)。
Q: X-MSG-Stub 头的作用?
A: 标记当前请求为桩(stub)模式,用于尚在开发中的端点(defi/, registry/, bridge/*),服务端返回模拟数据。
Q: 如何处理限流?
A: SDK 内置令牌桶算法(100 req/s, burst 200),自动排队等待。收到 429 响应时读取 Retry-After 头。
Q: 支持 Windows 吗?
A: 完全支持。reqwest 使用 rustls-tls 无需 OpenSSL,tokio 跨平台。
D. 版本兼容性
| SDK 版本 | MSG Chain | CosmWasm | Rust 最低 |
|---|---|---|---|
| 1.0.x | v1.x | 1.5 | 1.75 |
E. 贡献指南
- Fork 仓库
cargo test确保所有测试通过- 添加新功能时附带测试
- 遵循现有代码风格
- PR 描述需包含变更说明
本文档基于 MSG Chain 代码库核实的技术事实。
白皮书系统: https://msgchain.org/whitepaper/
