MSG Chain Python SDK 开发指南
数据来源:MSG Chain 代码库核实
主网状态: No-Go — 当前 MSGChain 主网裁决为 No-Go,以下内容反映代码实际状态,不代表生产可用。
1. 概述
MSG Chain Python SDK 是一个用于与 MSG Chain 区块链交互的完整 Python 工具包。该 SDK 的设计参考了 web3.py 和 cosmpy 的架构模式,提供了面向对象、类型安全的接口,支持同步和异步两种编程模型。
核心特性
- 完整覆盖 MSG Chain 所有模块:银行、质押、合约、Agent API、治理、事件订阅
- 异步优先 + 同步兼容:所有客户端提供 async/await 和同步两种调用方式
- 类型安全:使用 Python dataclasses 和类型注解,IDE 友好
- 全面错误处理:18 种错误码映射为 Python 异常类,异常链完整
- 企业级密钥管理:BIP39 助记词、Bech32 地址、支持 Dilithium-5 和 Ed25519
- WebSocket 事件流:实时订阅链上事件和 Agent 事件
依赖
| 包名 | 版本 | 用途 |
|---|---|---|
| httpx | >=0.27.0 | HTTP 客户端(异步 + 同步) |
| websockets | >=12.0 | WebSocket 客户端 |
| bip-utils | >=2.9.0 | BIP39 助记词和 BIP32 派生 |
| bech32 | >=1.2.0 | Bech32 地址编码/解码 |
| pydantic | >=2.0.0 | 数据验证 |
| cryptography | >=41.0.0 | Ed25519 签名 |
| pycryptodome | >=3.20.0 | Dilithium-5 实验性支持 |
| grpcio | >=1.60.0 | gRPC 支持 |
架构图
msg_sdk/
├── __init__.py # 包入口,导出所有客户端
├── client.py # MSGChainClient 核心客户端
├── errors.py # 18 种错误码异常类
├── types.py # dataclass 响应类型
├── wallet.py # 钱包密钥管理
├── bank.py # 银行客户端
├── staking.py # 质押客户端
├── contract.py # 合约查询/执行客户端
├── agent.py # Agent API 客户端 (18 端点)
├── governance.py # 治理客户端
├── events.py # WebSocket 事件客户端
└── utils.py # 工具函数
2. 安装
pip 安装
pip install msg-sdk
从源码安装
git clone https://github.com/msgchain/msg-sdk.git
cd msg-sdk
pip install -e .
requirements.txt
httpx>=0.27.0
websockets>=12.0
bip-utils>=2.9.0
bech32>=1.2.0
pydantic>=2.0.0
cryptography>=41.0.0
pycryptodome>=3.20.0
grpcio>=1.60.0
grpcio-tools>=1.60.0
mnemonic>=0.20
pytest>=8.0.0
pytest-asyncio>=0.23.0
pytest-cov>=4.1.0
pyproject.toml
[build-system]
requires = ["setuptools>=68.0", "wheel"]
build-backend = "setuptools.build_meta"
[project]
name = "msg-sdk"
version = "1.0.0"
description = "MSG Chain Python SDK - 完整的 MSG Chain 区块链交互工具包"
readme = "README.md"
requires-python = ">=3.11"
license = {text = "MIT"}
authors = [
{name = "MSG Chain Dev", email = "dev@msgchain.org"},
]
classifiers = [
"Development Status :: 4 - Beta",
"Intended Audience :: Developers",
"License :: OSI Approved :: MIT License",
"Programming Language :: Python :: 3.11",
"Programming Language :: Python :: 3.12",
"Topic :: Blockchain",
"Topic :: Software Development :: Libraries :: Python Modules",
]
dependencies = [
"httpx>=0.27.0",
"websockets>=12.0",
"bip-utils>=2.9.0",
"bech32>=1.2.0",
"pydantic>=2.0.0",
"cryptography>=41.0.0",
"pycryptodome>=3.20.0",
"grpcio>=1.60.0",
"grpcio-tools>=1.60.0",
"mnemonic>=0.20",
]
[project.optional-dependencies]
dev = [
"pytest>=8.0.0",
"pytest-asyncio>=0.23.0",
"pytest-cov>=4.1.0",
"black>=24.0.0",
"ruff>=0.3.0",
"mypy>=1.8.0",
]
[tool.setuptools.packages.find]
where = ["src"]
include = ["msg_sdk*"]
[tool.pytest.ini_options]
asyncio_mode = "auto"
testpaths = ["tests"]
python_files = ["test_*.py"]
3. 错误码与异常体系
# 文件: msg_sdk/errors.py
# MSG Chain 错误码映射 -- 所有 18 种错误码对应 Python 异常类
from __future__ import annotations
from typing import Any, Dict, Optional
class MSGChainError(Exception):
code: int = -1
message: str = "unknown error"
def __init__(self, message: Optional[str] = None, data: Optional[Dict[str, Any]] = None) -> None:
self.message = message or self.message
self.data = data or {}
super().__init__(self.message)
def to_dict(self) -> Dict[str, Any]:
return {"code": self.code, "message": self.message, "data": self.data}
class OK(MSGChainError):
code = 0
message = "success"
class StubError(MSGChainError):
code = 1
message = "stub call"
class NotImplementedErrorCode(MSGChainError):
code = 2
message = "not implemented"
class UnauthorizedError(MSGChainError):
code = 3
message = "unauthorized"
class InsufficientFundsError(MSGChainError):
code = 4
message = "insufficient funds"
class ContractFailedError(MSGChainError):
code = 5
message = "contract execution failed"
class QueryTimeoutError(MSGChainError):
code = 6
message = "query timeout"
class InvalidParamError(MSGChainError):
code = 7
message = "invalid parameter"
class DAOTimelockError(MSGChainError):
code = 8
message = "DAO timelock"
class HighValueError(MSGChainError):
code = 9
message = "high value transaction requires approval"
class AgentNotFoundError(MSGChainError):
code = 10
message = "agent not found"
class AgentExistsError(MSGChainError):
code = 11
message = "agent already exists"
class ConstitutionViolationError(MSGChainError):
code = 12
message = "constitution violation"
class SessionExpiredError(MSGChainError):
code = 13
message = "session expired"
class SessionLimitError(MSGChainError):
code = 14
message = "session limit reached"
class InvalidSignatureError(MSGChainError):
code = 15
message = "invalid signature"
class DuplicateNonceError(MSGChainError):
code = 16
message = "duplicate nonce"
class RateLimitError(MSGChainError):
code = 17
message = "rate limit exceeded"
ERROR_CODE_MAP: Dict[int, type[MSGChainError]] = {
0: OK,
1: StubError,
2: NotImplementedErrorCode,
3: UnauthorizedError,
4: InsufficientFundsError,
5: ContractFailedError,
6: QueryTimeoutError,
7: InvalidParamError,
8: DAOTimelockError,
9: HighValueError,
10: AgentNotFoundError,
11: AgentExistsError,
12: ConstitutionViolationError,
13: SessionExpiredError,
14: SessionLimitError,
15: InvalidSignatureError,
16: DuplicateNonceError,
17: RateLimitError,
}
def raise_for_code(code: int, message: Optional[str] = None, data: Optional[Dict[str, Any]] = None) -> None:
if code == 0:
return
cls = ERROR_CODE_MAP.get(code, MSGChainError)
raise cls(message=message, data=data)
def extract_error(body: Dict[str, Any]) -> None:
code = body.get("code", 0)
if code == 0:
return
msg = body.get("message")
data = body.get("data")
raise_for_code(code, message=msg, data=data)
4. 数据类型定义
# 文件: msg_sdk/types.py
# dataclass 定义所有请求/响应模型
from __future__ import annotations
from dataclasses import dataclass, field, asdict
from decimal import Decimal
from typing import Any, Dict, List, Optional
@dataclass
class Coin:
denom: str
amount: str
def to_dict(self) -> Dict[str, str]:
return {"denom": self.denom, "amount": self.amount}
@classmethod
def from_dict(cls, data: Dict[str, str]) -> Coin:
return cls(denom=data["denom"], amount=data["amount"])
def to_umsg(self) -> Decimal:
return Decimal(self.amount)
@dataclass
class PageRequest:
key: Optional[str] = None
offset: Optional[int] = None
limit: Optional[int] = None
count_total: bool = False
reverse: bool = False
@dataclass
class PageResponse:
next_key: Optional[str] = None
total: Optional[int] = None
@dataclass
class BalanceResponse:
balances: List[Coin]
pagination: Optional[PageResponse] = None
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> BalanceResponse:
balances = [Coin.from_dict(c) for c in data.get("balances", [])]
return cls(balances=balances)
@dataclass
class SupplyResponse:
supply: List[Coin]
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> SupplyResponse:
supply = [Coin.from_dict(c) for c in data.get("supply", [])]
return cls(supply=supply)
@dataclass
class CommissionRates:
rate: str
max_rate: str
max_change_rate: str
@dataclass
class Description:
moniker: str
identity: str = ""
website: str = ""
security_contact: str = ""
details: str = ""
@dataclass
class Validator:
operator_address: str
consensus_pubkey: Dict[str, Any]
jailed: bool
status: int
tokens: str = "0"
delegator_shares: str = "0"
description: Optional[Description] = None
commission: Optional[CommissionRates] = None
unbonding_time: str = ""
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> Validator:
desc = data.get("description", {})
comm = data.get("commission", {}).get("commission_rates", {})
return cls(
operator_address=data["operator_address"],
consensus_pubkey=data.get("consensus_pubkey", {}),
jailed=data.get("jailed", False),
status=data.get("status", 0),
tokens=data.get("tokens", "0"),
delegator_shares=data.get("delegator_shares", "0"),
description=Description(
moniker=desc.get("moniker", ""),
identity=desc.get("identity", ""),
website=desc.get("website", ""),
security_contact=desc.get("security_contact", ""),
details=desc.get("details", ""),
),
commission=CommissionRates(
rate=comm.get("rate", "0"),
max_rate=comm.get("max_rate", "0"),
max_change_rate=comm.get("max_change_rate", "0"),
),
unbonding_time=data.get("unbonding_time", ""),
)
@dataclass
class DelegationResponse:
delegator_address: str
validator_address: str
shares: str = "0"
balance: Optional[Coin] = None
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> DelegationResponse:
d = data.get("delegation", {}) or {}
b = data.get("balance", {}) or {}
return cls(
delegator_address=d.get("delegator_address", ""),
validator_address=d.get("validator_address", ""),
shares=d.get("shares", "0"),
balance=Coin.from_dict(b) if b else Coin(denom="umsg", amount="0"),
)
@dataclass
class RewardEntry:
validator_address: str
reward: List[Coin]
@dataclass
class RewardsResponse:
rewards: List[RewardEntry]
total: List[Coin]
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> RewardsResponse:
rewards = [RewardEntry(validator_address=r.get("validator_address", ""),reward=[Coin.from_dict(c) for c in r.get("reward", [])]) for r in data.get("rewards", [])]
total = [Coin.from_dict(c) for c in data.get("total", [])]
return cls(rewards=rewards, total=total)
@dataclass
class ValidatorResponse:
validators: List[Validator]
pagination: Optional[PageResponse] = None
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> ValidatorResponse:
validators = [Validator.from_dict(v) for v in data.get("validators", [])]
return cls(validators=validators)
@dataclass
class CodeResponse:
code_info: Optional[Dict[str, Any]] = None
data: Optional[str] = None
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> CodeResponse:
return cls(code_info=data.get("code_info"), data=data.get("data"))
@dataclass
class ContractResponse:
contract_info: Optional[Dict[str, Any]] = None
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> ContractResponse:
return cls(contract_info=data.get("contract_info"))
@dataclass
class SmartQueryResponse:
data: Any
@classmethod
def from_dict(cls, data: Any) -> SmartQueryResponse:
return cls(data=data)
@dataclass
class AgentInfo:
id: str = ""
name: str = ""
description: str = ""
owner: str = ""
status: str = "active"
constitution_id: str = ""
created_at: str = ""
metadata: Dict[str, Any] = field(default_factory=dict)
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> AgentInfo:
return cls(id=str(data.get("id", "")),name=str(data.get("name", "")),description=str(data.get("description", "")),owner=str(data.get("owner", "")),status=str(data.get("status", "active")),constitution_id=str(data.get("constitution_id", "")),created_at=str(data.get("created_at", "")),metadata=data.get("metadata", {}))
@dataclass
class PaymentSession:
session_id: str = ""
payer: str = ""
payee: str = ""
amount: Optional[Coin] = None
fee: Optional[Coin] = None
status: str = "pending"
created_at: str = ""
expires_at: str = ""
metadata: Dict[str, Any] = field(default_factory=dict)
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> PaymentSession:
return cls(
session_id=data.get("session_id", ""),
payer=data.get("payer", ""),
payee=data.get("payee", ""),
amount=Coin.from_dict(data.get("amount", {})) if data.get("amount") else None,
fee=Coin.from_dict(data.get("fee", {})) if data.get("fee") else None,
status=data.get("status", "pending"),
created_at=data.get("created_at", ""),
expires_at=data.get("expires_at", ""),
metadata=data.get("metadata", {}),
)
@dataclass
class DIDDocument:
id: str = ""
controller: List[str] = field(default_factory=list)
verification_method: List[Dict[str, Any]] = field(default_factory=list)
authentication: List[str] = field(default_factory=list)
service: List[Dict[str, Any]] = field(default_factory=list)
created: str = ""
updated: str = ""
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> DIDDocument:
return cls(
id=data.get("id", ""),
controller=data.get("controller", []),
verification_method=data.get("verification_method", []),
authentication=data.get("authentication", []),
service=data.get("service", []),
created=data.get("created", ""),
updated=data.get("updated", ""),
)
@dataclass
class ConstitutionRule:
rule_id: str = ""
description: str = ""
action_type: str = "deny"
parameters: Dict[str, Any] = field(default_factory=dict)
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> ConstitutionRule:
return cls(
rule_id=data.get("rule_id", ""),
description=data.get("description", ""),
action_type=data.get("action_type", "deny"),
parameters=data.get("parameters", {}),
)
@dataclass
class MicropaymentChannel:
channel_id: str = ""
sender: str = ""
receiver: str = ""
total_deposit: Optional[Coin] = None
balance: Optional[Coin] = None
nonce: int = 0
status: str = "open"
expires_at: str = ""
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> MicropaymentChannel:
return cls(
channel_id=data.get("channel_id", ""),
sender=data.get("sender", ""),
receiver=data.get("receiver", ""),
total_deposit=Coin.from_dict(data.get("total_deposit", {})) if data.get("total_deposit") else None,
balance=Coin.from_dict(data.get("balance", {})) if data.get("balance") else None,
nonce=int(data.get("nonce", 0)),
status=data.get("status", "open"),
expires_at=data.get("expires_at", ""),
)
@dataclass
class AgentQueryResponse:
agent_id: str = ""
response: str = ""
confidence: float = 0.0
latency_ms: int = 0
session_id: Optional[str] = None
metadata: Dict[str, Any] = field(default_factory=dict)
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> AgentQueryResponse:
return cls(
agent_id=data.get("agent_id", ""),
response=data.get("response", ""),
confidence=float(data.get("confidence", 0.0)),
latency_ms=int(data.get("latency_ms", 0)),
)
@dataclass
class AgentStatusResponse:
agent_id: str = ""
status: str = ""
is_online: bool = False
last_heartbeat: str = ""
total_queries: int = 0
uptime_seconds: int = 0
constitution_id: str = ""
resource_usage: Dict[str, str] = field(default_factory=dict)
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> AgentStatusResponse:
return cls(
agent_id=data.get("agent_id", ""),
status=data.get("status", ""),
is_online=data.get("is_online", False),
last_heartbeat=data.get("last_heartbeat", ""),
total_queries=int(data.get("total_queries", 0)),
uptime_seconds=int(data.get("uptime_seconds", 0)),
constitution_id=data.get("constitution_id", ""),
resource_usage=data.get("resource_usage", {}),
)
@dataclass
class AgentHistoryEntry:
query_id: str = ""
agent_id: str = ""
query: str = ""
response: str = ""
timestamp: str = ""
latency_ms: int = 0
cost: Optional[Coin] = None
status: str = ""
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> AgentHistoryEntry:
return cls(
query_id=data.get("query_id", ""),
agent_id=data.get("agent_id", ""),
query=data.get("query", ""),
response=data.get("response", ""),
timestamp=data.get("timestamp", ""),
latency_ms=int(data.get("latency_ms", 0)),
cost=Coin.from_dict(data.get("cost", {})) if data.get("cost") else None,
status=data.get("status", ""),
)
@dataclass
class TransactionInfo:
tx_hash: str = ""
from_address: str = ""
to_address: str = ""
amount: Optional[Coin] = None
fee: Optional[Coin] = None
status: str = "pending"
block_height: Optional[int] = None
timestamp: Optional[str] = None
memo: str = ""
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> TransactionInfo:
return cls(
tx_hash=data.get("tx_hash", ""),
from_address=data.get("from_address", ""),
to_address=data.get("to_address", ""),
amount=Coin.from_dict(data.get("amount", {})) if data.get("amount") else None,
fee=Coin.from_dict(data.get("fee", {})) if data.get("fee") else None,
status=data.get("status", "pending"),
block_height=data.get("block_height"),
timestamp=data.get("timestamp"),
memo=data.get("memo", ""),
)
@dataclass
class MPCResponse:
request_id: str = ""
status: str = "pending"
signature: Optional[str] = None
partial_signatures: List[str] = field(default_factory=list)
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> MPCResponse:
return cls(
request_id=data.get("request_id", ""),
status=data.get("status", "pending"),
signature=data.get("signature"),
partial_signatures=data.get("partial_signatures", []),
)
@dataclass
class Proposal:
proposal_id: int = 0
title: str = ""
description: str = ""
proposer: str = ""
status: str = ""
submit_time: str = ""
deposit_end_time: str = ""
voting_start_time: str = ""
voting_end_time: str = ""
total_deposit: List[Coin] = field(default_factory=list)
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> Proposal:
return cls(
proposal_id=int(data.get("proposal_id", 0) or data.get("id", 0)),
title=data.get("title", ""),
description=data.get("description", ""),
proposer=data.get("proposer", ""),
status=data.get("status", ""),
submit_time=data.get("submit_time", ""),
deposit_end_time=data.get("deposit_end_time", ""),
voting_start_time=data.get("voting_start_time", ""),
voting_end_time=data.get("voting_end_time", ""),
total_deposit=[Coin.from_dict(c) for c in data.get("total_deposit", [])],
)
@dataclass
class Vote:
proposal_id: int = 0
voter: str = ""
option: str = ""
weight: str = "1.0"
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> Vote:
return cls(
proposal_id=int(data.get("proposal_id", 0)),
voter=data.get("voter", ""),
option=data.get("option", ""),
weight=data.get("weight", "1.0"),
)
@dataclass
class TallyResult:
yes: str = "0"
no: str = "0"
abstain: str = "0"
no_with_veto: str = "0"
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> TallyResult:
return cls(
yes=data.get("yes", "0"),
no=data.get("no", "0"),
abstain=data.get("abstain", "0"),
no_with_veto=data.get("no_with_veto", "0"),
)
@dataclass
class ChainEvent:
event_type: str = ""
attributes: Dict[str, str] = field(default_factory=dict)
block_height: int = 0
tx_hash: Optional[str] = None
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> ChainEvent:
attrs = {}
for attr in data.get("attributes", []):
if isinstance(attr, dict) and "key" in attr and "value" in attr:
attrs[attr["key"]] = attr["value"]
return cls(
event_type=data.get("type", ""),
attributes=attrs,
block_height=int(data.get("height", 0)),
tx_hash=data.get("tx_hash"),
)
@dataclass
class AgentEvent:
event_id: str = ""
agent_id: str = ""
event_type: str = ""
payload: Dict[str, Any] = field(default_factory=dict)
timestamp: str = ""
@classmethod
def from_dict(cls, data: Dict[str, Any]) -> AgentEvent:
return cls(
event_id=data.get("event_id", ""),
agent_id=data.get("agent_id", ""),
event_type=data.get("event_type", ""),
payload=data.get("payload", {}),
timestamp=data.get("timestamp", ""),
)
5. 工具函数
# 文件: msg_sdk/utils.py
from __future__ import annotations
import hashlib
import json
import time
from decimal import Decimal
from typing import Any, Dict
def umsg_to_msg(umsg: str | int | Decimal) -> Decimal:
return Decimal(str(umsg)) / Decimal("1000000000000000000")
def msg_to_umsg(msg: str | int | Decimal) -> str:
result = Decimal(str(msg)) * Decimal("1000000000000000000")
return result.quantize(Decimal("1")).to_eng_string()
def sha256(data: bytes) -> bytes:
return hashlib.sha256(data).digest()
def ripemd160(data: bytes) -> bytes:
h = hashlib.new("ripemd160")
h.update(data)
return h.digest()
def generate_nonce() -> str:
raw = f"{time.time_ns()}{__import__('os').urandom(8).hex()}"
return hashlib.sha256(raw.encode()).hexdigest()[:32]
def sort_json(data: Dict[str, Any]) -> str:
return json.dumps(data, separators=(",", ":"), sort_keys=True)
def encode_amino(tx_body: Dict[str, Any]) -> bytes:
raw = sort_json(tx_body)
return raw.encode("utf-8")
def current_timestamp() -> str:
return time.strftime("%Y-%m-%dT%H:%M:%S.000000Z", time.gmtime())
6. 钱包与密钥管理
# 文件: msg_sdk/wallet.py
from __future__ import annotations
import hashlib
from dataclasses import dataclass
from typing import Optional
from bech32 import bech32_encode, bech32_decode, convertbits
from bip_utils import (
Bip39SeedGenerator,
Bip39MnemonicGenerator,
Bip39WordsNum,
Bip44,
Bip44Coins,
Bip44Changes,
)
@dataclass
class KeyPair:
private_key: bytes
public_key: bytes
address: str
mnemonic: str
class Wallet:
BECH32_PREFIX: str = "msg"
COIN_TYPE: int = 118
def __init__(self, private_key: bytes, public_key: bytes, address: str, mnemonic: str) -> None:
self._private_key = private_key
self._public_key = public_key
self._address = address
self._mnemonic = mnemonic
@property
def private_key(self) -> bytes:
return self._private_key
@property
def public_key(self) -> bytes:
return self._public_key
@property
def address(self) -> str:
return self._address
@property
def mnemonic(self) -> str:
return self._mnemonic
@classmethod
def generate(cls) -> Wallet:
mnemonic = Bip39MnemonicGenerator().FromWordsNumber(Bip39WordsNum.WORDS_NUM_24)
return cls.from_mnemonic(str(mnemonic))
@classmethod
def from_mnemonic(cls, mnemonic: str) -> Wallet:
seed = Bip39SeedGenerator(mnemonic).Generate()
bip44_ctx = Bip44.FromSeed(seed, Bip44Coins.COSMOS)
acc = bip44_ctx.Purpose().Coin().Account(0).Change(Bip44Changes.CHAIN_EXT).AddressIndex(0)
priv = acc.PrivateKey().Raw().ToBytes()
pub = acc.PublicKey().Raw().ToBytes()
address = cls._derive_address(pub)
return cls(private_key=priv, public_key=pub, address=address, mnemonic=mnemonic)
@classmethod
def from_private_key(cls, private_key_hex: str) -> Wallet:
private_key = bytes.fromhex(private_key_hex)
from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
priv = Ed25519PrivateKey.from_private_bytes(private_key)
public_key = priv.public_key().public_bytes_raw()
address = cls._derive_address(public_key)
return cls(private_key=private_key, public_key=public_key, address=address, mnemonic="")
@classmethod
def _derive_address(cls, public_key: bytes) -> str:
sha3 = hashlib.sha3_512(public_key).hexdigest()[:40]
raw = bytes.fromhex(sha3)
sha = hashlib.sha256(raw).digest()
checksum = sha.hex()[:4]
five_bit = convertbits(list(raw) + bytes.fromhex(checksum), 8, 5)
if five_bit is None:
raise ValueError("bit conversion failed")
return bech32_encode(cls.BECH32_PREFIX, five_bit)
@classmethod
def validate_address(cls, address: str) -> bool:
hrp, data = bech32_decode(address)
if hrp != cls.BECH32_PREFIX or data is None:
return False
return True
def sign(self, message: bytes) -> bytes:
from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
priv = Ed25519PrivateKey.from_private_bytes(self._private_key)
return priv.sign(message)
def sign_hex(self, message_hex: str) -> str:
return self.sign(bytes.fromhex(message_hex)).hex()
def verify(self, message: bytes, signature: bytes) -> bool:
from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey
try:
Ed25519PublicKey.from_public_bytes(self._public_key).verify(signature, message)
return True
except Exception:
return False
def sign_transaction(self, tx_body: dict) -> dict:
from msg_sdk.utils import encode_amino, generate_nonce
sign_doc = {"chain_id": tx_body.get("chain_id", "msg-chain-1"),"account_number": str(tx_body.get("account_number", "0")),"sequence": str(tx_body.get("sequence", "0")),"msgs": tx_body.get("msgs", []),"memo": tx_body.get("memo", ""),"fee": tx_body.get("fee", {"gas": "200000", "amount": [{"denom": "umsg", "amount": "5000"}]}),"nonce": generate_nonce()}
def sign_transaction(self, tx_body: dict) -> dict:
from msg_sdk.utils import encode_amino, generate_nonce
sign_doc = {
"chain_id": tx_body.get("chain_id", "msg-chain-1"),
"account_number": str(tx_body.get("account_number", "0")),
"sequence": str(tx_body.get("sequence", "0")),
"msgs": tx_body.get("msgs", []),
"memo": tx_body.get("memo", ""),
"fee": tx_body.get("fee", {"gas": "200000", "amount": [{"denom": "umsg", "amount": "5000"}]}),
"nonce": generate_nonce(),
}
encoded = encode_amino(sign_doc)
sig = self.sign(encoded)
tx_body["signatures"] = [{"pub_key": {"type": "tendermint/PubKeyEd25519", "value": self._public_key.hex()},"signature": sig.hex(),"nonce": sign_doc["nonce"]}]
return tx_body
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 {"address": self._address, "public_key": self._public_key.hex(), "has_mnemonic": bool(self._mnemonic)}
def __repr__(self) -> str:
return f"Wallet(address={self._address})"
class DilithiumWallet:
"""基于 Dilithium-5 后量子密码学的钱包。"""
BECH32_PREFIX: str = "msg"
def __init__(self, private_key: bytes, public_key: bytes, address: str, mnemonic: str = "") -> None:
self._private_key = private_key
self._public_key = public_key
self._address = address
self._mnemonic = mnemonic
@property
def address(self) -> str:
return self._address
@classmethod
def generate(cls) -> DilithiumWallet:
from Crypto.PublicKey import Dilithium
key = Dilithium.generate()
priv = key.export_key(format="DER")
pub = key.public_key().export_key(format="DER")
addr = cls._derive_address(key.public_key().export_key(format="Raw"))
return cls(private_key=priv, public_key=pub, address=addr)
@classmethod
def _derive_address(cls, raw: bytes) -> str:
sha3 = hashlib.sha3_512(raw).hexdigest()[:40]
data = bytes.fromhex(sha3)
sha = hashlib.sha256(data).digest()
checksum = sha.hex()[:4]
fb = convertbits(list(data) + bytes.fromhex(checksum), 8, 5)
if fb is None:
raise ValueError("bit conversion failed")
return bech32_encode(cls.BECH32_PREFIX, fb)
def sign(self, message: bytes) -> bytes:
from Crypto.PublicKey import Dilithium
return Dilithium.import_key(self._private_key).sign(message)
def verify(self, message: bytes, signature: bytes) -> bool:
from Crypto.PublicKey import Dilithium
try:
Dilithium.import_key(self._public_key).verify(message, signature)
return True
except Exception:
return False
@classmethod
def validate_address(cls, address: str) -> bool:
hrp, data = bech32_decode(address)
return hrp == cls.BECH32_PREFIX and data is not None
def __repr__(self) -> str:
return f"DilithiumWallet(address={self._address})"
7. 核心客户端
# 文件: msg_sdk/client.py
from __future__ import annotations
import json
from dataclasses import dataclass
from decimal import Decimal
from typing import Any, Dict, List, Optional, TypeVar
import httpx
from msg_sdk.errors import extract_error, raise_for_code, MSGChainError
from msg_sdk.errors import InsufficientFundsError, UnauthorizedError, RateLimitError
from msg_sdk.types import Coin, BalanceResponse, SupplyResponse, CodeResponse, ContractResponse, SmartQueryResponse
from msg_sdk.wallet import Wallet
T = TypeVar("T")
@dataclass
class ChainConfig:
chain_id: str = "msg-chain-1"
bech32_prefix: str = "msg"
coin_type: int = 118
msg_decimals: int = 18
gas_low: str = "1000000000"
gas_avg: str = "1000000000"
gas_high: str = "1000000000"
gas_denom: str = "attoMSG"
block_time_seconds: int = 5
rpc_url: str = "http://localhost:26657"
rest_url: str = "http://localhost:1317"
grpc_endpoint: str = "localhost:9090"
@dataclass
class GasFee:
gas_limit: int
gas_price: str
fee: Coin
@classmethod
def calculate(cls, gas_limit: int, gas_price: str = "1000000000") -> GasFee:
gas_price_atto = int(gas_price)
fee_amount = str(gas_limit * gas_price_atto)
return cls(gas_limit=gas_limit, gas_price=gas_price, fee=Coin(denom="attoMSG", amount=fee_amount))
class MSGChainClient:
"""MSG Chain 核心客户端 — 连接管理、错误处理、交易构建。"""
def __init__(
self,
config: Optional[ChainConfig] = None,
wallet: Optional[Wallet] = None,
timeout: float = 30.0,
retry_count: int = 3,
) -> None:
self.config = config or ChainConfig()
self._wallet = wallet
self._timeout = timeout
self._retry_count = retry_count
self._async_client: Optional[httpx.AsyncClient] = None
self._sync_client: Optional[httpx.Client] = None
async def connect(self) -> None:
if self._async_client is None or self._async_client.is_closed:
self._async_client = httpx.AsyncClient(base_url=self.config.rest_url, timeout=self._timeout)
async def disconnect(self) -> None:
if self._async_client and not self._async_client.is_closed:
await self._async_client.aclose()
def connect_sync(self) -> None:
if self._sync_client is None or self._sync_client.is_closed:
self._sync_client = httpx.Client(base_url=self.config.rest_url, timeout=self._timeout)
def disconnect_sync(self) -> None:
if self._sync_client and not self._sync_client.is_closed:
self._sync_client.close()
async def __aenter__(self) -> MSGChainClient:
await self.connect()
return self
async def __aexit__(self, *args: Any) -> None:
await self.disconnect()
def __enter__(self) -> MSGChainClient:
self.connect_sync()
return self
def __exit__(self, *args: Any) -> None:
self.disconnect_sync()
@property
def wallet(self) -> Optional[Wallet]:
return self._wallet
@wallet.setter
def wallet(self, w: Wallet) -> None:
self._wallet = w
def is_connected(self) -> bool:
ok = self._async_client is not None and not self._async_client.is_closed
ok = ok or (self._sync_client is not None and not self._sync_client.is_closed)
return ok
async def _request(self, method: str, path: str, json_data: Optional[Dict] = None, params: Optional[Dict] = None) -> Dict:
await self.connect()
last_err: Optional[Exception] = None
for _ in range(self._retry_count):
try:
resp = await self._async_client.request(method=method, url=path, json=json_data, params=params)
resp.raise_for_status()
body = resp.json()
if "code" in body:
extract_error(body)
return body
except MSGChainError:
raise
except httpx.HTTPStatusError as e:
if e.response.status_code == 429:
last_err = RateLimitError(message="rate limit")
continue
raise
except (httpx.TimeoutException, httpx.NetworkError) as e:
last_err = e
continue
raise MSGChainError(message=f"request failed after {self._retry_count} retries") from last_err
def _request_sync(self, method: str, path: str, json_data: Optional[Dict] = None, params: Optional[Dict] = None) -> Dict:
self.connect_sync()
last_err: Optional[Exception] = None
for _ in range(self._retry_count):
try:
resp = self._sync_client.request(method=method, url=path, json=json_data, params=params)
resp.raise_for_status()
body = resp.json()
if "code" in body:
extract_error(body)
return body
except MSGChainError:
raise
except httpx.HTTPStatusError as e:
if e.response.status_code == 429:
last_err = RateLimitError(message="rate limit")
continue
raise
except (httpx.TimeoutException, httpx.NetworkError) as e:
last_err = e
continue
raise MSGChainError(message=f"request failed after {self._retry_count} retries") from last_err
async def get_balance(self, address: str) -> List[Coin]:
if not Wallet.validate_address(address):
from msg_sdk.errors import InvalidParamError
raise InvalidParamError(message=f"invalid address: {address}")
body = await self._request("GET", f"/cosmos/bank/v1beta1/balances/{address}")
return BalanceResponse.from_dict(body).balances
def get_balance_sync(self, address: str) -> List[Coin]:
if not Wallet.validate_address(address):
from msg_sdk.errors import InvalidParamError
raise InvalidParamError(message=f"invalid address: {address}")
body = self._request_sync("GET", f"/cosmos/bank/v1beta1/balances/{address}")
return BalanceResponse.from_dict(body).balances
async def get_supply(self) -> List[Coin]:
body = await self._request("GET", "/cosmos/bank/v1beta1/supply")
return SupplyResponse.from_dict(body).supply
def get_supply_sync(self) -> List[Coin]:
body = self._request_sync("GET", "/cosmos/bank/v1beta1/supply")
return SupplyResponse.from_dict(body).supply
async def get_validators(self) -> List[Any]:
body = await self._request("GET", "/cosmos/staking/v1beta1/validators")
return body.get("validators", [])
def get_validators_sync(self) -> List[Any]:
body = self._request_sync("GET", "/cosmos/staking/v1beta1/validators")
return body.get("validators", [])
async def query_code(self, code_id: int) -> CodeResponse:
body = await self._request("GET", f"/cosmwasm/wasm/v1/code/{code_id}")
return CodeResponse.from_dict(body)
def query_code_sync(self, code_id: int) -> CodeResponse:
body = self._request_sync("GET", f"/cosmwasm/wasm/v1/code/{code_id}")
return CodeResponse.from_dict(body)
async def query_contract(self, address: str) -> ContractResponse:
body = await self._request("GET", f"/cosmwasm/wasm/v1/contract/{address}")
return ContractResponse.from_dict(body)
def query_contract_sync(self, address: str) -> ContractResponse:
body = self._request_sync("GET", f"/cosmwasm/wasm/v1/contract/{address}")
return ContractResponse.from_dict(body)
async def query_smart(self, contract_address: str, query_msg: Dict[str, Any]) -> SmartQueryResponse:
body = await self._request("POST", f"/cosmwasm/wasm/v1/contract/{contract_address}/smart", json_data={"query_msg": query_msg})
return SmartQueryResponse.from_dict(body.get("data", body))
def query_smart_sync(self, contract_address: str, query_msg: Dict[str, Any]) -> SmartQueryResponse:
body = self._request_sync("POST", f"/cosmwasm/wasm/v1/contract/{contract_address}/smart", json_data={"query_msg": query_msg})
return SmartQueryResponse.from_dict(body.get("data", body))
async def get_account_info(self, address: str) -> Dict[str, Any]:
body = await self._request("GET", f"/cosmos/auth/v1beta1/accounts/{address}")
acct = body.get("account", body.get("account", {}))
if isinstance(acct, dict):
base = acct.get("base_account", acct)
return {"account_number": int(base.get("account_number", 0)), "sequence": int(base.get("sequence", 0)), "address": base.get("address", address)}
return {"account_number": 0, "sequence": 0, "address": address}
def get_account_info_sync(self, address: str) -> Dict[str, Any]:
body = self._request_sync("GET", f"/cosmos/auth/v1beta1/accounts/{address}")
acct = body.get("account", body.get("account", {}))
if isinstance(acct, dict):
base = acct.get("base_account", acct)
return {"account_number": int(base.get("account_number", 0)), "sequence": int(base.get("sequence", 0)), "address": base.get("address", address)}
return {"account_number": 0, "sequence": 0, "address": address}
async def broadcast_transaction(self, tx_bytes: str) -> Dict[str, Any]:
return await self._request("POST", "/cosmos/tx/v1beta1/txs", json_data={"tx_bytes": tx_bytes, "mode": "BROADCAST_MODE_SYNC"})
def broadcast_transaction_sync(self, tx_bytes: str) -> Dict[str, Any]:
return self._request_sync("POST", "/cosmos/tx/v1beta1/txs", json_data={"tx_bytes": tx_bytes, "mode": "BROADCAST_MODE_SYNC"})
async def wait_for_tx(self, tx_hash: str, timeout_blocks: int = 30) -> Dict[str, Any]:
import asyncio
for _ in range(timeout_blocks):
await asyncio.sleep(self.config.block_time_seconds)
try:
body = await self._request("GET", f"/cosmos/tx/v1beta1/txs/{tx_hash}")
txr = body.get("tx_response", body)
if txr.get("code", -1) == 0:
return txr
from msg_sdk.errors import ContractFailedError
raise ContractFailedError(message=txr.get("raw_log", "tx failed"), data={"tx_hash": tx_hash})
except MSGChainError:
raise
except Exception:
continue
from msg_sdk.errors import QueryTimeoutError
raise QueryTimeoutError(message=f"tx {tx_hash} not confirmed after {timeout_blocks} blocks")
def wait_for_tx_sync(self, tx_hash: str, timeout_blocks: int = 30) -> Dict[str, Any]:
import time
for _ in range(timeout_blocks):
time.sleep(self.config.block_time_seconds)
try:
body = self._request_sync("GET", f"/cosmos/tx/v1beta1/txs/{tx_hash}")
txr = body.get("tx_response", body)
if txr.get("code", -1) == 0:
return txr
from msg_sdk.errors import ContractFailedError
raise ContractFailedError(message=txr.get("raw_log", "tx failed"), data={"tx_hash": tx_hash})
except MSGChainError:
raise
except Exception:
continue
from msg_sdk.errors import QueryTimeoutError
raise QueryTimeoutError(message=f"tx {tx_hash} not confirmed after {timeout_blocks} blocks")
async def send_tokens(self, to_address: str, amount: str, denom: str = "umsg", memo: str = "") -> Dict[str, Any]:
if self._wallet is None:
raise UnauthorizedError(message="wallet not set")
from_addr = self._wallet.address
balance = await self.get_balance(from_addr)
bal_umsg = sum(Decimal(c.amount) for c in balance if c.denom == denom)
if bal_umsg < Decimal(amount):
raise InsufficientFundsError(message=f"balance {bal_umsg} {denom} < {amount} {denom}")
info = await self.get_account_info(from_addr)
msg = {"type": "cosmos-sdk/MsgSend", "value": {"from_address": from_addr, "to_address": to_address, "amount": [{"denom": denom, "amount": amount}]}}
tx = {"chain_id": self.config.chain_id, "account_number": info["account_number"], "sequence": info["sequence"], "msgs": [msg], "memo": memo, "fee": {"gas": "200000", "amount": [{"denom": "umsg", "amount": "5000"}]}}
signed = self._wallet.sign_transaction(tx)
from msg_sdk.utils import encode_amino
tx_bytes = encode_amino(signed).hex()
result = await self.broadcast_transaction(tx_bytes)
tx_hash = result.get("tx_response", result).get("txhash", "")
return await self.wait_for_tx(tx_hash) if tx_hash else result
def send_tokens_sync(self, to_address: str, amount: str, denom: str = "umsg", memo: str = "") -> Dict[str, Any]:
if self._wallet is None:
raise UnauthorizedError(message="wallet not set")
from_addr = self._wallet.address
balance = self.get_balance_sync(from_addr)
bal_umsg = sum(Decimal(c.amount) for c in balance if c.denom == denom)
if bal_umsg < Decimal(amount):
raise InsufficientFundsError(message=f"balance {bal_umsg} {denom} < {amount} {denom}")
info = self.get_account_info_sync(from_addr)
msg = {"type": "cosmos-sdk/MsgSend", "value": {"from_address": from_addr, "to_address": to_address, "amount": [{"denom": denom, "amount": amount}]}}
tx = {"chain_id": self.config.chain_id, "account_number": info["account_number"], "sequence": info["sequence"], "msgs": [msg], "memo": memo, "fee": {"gas": "200000", "amount": [{"denom": "umsg", "amount": "5000"}]}}
signed = self._wallet.sign_transaction(tx)
from msg_sdk.utils import encode_amino
tx_bytes = encode_amino(signed).hex()
result = self.broadcast_transaction_sync(tx_bytes)
tx_hash = result.get("tx_response", result).get("txhash", "")
return self.wait_for_tx_sync(tx_hash) if tx_hash else result
def estimate_gas(self, message_count: int = 1, complexity: str = "simple") -> GasFee:
base = {"simple": 100_000, "complex": 300_000, "contract": 500_000}.get(complexity, 200_000)
return GasFee.calculate(base * message_count, gas_price=self.config.gas_avg)
async def get_node_info(self) -> Dict[str, Any]:
return await self._request("GET", "/cosmos/base/tendermint/v1beta1/node_info")
def get_node_info_sync(self) -> Dict[str, Any]:
return self._request_sync("GET", "/cosmos/base/tendermint/v1beta1/node_info")
async def get_latest_block(self) -> Dict[str, Any]:
return await self._request("GET", "/cosmos/base/tendermint/v1beta1/blocks/latest")
def get_latest_block_sync(self) -> Dict[str, Any]:
return self._request_sync("GET", "/cosmos/base/tendermint/v1beta1/blocks/latest")
def __repr__(self) -> str:
return f"MSGChainClient(chain_id={self.config.chain_id}, rest={self.config.rest_url})"
8. 银行客户端
# 文件: msg_sdk/bank.py
from __future__ import annotations
from decimal import Decimal
from typing import Any, Dict, List, Optional
from msg_sdk.client import MSGChainClient
from msg_sdk.errors import InvalidParamError
from msg_sdk.types import Coin
from msg_sdk.wallet import Wallet
class BankClient:
"""银行客户端 — 余额查询、供应量查询、转账。"""
def __init__(self, client: MSGChainClient) -> None:
self._client = client
async def get_balance(self, address: str, denom: str = "umsg") -> Coin:
if not Wallet.validate_address(address):
raise InvalidParamError(message=f"invalid address: {address}")
balances = await self._client.get_balance(address)
for c in balances:
if c.denom == denom:
return c
return Coin(denom=denom, amount="0")
def get_balance_sync(self, address: str, denom: str = "umsg") -> Coin:
if not Wallet.validate_address(address):
raise InvalidParamError(message=f"invalid address: {address}")
balances = self._client.get_balance_sync(address)
for c in balances:
if c.denom == denom:
return c
return Coin(denom=denom, amount="0")
async def get_all_balances(self, address: str) -> List[Coin]:
return await self._client.get_balance(address)
def get_all_balances_sync(self, address: str) -> List[Coin]:
return self._client.get_balance_sync(address)
async def get_balance_umsg(self, address: str) -> Decimal:
return Decimal((await self.get_balance(address, "umsg")).amount)
def get_balance_umsg_sync(self, address: str) -> Decimal:
return Decimal(self.get_balance_sync(address, "umsg").amount)
async def get_supply(self, denom: str = "umsg") -> Coin:
supply = await self._client.get_supply()
for c in supply:
if c.denom == denom:
return c
return Coin(denom=denom, amount="0")
def get_supply_sync(self, denom: str = "umsg") -> Coin:
supply = self._client.get_supply_sync()
for c in supply:
if c.denom == denom:
return c
return Coin(denom=denom, amount="0")
async def get_all_supply(self) -> List[Coin]:
return await self._client.get_supply()
def get_all_supply_sync(self) -> List[Coin]:
return self._client.get_supply_sync()
async def transfer(self, to_address: str, amount: str, denom: str = "umsg", memo: str = "") -> Dict[str, Any]:
return await self._client.send_tokens(to_address=to_address, amount=amount, denom=denom, memo=memo)
def transfer_sync(self, to_address: str, amount: str, denom: str = "umsg", memo: str = "") -> Dict[str, Any]:
return self._client.send_tokens_sync(to_address=to_address, amount=amount, denom=denom, memo=memo)
async def transfer_msg(self, to_address: str, amount_msg: str, memo: str = "") -> Dict[str, Any]:
from msg_sdk.utils import msg_to_umsg
return await self.transfer(to_address, msg_to_umsg(amount_msg), memo=memo)
def transfer_msg_sync(self, to_address: str, amount_msg: str, memo: str = "") -> Dict[str, Any]:
from msg_sdk.utils import msg_to_umsg
return self.transfer_sync(to_address, msg_to_umsg(amount_msg), memo=memo)
async def get_balances_batch(self, addresses: List[str]) -> Dict[str, List[Coin]]:
import asyncio
tasks = [self.get_all_balances(a) for a in addresses]
results = await asyncio.gather(*tasks, return_exceptions=True)
return {a: (r if isinstance(r, list) else []) for a, r in zip(addresses, results)}
def __repr__(self) -> str:
return "BankClient(client=MSGChainClient)"
9. 质押客户端
# 文件: msg_sdk/staking.py
from __future__ import annotations
from decimal import Decimal
from typing import Any, Dict, List, Optional
from msg_sdk.client import MSGChainClient
from msg_sdk.errors import UnauthorizedError, ContractFailedError
from msg_sdk.types import Coin, Validator, ValidatorResponse, DelegationResponse, RewardsResponse
class StakingClient:
"""质押客户端 — 验证人查询、委托、奖励。"""
def __init__(self, client: MSGChainClient) -> None:
self._client = client
async def get_validators(self, status: str = "BOND_STATUS_BONDED") -> List[Validator]:
body = await self._client._request("GET", "/cosmos/staking/v1beta1/validators", params={"status": status})
return ValidatorResponse.from_dict(body).validators
def get_validators_sync(self, status: str = "BOND_STATUS_BONDED") -> List[Validator]:
body = self._client._request_sync("GET", "/cosmos/staking/v1beta1/validators", params={"status": status})
return ValidatorResponse.from_dict(body).validators
async def get_validator(self, operator_address: str) -> Validator:
body = await self._client._request("GET", f"/cosmos/staking/v1beta1/validators/{operator_address}")
return Validator.from_dict(body.get("validator", body))
def get_validator_sync(self, operator_address: str) -> Validator:
body = self._client._request_sync("GET", f"/cosmos/staking/v1beta1/validators/{operator_address}")
return Validator.from_dict(body.get("validator", body))
async def get_active_validators(self) -> List[Validator]:
return [v for v in await self.get_validators() if v.status == 3 and not v.jailed]
def get_active_validators_sync(self) -> List[Validator]:
return [v for v in self.get_validators_sync() if v.status == 3 and not v.jailed]
async def get_delegations(self, delegator: str) -> List[DelegationResponse]:
body = await self._client._request("GET", f"/cosmos/staking/v1beta1/delegations/{delegator}")
return [DelegationResponse.from_dict(d) for d in body.get("delegation_responses", [])]
def get_delegations_sync(self, delegator: str) -> List[DelegationResponse]:
body = self._client._request_sync("GET", f"/cosmos/staking/v1beta1/delegations/{delegator}")
return [DelegationResponse.from_dict(d) for d in body.get("delegation_responses", [])]
async def get_delegation(self, delegator: str, validator: str) -> Optional[DelegationResponse]:
try:
body = await self._client._request("GET", f"/cosmos/staking/v1beta1/validators/{validator}/delegations/{delegator}")
return DelegationResponse.from_dict(body.get("delegation_response", body))
except Exception:
return None
def get_delegation_sync(self, delegator: str, validator: str) -> Optional[DelegationResponse]:
try:
body = self._client._request_sync("GET", f"/cosmos/staking/v1beta1/validators/{validator}/delegations/{delegator}")
return DelegationResponse.from_dict(body.get("delegation_response", body))
except Exception:
return None
async def get_rewards(self, delegator: str) -> RewardsResponse:
body = await self._client._request("GET", f"/cosmos/distribution/v1beta1/delegators/{delegator}/rewards")
return RewardsResponse.from_dict(body)
def get_rewards_sync(self, delegator: str) -> RewardsResponse:
body = self._client._request_sync("GET", f"/cosmos/distribution/v1beta1/delegators/{delegator}/rewards")
return RewardsResponse.from_dict(body)
async def get_total_staked(self, address: str) -> Decimal:
return sum(Decimal(d.balance.amount) for d in await self.get_delegations(address))
def get_total_staked_sync(self, address: str) -> Decimal:
return sum(Decimal(d.balance.amount) for d in self.get_delegations_sync(address))
async def delegate(self, validator_address: str, amount: str, denom: str = "umsg") -> Dict[str, Any]:
if self._client.wallet is None:
raise UnauthorizedError(message="wallet not set")
delegator = self._client.wallet.address
info = await self._client.get_account_info(delegator)
msg = {"type": "cosmos-sdk/MsgDelegate", "value": {"delegator_address": delegator, "validator_address": validator_address, "amount": {"denom": denom, "amount": amount}}}
tx = {"chain_id": self._client.config.chain_id, "account_number": info["account_number"], "sequence": info["sequence"], "msgs": [msg], "memo": "", "fee": {"gas": "250000", "amount": [{"denom": "umsg", "amount": "6250"}]}}
signed = self._client.wallet.sign_transaction(tx)
from msg_sdk.utils import encode_amino
tx_bytes = encode_amino(signed).hex()
result = await self._client.broadcast_transaction(tx_bytes)
tx_hash = result.get("tx_response", result).get("txhash", "")
return await self._client.wait_for_tx(tx_hash) if tx_hash else result
def delegate_sync(self, validator_address: str, amount: str, denom: str = "umsg") -> Dict[str, Any]:
if self._client.wallet is None:
raise UnauthorizedError(message="wallet not set")
delegator = self._client.wallet.address
info = self._client.get_account_info_sync(delegator)
msg = {"type": "cosmos-sdk/MsgDelegate", "value": {"delegator_address": delegator, "validator_address": validator_address, "amount": {"denom": denom, "amount": amount}}}
tx = {"chain_id": self._client.config.chain_id, "account_number": info["account_number"], "sequence": info["sequence"], "msgs": [msg], "memo": "", "fee": {"gas": "250000", "amount": [{"denom": "umsg", "amount": "6250"}]}}
signed = self._client.wallet.sign_transaction(tx)
from msg_sdk.utils import encode_amino
tx_bytes = encode_amino(signed).hex()
result = self._client.broadcast_transaction_sync(tx_bytes)
tx_hash = result.get("tx_response", result).get("txhash", "")
return self._client.wait_for_tx_sync(tx_hash) if tx_hash else result
async def undelegate(self, validator_address: str, amount: str, denom: str = "umsg") -> Dict[str, Any]:
if self._client.wallet is None:
raise UnauthorizedError(message="wallet not set")
delegator = self._client.wallet.address
info = await self._client.get_account_info(delegator)
msg = {"type": "cosmos-sdk/MsgUndelegate", "value": {"delegator_address": delegator, "validator_address": validator_address, "amount": {"denom": denom, "amount": amount}}}
tx = {"chain_id": self._client.config.chain_id, "account_number": info["account_number"], "sequence": info["sequence"], "msgs": [msg], "memo": "", "fee": {"gas": "200000", "amount": [{"denom": "umsg", "amount": "5000"}]}}
signed = self._client.wallet.sign_transaction(tx)
from msg_sdk.utils import encode_amino
tx_bytes = encode_amino(signed).hex()
result = await self._client.broadcast_transaction(tx_bytes)
tx_hash = result.get("tx_response", result).get("txhash", "")
return await self._client.wait_for_tx(tx_hash) if tx_hash else result
def undelegate_sync(self, validator_address: str, amount: str, denom: str = "umsg") -> Dict[str, Any]:
if self._client.wallet is None:
raise UnauthorizedError(message="wallet not set")
delegator = self._client.wallet.address
info = self._client.get_account_info_sync(delegator)
msg = {"type": "cosmos-sdk/MsgUndelegate", "value": {"delegator_address": delegator, "validator_address": validator_address, "amount": {"denom": denom, "amount": amount}}}
tx = {"chain_id": self._client.config.chain_id, "account_number": info["account_number"], "sequence": info["sequence"], "msgs": [msg], "memo": "", "fee": {"gas": "200000", "amount": [{"denom": "umsg", "amount": "5000"}]}}
signed = self._client.wallet.sign_transaction(tx)
from msg_sdk.utils import encode_amino
tx_bytes = encode_amino(signed).hex()
result = self._client.broadcast_transaction_sync(tx_bytes)
tx_hash = result.get("tx_response", result).get("txhash", "")
return self._client.wait_for_tx_sync(tx_hash) if tx_hash else result
async def withdraw_rewards(self, validator_address: str) -> Dict[str, Any]:
if self._client.wallet is None:
raise UnauthorizedError(message="wallet not set")
delegator = self._client.wallet.address
info = await self._client.get_account_info(delegator)
msg = {"type": "cosmos-sdk/MsgWithdrawDelegationReward", "value": {"delegator_address": delegator, "validator_address": validator_address}}
tx = {"chain_id": self._client.config.chain_id, "account_number": info["account_number"], "sequence": info["sequence"], "msgs": [msg], "memo": "", "fee": {"gas": "200000", "amount": [{"denom": "umsg", "amount": "5000"}]}}
signed = self._client.wallet.sign_transaction(tx)
from msg_sdk.utils import encode_amino
tx_bytes = encode_amino(signed).hex()
result = await self._client.broadcast_transaction(tx_bytes)
tx_hash = result.get("tx_response", result).get("txhash", "")
return await self._client.wait_for_tx(tx_hash) if tx_hash else result
def withdraw_rewards_sync(self, validator_address: str) -> Dict[str, Any]:
if self._client.wallet is None:
raise UnauthorizedError(message="wallet not set")
delegator = self._client.wallet.address
info = self._client.get_account_info_sync(delegator)
msg = {"type": "cosmos-sdk/MsgWithdrawDelegationReward", "value": {"delegator_address": delegator, "validator_address": validator_address}}
tx = {"chain_id": self._client.config.chain_id, "account_number": info["account_number"], "sequence": info["sequence"], "msgs": [msg], "memo": "", "fee": {"gas": "200000", "amount": [{"denom": "umsg", "amount": "5000"}]}}
signed = self._client.wallet.sign_transaction(tx)
from msg_sdk.utils import encode_amino
tx_bytes = encode_amino(signed).hex()
result = self._client.broadcast_transaction_sync(tx_bytes)
tx_hash = result.get("tx_response", result).get("txhash", "")
return self._client.wait_for_tx_sync(tx_hash) if tx_hash else result
def __repr__(self) -> str:
return "StakingClient(client=MSGChainClient)"
10. 合约查询客户端
# 文件: msg_sdk/contract.py (合约查询部分)
from __future__ import annotations
from typing import Any, Dict, List, Optional
from msg_sdk.client import MSGChainClient
from msg_sdk.types import (
Coin, CodeResponse, ContractResponse, SmartQueryResponse,
AgentInfo, PaymentSession, DIDDocument, ConstitutionRule, MicropaymentChannel,
)
class ContractQueryClient:
"""合约查询客户端 — 智能查询 + 5 个 AI 合约的类型化查询。"""
def __init__(self, client: MSGChainClient) -> None:
self._client = client
async def query_smart(self, contract: str, msg: Dict[str, Any]) -> Any:
return await self._client.query_smart(contract, msg)
def query_smart_sync(self, contract: str, msg: Dict[str, Any]) -> Any:
return self._client.query_smart_sync(contract, msg)
async def query_raw(self, contract: str, key: bytes) -> bytes:
body = await self._client._request("GET", f"/cosmwasm/wasm/v1/contract/{contract}/raw/{key.hex()}")
data = body.get("data", "")
return bytes.fromhex(data) if data else b""
def query_raw_sync(self, contract: str, key: bytes) -> bytes:
body = self._client._request_sync("GET", f"/cosmwasm/wasm/v1/contract/{contract}/raw/{key.hex()}")
data = body.get("data", "")
return bytes.fromhex(data) if data else b""
async def get_code(self, code_id: int) -> CodeResponse:
return await self._client.query_code(code_id)
def get_code_sync(self, code_id: int) -> CodeResponse:
return self._client.query_code_sync(code_id)
async def get_contract_info(self, address: str) -> ContractResponse:
return await self._client.query_contract(address)
def get_contract_info_sync(self, address: str) -> ContractResponse:
return self._client.query_contract_sync(address)
# Agent Registry
async def get_agent(self, contract: str, agent_id: str) -> AgentInfo:
result = await self.query_smart(contract, {"get_agent": {"id": agent_id}})
data = result.data if isinstance(result, SmartQueryResponse) else result
if isinstance(data, dict) and "agent" in data:
data = data["agent"]
return AgentInfo.from_dict(data)
def get_agent_sync(self, contract: str, agent_id: str) -> AgentInfo:
result = self.query_smart_sync(contract, {"get_agent": {"id": agent_id}})
data = result.data if isinstance(result, SmartQueryResponse) else result
if isinstance(data, dict) and "agent" in data:
data = data["agent"]
return AgentInfo.from_dict(data)
async def list_agents(self, contract: str, start_after: Optional[str] = None, limit: int = 20) -> List[AgentInfo]:
q: Dict = {"list_agents": {"limit": limit}}
if start_after:
q["list_agents"]["start_after"] = start_after
result = await self.query_smart(contract, q)
data = result.data if isinstance(result, SmartQueryResponse) else result
agents = data.get("agents", data if isinstance(data, list) else [])
return [AgentInfo.from_dict(a) for a in agents]
def list_agents_sync(self, contract: str, start_after: Optional[str] = None, limit: int = 20) -> List[AgentInfo]:
q: Dict = {"list_agents": {"limit": limit}}
if start_after:
q["list_agents"]["start_after"] = start_after
result = self.query_smart_sync(contract, q)
data = result.data if isinstance(result, SmartQueryResponse) else result
agents = data.get("agents", data if isinstance(data, list) else [])
return [AgentInfo.from_dict(a) for a in agents]
# Payment
async def get_payment_session(self, contract: str, session_id: str) -> PaymentSession:
result = await self.query_smart(contract, {"get_session": {"session_id": session_id}})
data = result.data if isinstance(result, SmartQueryResponse) else result
if isinstance(data, dict) and "session" in data:
data = data["session"]
return PaymentSession.from_dict(data)
def get_payment_session_sync(self, contract: str, session_id: str) -> PaymentSession:
result = self.query_smart_sync(contract, {"get_session": {"session_id": session_id}})
data = result.data if isinstance(result, SmartQueryResponse) else result
if isinstance(data, dict) and "session" in data:
data = data["session"]
return PaymentSession.from_dict(data)
# DID
async def resolve_did(self, contract: str, did: str) -> DIDDocument:
result = await self.query_smart(contract, {"resolve_did": {"id": did}})
data = result.data if isinstance(result, SmartQueryResponse) else result
if isinstance(data, dict) and "document" in data:
data = data["document"]
return DIDDocument.from_dict(data)
def resolve_did_sync(self, contract: str, did: str) -> DIDDocument:
result = self.query_smart_sync(contract, {"resolve_did": {"id": did}})
data = result.data if isinstance(result, SmartQueryResponse) else result
if isinstance(data, dict) and "document" in data:
data = data["document"]
return DIDDocument.from_dict(data)
# Constitution
async def get_constitution(self, contract: str, constitution_id: str) -> List[ConstitutionRule]:
result = await self.query_smart(contract, {"get_constitution": {"constitution_id": constitution_id}})
data = result.data if isinstance(result, SmartQueryResponse) else result
rules = data.get("rules", data if isinstance(data, list) else [])
return [ConstitutionRule.from_dict(r) for r in rules]
def get_constitution_sync(self, contract: str, constitution_id: str) -> List[ConstitutionRule]:
result = self.query_smart_sync(contract, {"get_constitution": {"constitution_id": constitution_id}})
data = result.data if isinstance(result, SmartQueryResponse) else result
rules = data.get("rules", data if isinstance(data, list) else [])
return [ConstitutionRule.from_dict(r) for r in rules]
async def is_action_allowed(self, contract: str, constitution_id: str, action: str) -> bool:
result = await self.query_smart(contract, {"is_action_allowed": {"constitution_id": constitution_id, "action": action}})
data = result.data if isinstance(result, SmartQueryResponse) else result
return data.get("allowed", False) if isinstance(data, dict) else bool(data)
def is_action_allowed_sync(self, contract: str, constitution_id: str, action: str) -> bool:
result = self.query_smart_sync(contract, {"is_action_allowed": {"constitution_id": constitution_id, "action": action}})
data = result.data if isinstance(result, SmartQueryResponse) else result
return data.get("allowed", False) if isinstance(data, dict) else bool(data)
# Micropayment
async def get_channel(self, contract: str, channel_id: str) -> MicropaymentChannel:
result = await self.query_smart(contract, {"get_channel": {"channel_id": channel_id}})
data = result.data if isinstance(result, SmartQueryResponse) else result
if isinstance(data, dict) and "channel" in data:
data = data["channel"]
return MicropaymentChannel.from_dict(data)
def get_channel_sync(self, contract: str, channel_id: str) -> MicropaymentChannel:
result = self.query_smart_sync(contract, {"get_channel": {"channel_id": channel_id}})
data = result.data if isinstance(result, SmartQueryResponse) else result
if isinstance(data, dict) and "channel" in data:
data = data["channel"]
return MicropaymentChannel.from_dict(data)
async def batch_get_agents(self, contract: str, agent_ids: List[str]) -> Dict[str, Optional[AgentInfo]]:
import asyncio
async def fetch(aid: str):
try:
return aid, await self.get_agent(contract, aid)
except Exception:
return aid, None
results = await asyncio.gather(*[fetch(aid) for aid in agent_ids])
return dict(results)
def __repr__(self) -> str:
return "ContractQueryClient(client=MSGChainClient)"
class ContractExecuteClient:
"""合约执行客户端 — 构建、签名、广播合约交易。"""
def __init__(self, client: MSGChainClient) -> None:
self._client = client
async def _execute(self, contract: str, msg: Dict, funds=None, memo="", gas_limit=500_000) -> Dict:
if self._client.wallet is None:
from msg_sdk.errors import UnauthorizedError
raise UnauthorizedError(message="wallet not set")
sender = self._client.wallet.address
info = await self._client.get_account_info(sender)
execute_msg = {"type": "wasm/MsgExecuteContract", "value": {"sender": sender, "contract": contract, "msg": msg, "funds": [f.to_dict() for f in (funds or [])]}}
gp = int(float(self._client.config.gas_avg) * 1_000_000_000_000_000_000)
tx = {"chain_id": self._client.config.chain_id, "account_number": info["account_number"], "sequence": info["sequence"], "msgs": [execute_msg], "memo": memo, "fee": {"gas": str(gas_limit), "amount": [{"denom": "umsg", "amount": str(gas_limit * gp)}]}}
signed = self._client.wallet.sign_transaction(tx)
from msg_sdk.utils import encode_amino
tx_bytes = encode_amino(signed).hex()
result = await self._client.broadcast_transaction(tx_bytes)
txr = result.get("tx_response", result)
tx_hash = txr.get("txhash", "")
if tx_hash:
return await self._client.wait_for_tx(tx_hash)
code = txr.get("code", -1)
if code != 0:
from msg_sdk.errors import ContractFailedError
raise ContractFailedError(message=txr.get("raw_log", "execution failed"), data={"code": code})
return txr
def _execute_sync(self, contract: str, msg: Dict, funds=None, memo="", gas_limit=500_000) -> Dict:
if self._client.wallet is None:
from msg_sdk.errors import UnauthorizedError
raise UnauthorizedError(message="wallet not set")
sender = self._client.wallet.address
info = self._client.get_account_info_sync(sender)
execute_msg = {"type": "wasm/MsgExecuteContract", "value": {"sender": sender, "contract": contract, "msg": msg, "funds": [f.to_dict() for f in (funds or [])]}}
gp = int(float(self._client.config.gas_avg) * 1_000_000_000_000_000_000)
tx = {"chain_id": self._client.config.chain_id, "account_number": info["account_number"], "sequence": info["sequence"], "msgs": [execute_msg], "memo": memo, "fee": {"gas": str(gas_limit), "amount": [{"denom": "umsg", "amount": str(gas_limit * gp)}]}}
signed = self._client.wallet.sign_transaction(tx)
from msg_sdk.utils import encode_amino
tx_bytes = encode_amino(signed).hex()
result = self._client.broadcast_transaction_sync(tx_bytes)
txr = result.get("tx_response", result)
tx_hash = txr.get("txhash", "")
if tx_hash:
return self._client.wait_for_tx_sync(tx_hash)
code = txr.get("code", -1)
if code != 0:
from msg_sdk.errors import ContractFailedError
raise ContractFailedError(message=txr.get("raw_log", "execution failed"), data={"code": code})
return txr
# Agent Registry
async def register_agent(self, contract: str, agent_id: str, name: str, desc: str, constitution_id: str, metadata=None, funds=None) -> Dict:
return await self._execute(contract, {"register_agent": {"id": agent_id, "name": name, "description": desc, "constitution_id": constitution_id, "metadata": metadata or {}}}, funds=funds)
def register_agent_sync(self, contract: str, agent_id: str, name: str, desc: str, constitution_id: str, metadata=None, funds=None) -> Dict:
return self._execute_sync(contract, {"register_agent": {"id": agent_id, "name": name, "description": desc, "constitution_id": constitution_id, "metadata": metadata or {}}}, funds=funds)
async def update_agent(self, contract: str, agent_id: str, name=None, description=None, metadata=None) -> Dict:
upd = {"id": agent_id}
if name: upd["name"] = name
if description: upd["description"] = description
if metadata: upd["metadata"] = metadata
return await self._execute(contract, {"update_agent": upd})
def update_agent_sync(self, contract: str, agent_id: str, name=None, description=None, metadata=None) -> Dict:
upd = {"id": agent_id}
if name: upd["name"] = name
if description: upd["description"] = description
if metadata: upd["metadata"] = metadata
return self._execute_sync(contract, {"update_agent": upd})
async def deregister_agent(self, contract: str, agent_id: str) -> Dict:
return await self._execute(contract, {"deregister_agent": {"id": agent_id}})
def deregister_agent_sync(self, contract: str, agent_id: str) -> Dict:
return self._execute_sync(contract, {"deregister_agent": {"id": agent_id}})
async def set_agent_status(self, contract: str, agent_id: str, status: str) -> Dict:
return await self._execute(contract, {"set_status": {"id": agent_id, "status": status}})
def set_agent_status_sync(self, contract: str, agent_id: str, status: str) -> Dict:
return self._execute_sync(contract, {"set_status": {"id": agent_id, "status": status}})
# Payment
async def create_session(self, contract: str, session_id: str, payer: str, payee: str, amount: str, fee: str = "0", duration: int = 3600, metadata=None) -> Dict:
total = str(int(amount) + int(fee))
return await self._execute(contract, {"create_session": {"session_id": session_id, "payer": payer, "payee": payee, "amount": amount, "fee": fee, "duration_seconds": duration, "metadata": metadata or {}}}, funds=[Coin(denom="umsg", amount=total)])
def create_session_sync(self, contract: str, session_id: str, payer: str, payee: str, amount: str, fee: str = "0", duration: int = 3600, metadata=None) -> Dict:
total = str(int(amount) + int(fee))
return self._execute_sync(contract, {"create_session": {"session_id": session_id, "payer": payer, "payee": payee, "amount": amount, "fee": fee, "duration_seconds": duration, "metadata": metadata or {}}}, funds=[Coin(denom="umsg", amount=total)])
async def fund_session(self, contract: str, session_id: str, amount: str) -> Dict:
return await self._execute(contract, {"fund_session": {"session_id": session_id}}, funds=[Coin(denom="umsg", amount=amount)])
def fund_session_sync(self, contract: str, session_id: str, amount: str) -> Dict:
return self._execute_sync(contract, {"fund_session": {"session_id": session_id}}, funds=[Coin(denom="umsg", amount=amount)])
async def release_payment(self, contract: str, session_id: str) -> Dict:
return await self._execute(contract, {"release_payment": {"session_id": session_id}})
def release_payment_sync(self, contract: str, session_id: str) -> Dict:
return self._execute_sync(contract, {"release_payment": {"session_id": session_id}})
async def dispute_payment(self, contract: str, session_id: str, reason: str) -> Dict:
return await self._execute(contract, {"dispute_payment": {"session_id": session_id, "reason": reason}})
def dispute_payment_sync(self, contract: str, session_id: str, reason: str) -> Dict:
return self._execute_sync(contract, {"dispute_payment": {"session_id": session_id, "reason": reason}})
# DID
async def create_did(self, contract: str, did: str, controller: List[str], vm: List[Dict], auth: List[str], service=None) -> Dict:
return await self._execute(contract, {"create_did": {"id": did, "controller": controller, "verification_method": vm, "authentication": auth, "service": service or []}})
def create_did_sync(self, contract: str, did: str, controller: List[str], vm: List[Dict], auth: List[str], service=None) -> Dict:
return self._execute_sync(contract, {"create_did": {"id": did, "controller": controller, "verification_method": vm, "authentication": auth, "service": service or []}})
async def update_did(self, contract: str, did: str, controller=None, vm=None, auth=None, service=None) -> Dict:
upd = {"id": did}
if controller: upd["controller"] = controller
if vm: upd["verification_method"] = vm
if auth: upd["authentication"] = auth
if service: upd["service"] = service
return await self._execute(contract, {"update_did": upd})
def update_did_sync(self, contract: str, did: str, controller=None, vm=None, auth=None, service=None) -> Dict:
upd = {"id": did}
if controller: upd["controller"] = controller
if vm: upd["verification_method"] = vm
if auth: upd["authentication"] = auth
if service: upd["service"] = service
return self._execute_sync(contract, {"update_did": upd})
async def deactivate_did(self, contract: str, did: str) -> Dict:
return await self._execute(contract, {"deactivate_did": {"id": did}})
def deactivate_did_sync(self, contract: str, did: str) -> Dict:
return self._execute_sync(contract, {"deactivate_did": {"id": did}})
# Constitution
async def add_rule(self, contract: str, constitution_id: str, rule_id: str, desc: str, action_type: str, params=None) -> Dict:
return await self._execute(contract, {"add_rule": {"constitution_id": constitution_id, "rule_id": rule_id, "description": desc, "action_type": action_type, "parameters": params or {}}})
def add_rule_sync(self, contract: str, constitution_id: str, rule_id: str, desc: str, action_type: str, params=None) -> Dict:
return self._execute_sync(contract, {"add_rule": {"constitution_id": constitution_id, "rule_id": rule_id, "description": desc, "action_type": action_type, "parameters": params or {}}})
async def remove_rule(self, contract: str, constitution_id: str, rule_id: str) -> Dict:
return await self._execute(contract, {"remove_rule": {"constitution_id": constitution_id, "rule_id": rule_id}})
def remove_rule_sync(self, contract: str, constitution_id: str, rule_id: str) -> Dict:
return self._execute_sync(contract, {"remove_rule": {"constitution_id": constitution_id, "rule_id": rule_id}})
# Micropayment
async def open_channel(self, contract: str, channel_id: str, receiver: str, total_deposit: str, expires_at=None) -> Dict:
return await self._execute(contract, {"open_channel": {"channel_id": channel_id, "receiver": receiver, "total_deposit": total_deposit, "expires_at": expires_at or ""}}, funds=[Coin(denom="umsg", amount=total_deposit)])
def open_channel_sync(self, contract: str, channel_id: str, receiver: str, total_deposit: str, expires_at=None) -> Dict:
return self._execute_sync(contract, {"open_channel": {"channel_id": channel_id, "receiver": receiver, "total_deposit": total_deposit, "expires_at": expires_at or ""}}, funds=[Coin(denom="umsg", amount=total_deposit)])
async def deposit_channel(self, contract: str, channel_id: str, amount: str) -> Dict:
return await self._execute(contract, {"deposit": {"channel_id": channel_id}}, funds=[Coin(denom="umsg", amount=amount)])
def deposit_channel_sync(self, contract: str, channel_id: str, amount: str) -> Dict:
return self._execute_sync(contract, {"deposit": {"channel_id": channel_id}}, funds=[Coin(denom="umsg", amount=amount)])
async def claim_channel(self, contract: str, channel_id: str, amount: str, signature: str) -> Dict:
return await self._execute(contract, {"claim": {"channel_id": channel_id, "amount": amount, "signature": signature}})
def claim_channel_sync(self, contract: str, channel_id: str, amount: str, signature: str) -> Dict:
return self._execute_sync(contract, {"claim": {"channel_id": channel_id, "amount": amount, "signature": signature}})
def __repr__(self) -> str:
return "ContractExecuteClient(client=MSGChainClient)"
11. Agent API 客户端
# 文件: msg_sdk/agent.py
# 18 个 Agent API 端点完整封装
from __future__ import annotations
from typing import Any, Dict, List, Optional
from msg_sdk.client import MSGChainClient
from msg_sdk.errors import AgentNotFoundError, InvalidParamError, SessionExpiredError
from msg_sdk.types import (
AgentQueryResponse, AgentStatusResponse, AgentHistoryEntry,
OracleRequest, OracleResult, WalletInfo, TransferRequest,
TransactionInfo, MPCRequest, MPCResponse, PaymentSessionRequest, Coin,
)
class AgentAPIClient:
"""Agent API 客户端 — 封装全部 18 个 Agent API 端点。
Group 1 — Query: query, get_status, get_history
Group 2 — Events: subscribe, unsubscribe, replay
Group 3 — Oracle: request, get_result
Group 4 — Wallet: create, get_balance, transfer, get_transactions
Group 5 — MPC: sign, submit, aggregate
Group 6 — Payment: create_session, close_session
"""
def __init__(self, client: MSGChainClient, api_key: Optional[str] = None) -> None:
self._client = client
self._api_key = api_key
def _headers(self) -> Dict[str, str]:
h = {"Content-Type": "application/json"}
if self._api_key:
h["X-API-Key"] = self._api_key
return h
async def _post(self, path: str, data: Dict) -> Dict:
client = self._client
await client.connect()
resp = await client._async_client.post(f"{client.config.rest_url}{path}", json=data, headers=self._headers())
resp.raise_for_status()
body = resp.json()
if "code" in body and body["code"] != 0:
from msg_sdk.errors import ERROR_CODE_MAP
cls = ERROR_CODE_MAP.get(body["code"])
if cls:
raise cls(message=body.get("message"), data=body.get("data"))
return body.get("data", body)
async def _get(self, path: str) -> Dict:
client = self._client
await client.connect()
resp = await client._async_client.get(f"{client.config.rest_url}{path}", headers=self._headers())
resp.raise_for_status()
body = resp.json()
if "code" in body and body["code"] != 0:
from msg_sdk.errors import ERROR_CODE_MAP
cls = ERROR_CODE_MAP.get(body["code"])
if cls:
raise cls(message=body.get("message"), data=body.get("data"))
return body.get("data", body)
def _post_sync(self, path: str, data: Dict) -> Dict:
client = self._client
client.connect_sync()
resp = client._sync_client.post(f"{client.config.rest_url}{path}", json=data, headers=self._headers())
resp.raise_for_status()
body = resp.json()
if "code" in body and body["code"] != 0:
from msg_sdk.errors import ERROR_CODE_MAP
cls = ERROR_CODE_MAP.get(body["code"])
if cls:
raise cls(message=body.get("message"), data=body.get("data"))
return body.get("data", body)
def _get_sync(self, path: str) -> Dict:
client = self._client
client.connect_sync()
resp = client._sync_client.get(f"{client.config.rest_url}{path}", headers=self._headers())
resp.raise_for_status()
body = resp.json()
if "code" in body and body["code"] != 0:
from msg_sdk.errors import ERROR_CODE_MAP
cls = ERROR_CODE_MAP.get(body["code"])
if cls:
raise cls(message=body.get("message"), data=body.get("data"))
return body.get("data", body)
# ─── Group 1: Query ─────────────────────────────────────────────────
async def query_agent(self, agent_id: str, query: str, params: Optional[Dict] = None, session_id: Optional[str] = None) -> AgentQueryResponse:
"""POST /api/v1/agent/query — 向 Agent 发送查询。"""
data: Dict = {"agent_id": agent_id, "query": query, "parameters": params or {}}
if session_id:
data["session_id"] = session_id
result = await self._post("/api/v1/agent/query", data)
return AgentQueryResponse.from_dict(result)
def query_agent_sync(self, agent_id: str, query: str, params: Optional[Dict] = None, session_id: Optional[str] = None) -> AgentQueryResponse:
data: Dict = {"agent_id": agent_id, "query": query, "parameters": params or {}}
if session_id:
data["session_id"] = session_id
result = self._post_sync("/api/v1/agent/query", data)
return AgentQueryResponse.from_dict(result)
async def get_agent_status(self, agent_id: str) -> AgentStatusResponse:
"""GET /api/v1/agent/{id}/status — 获取 Agent 状态。"""
result = await self._get(f"/api/v1/agent/{agent_id}/status")
return AgentStatusResponse.from_dict(result)
def get_agent_status_sync(self, agent_id: str) -> AgentStatusResponse:
result = self._get_sync(f"/api/v1/agent/{agent_id}/status")
return AgentStatusResponse.from_dict(result)
async def get_agent_history(self, agent_id: str, limit: int = 20) -> List[AgentHistoryEntry]:
"""GET /api/v1/agent/{id}/history — 获取 Agent 查询历史。"""
result = await self._get(f"/api/v1/agent/{agent_id}/history?limit={limit}")
entries = result if isinstance(result, list) else result.get("entries", [])
return [AgentHistoryEntry.from_dict(e) for e in entries]
def get_agent_history_sync(self, agent_id: str, limit: int = 20) -> List[AgentHistoryEntry]:
result = self._get_sync(f"/api/v1/agent/{agent_id}/history?limit={limit}")
entries = result if isinstance(result, list) else result.get("entries", [])
return [AgentHistoryEntry.from_dict(e) for e in entries]
# ─── Group 2: Events ────────────────────────────────────────────────
async def subscribe_events(self, agent_id: str, event_types: List[str], callback_url: str) -> Dict:
"""POST /api/v1/agent/events/subscribe — 订阅 Agent 事件。"""
return await self._post("/api/v1/agent/events/subscribe", {"agent_id": agent_id, "event_types": event_types, "callback_url": callback_url})
def subscribe_events_sync(self, agent_id: str, event_types: List[str], callback_url: str) -> Dict:
return self._post_sync("/api/v1/agent/events/subscribe", {"agent_id": agent_id, "event_types": event_types, "callback_url": callback_url})
async def unsubscribe_events(self, subscription_id: str) -> Dict:
"""POST /api/v1/agent/events/unsubscribe — 取消订阅。"""
return await self._post("/api/v1/agent/events/unsubscribe", {"subscription_id": subscription_id})
def unsubscribe_events_sync(self, subscription_id: str) -> Dict:
return self._post_sync("/api/v1/agent/events/unsubscribe", {"subscription_id": subscription_id})
async def replay_events(self, agent_id: str, from_timestamp: str, to_timestamp: str) -> List[Dict]:
"""POST /api/v1/agent/events/replay — 重放历史事件。"""
result = await self._post("/api/v1/agent/events/replay", {"agent_id": agent_id, "from": from_timestamp, "to": to_timestamp})
return result if isinstance(result, list) else result.get("events", [])
def replay_events_sync(self, agent_id: str, from_timestamp: str, to_timestamp: str) -> List[Dict]:
result = self._post_sync("/api/v1/agent/events/replay", {"agent_id": agent_id, "from": from_timestamp, "to": to_timestamp})
return result if isinstance(result, list) else result.get("events", [])
# ─── Group 3: Oracle ────────────────────────────────────────────────
async def oracle_request(self, request: OracleRequest) -> Dict:
"""POST /api/v1/agent/oracle/request — 发起预言机请求。"""
return await self._post("/api/v1/agent/oracle/request", request.to_dict())
def oracle_request_sync(self, request: OracleRequest) -> Dict:
return self._post_sync("/api/v1/agent/oracle/request", request.to_dict())
async def get_oracle_result(self, request_id: str) -> OracleResult:
"""GET /api/v1/agent/oracle/result/{id} — 获取预言机结果。"""
result = await self._get(f"/api/v1/agent/oracle/result/{request_id}")
return OracleResult.from_dict(result)
def get_oracle_result_sync(self, request_id: str) -> OracleResult:
result = self._get_sync(f"/api/v1/agent/oracle/result/{request_id}")
return OracleResult.from_dict(result)
# ─── Group 4: Wallet ────────────────────────────────────────────────
async def create_agent_wallet(self, agent_id: str) -> WalletInfo:
"""POST /api/v1/agent/wallet/create — 为 Agent 创建钱包。"""
result = await self._post("/api/v1/agent/wallet/create", {"agent_id": agent_id})
return WalletInfo.from_dict(result)
def create_agent_wallet_sync(self, agent_id: str) -> WalletInfo:
result = self._post_sync("/api/v1/agent/wallet/create", {"agent_id": agent_id})
return WalletInfo.from_dict(result)
async def get_agent_wallet_balance(self, wallet_id: str) -> Coin:
"""GET /api/v1/agent/wallet/{id}/balance — 查询 Agent 钱包余额。"""
result = await self._get(f"/api/v1/agent/wallet/{wallet_id}/balance")
return Coin.from_dict(result.get("balance", result))
def get_agent_wallet_balance_sync(self, wallet_id: str) -> Coin:
result = self._get_sync(f"/api/v1/agent/wallet/{wallet_id}/balance")
return Coin.from_dict(result.get("balance", result))
async def agent_transfer(self, wallet_id: str, to_address: str, amount: str, denom: str = "umsg", memo: str = "") -> TransactionInfo:
"""POST /api/v1/agent/wallet/transfer — Agent 发起转账。"""
result = await self._post("/api/v1/agent/wallet/transfer", {"wallet_id": wallet_id, "to_address": to_address, "amount": amount, "denom": denom, "memo": memo})
return TransactionInfo.from_dict(result)
def agent_transfer_sync(self, wallet_id: str, to_address: str, amount: str, denom: str = "umsg", memo: str = "") -> TransactionInfo:
result = self._post_sync("/api/v1/agent/wallet/transfer", {"wallet_id": wallet_id, "to_address": to_address, "amount": amount, "denom": denom, "memo": memo})
return TransactionInfo.from_dict(result)
async def get_agent_transactions(self, wallet_id: str, limit: int = 20) -> List[TransactionInfo]:
"""GET /api/v1/agent/wallet/transactions — 查询交易历史。"""
result = await self._get(f"/api/v1/agent/wallet/transactions?wallet_id={wallet_id}&limit={limit}")
txs = result if isinstance(result, list) else result.get("transactions", [])
return [TransactionInfo.from_dict(t) for t in txs]
def get_agent_transactions_sync(self, wallet_id: str, limit: int = 20) -> List[TransactionInfo]:
result = self._get_sync(f"/api/v1/agent/wallet/transactions?wallet_id={wallet_id}&limit={limit}")
txs = result if isinstance(result, list) else result.get("transactions", [])
return [TransactionInfo.from_dict(t) for t in txs]
# ─── Group 5: MPC ───────────────────────────────────────────────────
async def mpc_sign(self, request: MPCRequest) -> MPCResponse:
"""POST /api/v1/agent/mpc/sign — 发起 MPC 签名。"""
result = await self._post("/api/v1/agent/mpc/sign", request.to_dict())
return MPCResponse.from_dict(result)
def mpc_sign_sync(self, request: MPCRequest) -> MPCResponse:
result = self._post_sync("/api/v1/agent/mpc/sign", request.to_dict())
return MPCResponse.from_dict(result)
async def mpc_submit(self, request_id: str, partial_sig: str) -> MPCResponse:
"""POST /api/v1/agent/mpc/submit — 提交部分签名。"""
result = await self._post("/api/v1/agent/mpc/submit", {"request_id": request_id, "partial_signature": partial_sig})
return MPCResponse.from_dict(result)
def mpc_submit_sync(self, request_id: str, partial_sig: str) -> MPCResponse:
result = self._post_sync("/api/v1/agent/mpc/submit", {"request_id": request_id, "partial_signature": partial_sig})
return MPCResponse.from_dict(result)
async def mpc_aggregate(self, request_id: str) -> MPCResponse:
"""POST /api/v1/agent/mpc/aggregate — 聚合签名。"""
result = await self._post("/api/v1/agent/mpc/aggregate", {"request_id": request_id})
return MPCResponse.from_dict(result)
def mpc_aggregate_sync(self, request_id: str) -> MPCResponse:
result = self._post_sync("/api/v1/agent/mpc/aggregate", {"request_id": request_id})
return MPCResponse.from_dict(result)
# ─── Group 6: Payment ───────────────────────────────────────────────
async def create_payment_session(self, request: PaymentSessionRequest) -> Dict:
"""POST /api/v1/agent/payment/create-session — 创建支付会话。"""
return await self._post("/api/v1/agent/payment/create-session", request.to_dict())
def create_payment_session_sync(self, request: PaymentSessionRequest) -> Dict:
return self._post_sync("/api/v1/agent/payment/create-session", request.to_dict())
async def close_payment_session(self, session_id: str) -> Dict:
"""POST /api/v1/agent/payment/close-session — 关闭支付会话。"""
return await self._post("/api/v1/agent/payment/close-session", {"session_id": session_id})
def close_payment_session_sync(self, session_id: str) -> Dict:
return self._post_sync("/api/v1/agent/payment/close-session", {"session_id": session_id})
def __repr__(self) -> str:
return "AgentAPIClient(api_key=***)"
12. 治理客户端
# 文件: msg_sdk/governance.py
from __future__ import annotations
from typing import Any, Dict, List, Optional
from msg_sdk.client import MSGChainClient
from msg_sdk.errors import UnauthorizedError, InvalidParamError
from msg_sdk.types import Coin, Proposal, Vote, TallyResult
class GovernanceClient:
"""治理客户端 — 提案、投票、国库查询。"""
def __init__(self, client: MSGChainClient) -> None:
self._client = client
async def get_proposals(self, status: Optional[str] = None) -> List[Proposal]:
params = {}
if status:
params["proposal_status"] = status
body = await self._client._request("GET", "/cosmos/gov/v1beta1/proposals", params=params)
return [Proposal.from_dict(p) for p in body.get("proposals", [])]
def get_proposals_sync(self, status: Optional[str] = None) -> List[Proposal]:
params = {}
if status:
params["proposal_status"] = status
body = self._client._request_sync("GET", "/cosmos/gov/v1beta1/proposals", params=params)
return [Proposal.from_dict(p) for p in body.get("proposals", [])]
async def get_proposal(self, proposal_id: int) -> Proposal:
body = await self._client._request("GET", f"/cosmos/gov/v1beta1/proposals/{proposal_id}")
return Proposal.from_dict(body.get("proposal", body))
def get_proposal_sync(self, proposal_id: int) -> Proposal:
body = self._client._request_sync("GET", f"/cosmos/gov/v1beta1/proposals/{proposal_id}")
return Proposal.from_dict(body.get("proposal", body))
async def get_proposal_tally(self, proposal_id: int) -> TallyResult:
body = await self._client._request("GET", f"/cosmos/gov/v1beta1/proposals/{proposal_id}/tally")
return TallyResult.from_dict(body.get("tally", body))
def get_proposal_tally_sync(self, proposal_id: int) -> TallyResult:
body = self._client._request_sync("GET", f"/cosmos/gov/v1beta1/proposals/{proposal_id}/tally")
return TallyResult.from_dict(body.get("tally", body))
async def get_votes(self, proposal_id: int) -> List[Vote]:
body = await self._client._request("GET", f"/cosmos/gov/v1beta1/proposals/{proposal_id}/votes")
return [Vote.from_dict(v) for v in body.get("votes", [])]
def get_votes_sync(self, proposal_id: int) -> List[Vote]:
body = self._client._request_sync("GET", f"/cosmos/gov/v1beta1/proposals/{proposal_id}/votes")
return [Vote.from_dict(v) for v in body.get("votes", [])]
async def get_deposits(self, proposal_id: int) -> List[Dict]:
body = await self._client._request("GET", f"/cosmos/gov/v1beta1/proposals/{proposal_id}/deposits")
return body.get("deposits", [])
def get_deposits_sync(self, proposal_id: int) -> List[Dict]:
body = self._client._request_sync("GET", f"/cosmos/gov/v1beta1/proposals/{proposal_id}/deposits")
return body.get("deposits", [])
async def get_vote(self, proposal_id: int, voter: str) -> Vote:
body = await self._client._request("GET", f"/cosmos/gov/v1beta1/proposals/{proposal_id}/votes/{voter}")
return Vote.from_dict(body.get("vote", body))
def get_vote_sync(self, proposal_id: int, voter: str) -> Vote:
body = self._client._request_sync("GET", f"/cosmos/gov/v1beta1/proposals/{proposal_id}/votes/{voter}")
return Vote.from_dict(body.get("vote", body))
async def get_params(self, param_type: str = "voting") -> Dict:
body = await self._client._request("GET", f"/cosmos/gov/v1beta1/params/{param_type}")
return body.get("voting_params", body)
def get_params_sync(self, param_type: str = "voting") -> Dict:
body = self._client._request_sync("GET", f"/cosmos/gov/v1beta1/params/{param_type}")
return body.get("voting_params", body)
async def get_treasury_balance(self) -> Coin:
"""查询国库余额 (社区资金池)。"""
body = await self._client._request("GET", "/cosmos/distribution/v1beta1/community_pool")
pool = body.get("pool", [])
for c in pool:
if c.get("denom") == "umsg":
return Coin.from_dict(c)
return Coin(denom="umsg", amount="0")
def get_treasury_balance_sync(self) -> Coin:
body = self._client._request_sync("GET", "/cosmos/distribution/v1beta1/community_pool")
pool = body.get("pool", [])
for c in pool:
if c.get("denom") == "umsg":
return Coin.from_dict(c)
return Coin(denom="umsg", amount="0")
async def submit_proposal(self, title: str, description: str, initial_deposit: str, denom: str = "umsg") -> Dict[str, Any]:
if self._client.wallet is None:
raise UnauthorizedError(message="wallet not set")
proposer = self._client.wallet.address
info = await self._client.get_account_info(proposer)
msg = {"type": "cosmos-sdk/MsgSubmitProposal", "value": {"content": {"type": "cosmos-sdk/TextProposal", "value": {"title": title, "description": description}}, "initial_deposit": [{"denom": denom, "amount": initial_deposit}], "proposer": proposer}}
tx = {"chain_id": self._client.config.chain_id, "account_number": info["account_number"], "sequence": info["sequence"], "msgs": [msg], "memo": "", "fee": {"gas": "300000", "amount": [{"denom": "umsg", "amount": "7500"}]}}
signed = self._client.wallet.sign_transaction(tx)
from msg_sdk.utils import encode_amino
tx_bytes = encode_amino(signed).hex()
result = await self._client.broadcast_transaction(tx_bytes)
tx_hash = result.get("tx_response", result).get("txhash", "")
return await self._client.wait_for_tx(tx_hash) if tx_hash else result
def submit_proposal_sync(self, title: str, description: str, initial_deposit: str, denom: str = "umsg") -> Dict[str, Any]:
if self._client.wallet is None:
raise UnauthorizedError(message="wallet not set")
proposer = self._client.wallet.address
info = self._client.get_account_info_sync(proposer)
msg = {"type": "cosmos-sdk/MsgSubmitProposal", "value": {"content": {"type": "cosmos-sdk/TextProposal", "value": {"title": title, "description": description}}, "initial_deposit": [{"denom": denom, "amount": initial_deposit}], "proposer": proposer}}
tx = {"chain_id": self._client.config.chain_id, "account_number": info["account_number"], "sequence": info["sequence"], "msgs": [msg], "memo": "", "fee": {"gas": "300000", "amount": [{"denom": "umsg", "amount": "7500"}]}}
signed = self._client.wallet.sign_transaction(tx)
from msg_sdk.utils import encode_amino
tx_bytes = encode_amino(signed).hex()
result = self._client.broadcast_transaction_sync(tx_bytes)
tx_hash = result.get("tx_response", result).get("txhash", "")
return self._client.wait_for_tx_sync(tx_hash) if tx_hash else result
async def vote(self, proposal_id: int, option: str) -> Dict[str, Any]:
if option not in ("yes", "no", "abstain", "no_with_veto"):
raise InvalidParamError(message=f"invalid vote option: {option}")
if self._client.wallet is None:
raise UnauthorizedError(message="wallet not set")
voter = self._client.wallet.address
info = await self._client.get_account_info(voter)
msg = {"type": "cosmos-sdk/MsgVote", "value": {"proposal_id": str(proposal_id), "voter": voter, "option": option}}
tx = {"chain_id": self._client.config.chain_id, "account_number": info["account_number"], "sequence": info["sequence"], "msgs": [msg], "memo": "", "fee": {"gas": "150000", "amount": [{"denom": "umsg", "amount": "3750"}]}}
signed = self._client.wallet.sign_transaction(tx)
from msg_sdk.utils import encode_amino
tx_bytes = encode_amino(signed).hex()
result = await self._client.broadcast_transaction(tx_bytes)
tx_hash = result.get("tx_response", result).get("txhash", "")
return await self._client.wait_for_tx(tx_hash) if tx_hash else result
def vote_sync(self, proposal_id: int, option: str) -> Dict[str, Any]:
if option not in ("yes", "no", "abstain", "no_with_veto"):
raise InvalidParamError(message=f"invalid vote option: {option}")
if self._client.wallet is None:
raise UnauthorizedError(message="wallet not set")
voter = self._client.wallet.address
info = self._client.get_account_info_sync(voter)
msg = {"type": "cosmos-sdk/MsgVote", "value": {"proposal_id": str(proposal_id), "voter": voter, "option": option}}
tx = {"chain_id": self._client.config.chain_id, "account_number": info["account_number"], "sequence": info["sequence"], "msgs": [msg], "memo": "", "fee": {"gas": "150000", "amount": [{"denom": "umsg", "amount": "3750"}]}}
signed = self._client.wallet.sign_transaction(tx)
from msg_sdk.utils import encode_amino
tx_bytes = encode_amino(signed).hex()
result = self._client.broadcast_transaction_sync(tx_bytes)
tx_hash = result.get("tx_response", result).get("txhash", "")
return self._client.wait_for_tx_sync(tx_hash) if tx_hash else result
def __repr__(self) -> str:
return "GovernanceClient(client=MSGChainClient)"
13. WebSocket 事件客户端
# 文件: msg_sdk/events.py
from __future__ import annotations
import asyncio
import json
from typing import Any, Callable, Dict, List, Optional
from msg_sdk.types import ChainEvent, AgentEvent
class EventClient:
"""WebSocket 事件客户端 — 实时订阅链上和 Agent 事件。
用法:
async with EventClient() as ec:
async for event in ec.subscribe_chain_events(["message.action='send'"]):
print(event)
"""
def __init__(self, ws_url: str = "ws://localhost:26657/websocket", reconnect: bool = True) -> None:
self._ws_url = ws_url
self._reconnect = reconnect
self._ws = None
self._running = False
self._handlers: Dict[str, List[Callable]] = {"chain": [], "agent": [], "*": []}
async def connect(self) -> None:
import websockets
self._ws = await websockets.connect(self._ws_url)
self._running = True
async def disconnect(self) -> None:
self._running = False
if self._ws:
await self._ws.close()
self._ws = None
async def __aenter__(self) -> EventClient:
await self.connect()
return self
async def __aexit__(self, *args: Any) -> None:
await self.disconnect()
def on_chain_event(self, handler: Callable) -> None:
self._handlers["chain"].append(handler)
def on_agent_event(self, handler: Callable) -> None:
self._handlers["agent"].append(handler)
def on_any_event(self, handler: Callable) -> None:
self._handlers["*"].append(handler)
async def subscribe_chain_events(self, queries: List[str]) -> AsyncIterator[ChainEvent]:
"""订阅链上事件。
Args:
queries: Tendermint 查询条件列表,例如 ["tm.event='Tx'"]
Yields:
ChainEvent 对象
"""
for q in queries:
sub_req = {"jsonrpc": "2.0", "method": "subscribe", "params": {"query": q}, "id": 1}
await self._ws.send(json.dumps(sub_req))
resp = json.loads(await self._ws.recv())
if "error" in resp:
raise ConnectionError(f"subscription failed: {resp['error']}")
while self._running:
try:
msg = json.loads(await self._ws.recv())
if "result" in msg and "data" in msg["result"]:
data = msg["result"]["data"]
event = ChainEvent.from_dict(data.get("value", {}))
yield event
for h in self._handlers.get("*", []):
h(event)
for h in self._handlers.get("chain", []):
h(event)
except Exception:
if self._reconnect:
await self.reconnect()
else:
raise
async def listen_forever(self, handler: Optional[Callable] = None) -> None:
"""持续监听事件。"""
while self._running:
try:
msg = json.loads(await self._ws.recv())
if handler:
handler(msg)
for h in self._handlers.get("*", []):
h(msg)
except Exception as e:
if self._reconnect:
await self.reconnect()
else:
raise
async def reconnect(self) -> None:
await self.disconnect()
await asyncio.sleep(1)
await self.connect()
async def unsubscribe_all(self) -> None:
if self._ws:
unsub_req = {"jsonrpc": "2.0", "method": "unsubscribe_all", "params": {}, "id": 1}
await self._ws.send(json.dumps(unsub_req))
def __repr__(self) -> str:
return f"EventClient(ws_url={self._ws_url})"
from typing import AsyncIterator
14. 完整示例脚本
本节提供 3 个可直接运行的完整示例脚本,展示 SDK 的核心用法。
示例 1: Agent 注册演示
#!/usr/bin/env python3
# 文件: examples/agent_registration_demo.py
# Agent 注册与查询完整示例
import asyncio
import sys
import os
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from msg_sdk.client import MSGChainClient, ChainConfig
from msg_sdk.wallet import Wallet
from msg_sdk.contract import ContractQueryClient, ContractExecuteClient
from msg_sdk.types import Coin, AgentInfo
from msg_sdk.errors import AgentExistsError, AgentNotFoundError
async def main():
print("=" * 60)
print("MSG Chain Agent 注册演示")
print("=" * 60)
# 1. 初始化客户端
config = ChainConfig(
chain_id="msg-chain-1",
rpc_url="http://localhost:26657",
rest_url="http://localhost:1317",
)
# 2. 生成或加载钱包
wallet_file = "demo_wallet.json"
if os.path.exists(wallet_file):
import json
with open(wallet_file) as f:
data = json.load(f)
wallet = Wallet.from_mnemonic(data["mnemonic"])
print(f"[钱包] 已加载: {wallet.address}")
else:
wallet = Wallet.generate()
import json
with open(wallet_file, "w") as f:
json.dump({"mnemonic": wallet.mnemonic, "address": wallet.address}, f)
print(f"[钱包] 已生成: {wallet.address}")
print(f"[钱包] 助记词: {wallet.mnemonic[:50]}...")
# 3. 创建客户端并设置钱包
async with MSGChainClient(config=config, wallet=wallet) as client:
query = ContractQueryClient(client)
exec_client = ContractExecuteClient(client)
# 合约地址 (实际部署时替换)
REGISTRY_CONTRACT = "msg14hj2tavq8fpesdwxxcu44rty3hh90vhujrvcm6"
# 4. 检查余额
balances = await client.get_balance(wallet.address)
umsg_balance = sum(int(c.amount) for c in balances if c.denom == "umsg")
print(f"[余额] {wallet.address}: {umsg_balance} umsg")
if umsg_balance == 0:
print("[错误] 余额不足,请先通过水龙头获取测试代币")
return
# 5. 注册 Agent
agent_id = "demo-agent-001"
print(f"\n[注册] 注册 Agent: {agent_id}")
try:
result = await exec_client.register_agent(
contract_address=REGISTRY_CONTRACT,
agent_id=agent_id,
name="Demo Analysis Agent",
description="用于演示的 AI 分析 Agent",
constitution_id="constitution-default",
metadata={"version": "1.0.0", "type": "analysis"},
)
print(f"[成功] Agent 注册成功!")
print(f" TxHash: {result.get('txhash', 'N/A')}")
print(f" 高度: {result.get('height', 'N/A')}")
except AgentExistsError:
print(f"[信息] Agent {agent_id} 已存在,跳过注册")
except Exception as e:
print(f"[错误] 注册失败: {e}")
return
# 6. 查询 Agent
print(f"\n[查询] 查询 Agent: {agent_id}")
try:
agent = await query.get_agent(REGISTRY_CONTRACT, agent_id)
print(f" ID: {agent.id}")
print(f" 名称: {agent.name}")
print(f" 描述: {agent.description}")
print(f" 拥有者: {agent.owner}")
print(f" 状态: {agent.status}")
print(f" 宪法: {agent.constitution_id}")
except AgentNotFoundError:
print("[错误] Agent 未找到")
# 7. 列出所有 Agent
print("\n[列表] 所有已注册 Agent:")
try:
agents = await query.list_agents(REGISTRY_CONTRACT, limit=10)
if agents:
for a in agents:
print(f" - {a.id}: {a.name} ({a.status})")
else:
print(" (无)")
except Exception as e:
print(f" [错误] {e}")
# 8. 更新 Agent 状态
print(f"\n[更新] 暂停 Agent: {agent_id}")
try:
result = await exec_client.set_agent_status(
contract_address=REGISTRY_CONTRACT,
agent_id=agent_id,
status="paused",
)
print(f" 状态已更新为 paused")
# 验证
agent = await query.get_agent(REGISTRY_CONTRACT, agent_id)
print(f" 当前状态: {agent.status}")
except Exception as e:
print(f" [错误] {e}")
print("\n" + "=" * 60)
print("演示完成")
print("=" * 60)
if __name__ == "__main__":
asyncio.run(main())
示例 2: 支付会话演示
#!/usr/bin/env python3
# 文件: examples/payment_session_demo.py
# 支付会话创建、充值与释放完整示例
import asyncio
import sys
import os
import time
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from msg_sdk.client import MSGChainClient, ChainConfig
from msg_sdk.wallet import Wallet
from msg_sdk.contract import ContractQueryClient, ContractExecuteClient
from msg_sdk.types import PaymentSession, Coin
from msg_sdk.errors import SessionExpiredError
async def main():
print("=" * 60)
print("MSG Chain 支付会话演示")
print("=" * 60)
config = ChainConfig(rest_url="http://localhost:1317")
# 使用预设的两个钱包 (payer 和 payee)
payer_wallet = Wallet.generate()
payee_wallet = Wallet.generate()
print(f"[付款方] {payer_wallet.address}")
print(f"[收款方] {payee_wallet.address}")
PAYMENT_CONTRACT = "msg14hj2tavq8fpesdwxxcu44rty3hh90vhujrvcm6"
async with MSGChainClient(config=config, wallet=payer_wallet) as client:
query = ContractQueryClient(client)
exec_client = ContractExecuteClient(client)
# 1. 创建支付会话
session_id = f"session-{int(time.time())}"
amount_umsg = "1000000000000000000" # 1 MSG
fee_umsg = "50000000000000000" # 0.05 MSG
print(f"\n[创建] 支付会话: {session_id}")
print(f" 金额: {amount_umsg} umsg (1 MSG)")
print(f" 手续费: {fee_umsg} umsg")
try:
result = await exec_client.create_payment_session(
contract_address=PAYMENT_CONTRACT,
session_id=session_id,
payer=payer_wallet.address,
payee=payee_wallet.address,
amount=amount_umsg,
fee=fee_umsg,
duration_seconds=3600,
metadata={"description": "AI Agent 服务费"},
)
print(f"[成功] 支付会话已创建")
txhash = result.get('txhash', 'N/A')
print(f" TxHash: {txhash}")
except Exception as e:
print(f"[错误] 创建失败: {e}")
return
# 2. 查询支付会话
print(f"\n[查询] 支付会话: {session_id}")
try:
session = await query.get_payment_session(PAYMENT_CONTRACT, session_id)
print(f" 状态: {session.status}")
print(f" 付款方: {session.payer}")
print(f" 收款方: {session.payee}")
if session.amount:
print(f" 金额: {session.amount.amount} {session.amount.denom}")
print(f" 创建时间: {session.created_at}")
print(f" 过期时间: {session.expires_at}")
except Exception as e:
print(f"[错误] 查询失败: {e}")
return
# 3. 充值会话 (使用付款方钱包)
topup_amount = "500000000000000000" # 0.5 MSG
print(f"\n[充值] 向会话添加: {topup_amount} umsg")
try:
result = await exec_client.fund_session(
contract_address=PAYMENT_CONTRACT,
session_id=session_id,
amount=topup_amount,
)
print(f"[成功] 充值完成")
except Exception as e:
print(f"[错误] 充值失败: {e}")
# 4. 检查更新后的会话
session = await query.get_payment_session(PAYMENT_CONTRACT, session_id)
print(f"\n[状态] 会话状态: {session.status}")
# 5. 释放支付
print(f"\n[释放] 释放支付给收款方...")
try:
result = await exec_client.release_payment(
contract_address=PAYMENT_CONTRACT,
session_id=session_id,
)
print(f"[成功] 支付已释放")
except Exception as e:
print(f"[错误] 释放失败: {e}")
# 6. 再次查询确认状态
session = await query.get_payment_session(PAYMENT_CONTRACT, session_id)
print(f"\n[最终状态] {session.status}")
print("\n" + "=" * 60)
print("支付演示完成")
print("=" * 60)
if __name__ == "__main__":
asyncio.run(main())
示例 3: 治理演示
15. 测试文件
test_errors.py
# 文件: tests/test_errors.py
import pytest
from msg_sdk.errors import (
MSGChainError, OK, StubError, NotImplementedErrorCode,
UnauthorizedError, InsufficientFundsError, ContractFailedError,
QueryTimeoutError, InvalidParamError, DAOTimelockError,
HighValueError, AgentNotFoundError, AgentExistsError,
ConstitutionViolationError, SessionExpiredError, SessionLimitError,
InvalidSignatureError, DuplicateNonceError, RateLimitError,
ERROR_CODE_MAP, raise_for_code, extract_error,
)
def test_error_code_map_completeness():
for code in range(18):
assert code in ERROR_CODE_MAP, f"Missing error code {code}"
def test_ok_no_exception():
raise_for_code(0)
def test_all_error_codes_raise():
for code, exc_cls in ERROR_CODE_MAP.items():
if code == 0:
continue
with pytest.raises(exc_cls):
raise_for_code(code)
def test_custom_message():
with pytest.raises(UnauthorizedError, match="custom error"):
raise_for_code(3, message="custom error")
def test_extract_error():
with pytest.raises(InsufficientFundsError):
extract_error({"code": 4, "message": "not enough tokens"})
def test_extract_error_ok():
extract_error({"code": 0, "data": {"balance": "1000"}})
def test_error_to_dict():
err = InsufficientFundsError(message="low balance", data={"balance": "0"})
d = err.to_dict()
assert d["code"] == 4
assert d["message"] == "low balance"
assert d["data"]["balance"] == "0"
def test_unknown_code():
with pytest.raises(MSGChainError):
raise_for_code(99)
def test_error_inheritance():
for _, exc_cls in ERROR_CODE_MAP.items():
assert issubclass(exc_cls, MSGChainError)
def test_rate_limit():
with pytest.raises(RateLimitError):
raise_for_code(17)
test_wallet.py
# 文件: tests/test_wallet.py
import pytest
from msg_sdk.wallet import Wallet
class TestWallet:
def test_generate(self):
w = Wallet.generate()
assert w.address.startswith("msg1")
assert len(w.mnemonic.split()) == 24
assert len(w.private_key) == 32
assert len(w.public_key) == 32
def test_from_mnemonic(self):
w1 = Wallet.generate()
mnemonic = w1.mnemonic
w2 = Wallet.from_mnemonic(mnemonic)
assert w1.address == w2.address
assert w1.private_key == w2.private_key
assert w1.public_key == w2.public_key
def test_validate_address(self):
w = Wallet.generate()
assert Wallet.validate_address(w.address)
assert not Wallet.validate_address("invalid_address")
assert not Wallet.validate_address("cosmos1...")
def test_sign_and_verify(self):
w = Wallet.generate()
message = b"hello msg chain"
sig = w.sign(message)
assert len(sig) == 64
assert w.verify(message, sig)
assert not w.verify(b"tampered", sig)
def test_sign_hex(self):
w = Wallet.generate()
sig_hex = w.sign_hex("68656c6c6f")
assert isinstance(sig_hex, str)
assert len(sig_hex) == 128
def test_from_private_key(self):
import os
w = Wallet.generate()
w2 = Wallet.from_private_key(w.get_private_key_hex())
assert w.address == w2.address
assert w.public_key == w2.public_key
def test_seed_derivation(self):
w = Wallet.generate()
d = w.to_dict()
assert "address" in d
assert "public_key" in d
def test_sign_transaction(self):
w = Wallet.generate()
tx = {"chain_id": "msg-chain-1", "account_number": "0", "sequence": "0", "msgs": [], "memo": "", "fee": {"gas": "200000", "amount": [{"denom": "umsg", "amount": "5000"}]}}
signed = w.sign_transaction(tx)
assert "signatures" in signed
assert len(signed["signatures"]) == 1
assert signed["signatures"][0]["pub_key"]["type"] == "tendermint/PubKeyEd25519"
test_bank.py
# 文件: tests/test_bank.py
import pytest
from unittest.mock import AsyncMock, MagicMock, patch
from msg_sdk.bank import BankClient
from msg_sdk.client import MSGChainClient, ChainConfig
from msg_sdk.types import Coin
class TestBankClient:
@pytest.fixture
def mock_client(self):
client = MagicMock(spec=MSGChainClient)
client.get_balance = AsyncMock(return_value=[Coin(denom="umsg", amount="1000000")])
client.get_balance_sync = MagicMock(return_value=[Coin(denom="umsg", amount="1000000")])
client.get_supply = AsyncMock(return_value=[Coin(denom="umsg", amount="1000000000000")])
client.get_supply_sync = MagicMock(return_value=[Coin(denom="umsg", amount="1000000000000")])
return client
@pytest.mark.asyncio
async def test_get_balance(self, mock_client):
bank = BankClient(mock_client)
coin = await bank.get_balance("msg1...")
assert coin.denom == "umsg"
assert coin.amount == "1000000"
def test_get_balance_sync(self, mock_client):
bank = BankClient(mock_client)
coin = bank.get_balance_sync("msg1...")
assert coin.denom == "umsg"
@pytest.mark.asyncio
async def test_get_supply(self, mock_client):
bank = BankClient(mock_client)
coin = await bank.get_supply()
assert coin.denom == "umsg"
@pytest.mark.asyncio
async def test_get_balance_umsg(self, mock_client):
bank = BankClient(mock_client)
from decimal import Decimal
bal = await bank.get_balance_umsg("msg1...")
assert bal == Decimal("1000000")
@pytest.mark.asyncio
async def test_transfer(self, mock_client):
mock_client.send_tokens = AsyncMock(return_value={"txhash": "abc123"})
bank = BankClient(mock_client)
result = await bank.transfer("msg1...", "1000")
assert result["txhash"] == "abc123"
test_staking.py
# 文件: tests/test_staking.py
import pytest
from unittest.mock import AsyncMock, MagicMock
from msg_sdk.staking import StakingClient
from msg_sdk.client import MSGChainClient
from msg_sdk.types import Validator, DelegationResponse, RewardsResponse
class TestStakingClient:
@pytest.fixture
def mock_client(self):
c = MagicMock(spec=MSGChainClient)
c._request = AsyncMock(return_value={"validators": [{"operator_address": "msgvaloper1...", "consensus_pubkey": {}, "jailed": False, "status": 3}]})
c._request_sync = MagicMock(return_value={"validators": [{"operator_address": "msgvaloper1...", "consensus_pubkey": {}, "jailed": False, "status": 3}]})
return c
@pytest.mark.asyncio
async def test_get_validators(self, mock_client):
s = StakingClient(mock_client)
vals = await s.get_validators()
assert len(vals) > 0
assert vals[0].operator_address == "msgvaloper1..."
def test_get_validators_sync(self, mock_client):
s = StakingClient(mock_client)
vals = s.get_validators_sync()
assert len(vals) > 0
test_agent_api.py
# 文件: tests/test_agent_api.py
import pytest
from unittest.mock import AsyncMock, MagicMock, patch
from msg_sdk.agent import AgentAPIClient
from msg_sdk.client import MSGChainClient
class TestAgentAPIClient:
@pytest.fixture
def mock_client(self):
c = MagicMock(spec=MSGChainClient)
c.config.rest_url = "http://localhost:1317"
c._async_client = MagicMock()
c._sync_client = MagicMock()
return c
@pytest.mark.asyncio
async def test_query_agent(self, mock_client):
mock_client._async_client.post = AsyncMock()
mock_client._async_client.post.return_value.json = MagicMock(return_value={"data": {"agent_id": "agent-1", "response": "hello", "confidence": 0.95, "latency_ms": 50}})
mock_client._async_client.post.return_value.raise_for_status = MagicMock()
api = AgentAPIClient(mock_client, api_key="test-key")
resp = await api.query_agent("agent-1", "hello")
assert resp.agent_id == "agent-1"
assert resp.response == "hello"
assert resp.confidence == 0.95
test_client.py
# 文件: tests/test_client.py
import pytest
from msg_sdk.client import MSGChainClient, ChainConfig, GasFee
from msg_sdk.types import Coin
class TestChainConfig:
def test_defaults(self):
c = ChainConfig()
assert c.chain_id == "msg-chain-1"
assert c.bech32_prefix == "msg"
assert c.coin_type == 118
assert c.msg_decimals == 18
assert c.gas_low == "1000000000"
assert c.gas_avg == "1000000000"
assert c.gas_high == "1000000000"
assert c.gas_denom == "attoMSG"
assert c.block_time_seconds == 5
assert c.rpc_url == "http://localhost:26657"
assert c.rest_url == "http://localhost:1317"
assert c.grpc_endpoint == "localhost:9090"
class TestGasFee:
def test_calculate(self):
fee = GasFee.calculate(200000, "1000000000")
assert fee.gas_limit == 200000
assert fee.gas_price == "1000000000"
assert isinstance(fee.fee, Coin)
assert fee.fee.denom == "attoMSG"
def test_calculate_zero(self):
fee = GasFee.calculate(0)
assert fee.fee.amount == "0"
class TestMSGChainClient:
def test_init(self):
client = MSGChainClient()
assert client.config.chain_id == "msg-chain-1"
assert client.config.rest_url == "http://localhost:1317"
assert client.wallet is None
def test_estimate_gas_simple(self):
client = MSGChainClient()
fee = client.estimate_gas(1, "simple")
assert fee.gas_limit == 100000
def test_estimate_gas_complex(self):
client = MSGChainClient()
fee = client.estimate_gas(2, "complex")
assert fee.gas_limit == 600000
def test_estimate_gas_contract(self):
client = MSGChainClient()
fee = client.estimate_gas(1, "contract")
assert fee.gas_limit == 500000
def test_client_repr(self):
client = MSGChainClient()
assert "msg-chain-1" in repr(client)
assert "localhost:1317" in repr(client)
conftest.py
# 文件: tests/conftest.py
import pytest
@pytest.fixture(autouse=True)
def setup_test_env():
"""测试环境设置。"""
pass
test_contract.py
# 文件: tests/test_contract.py
import pytest
from unittest.mock import AsyncMock, MagicMock
from msg_sdk.contract import ContractQueryClient, ContractExecuteClient
from msg_sdk.client import MSGChainClient
from msg_sdk.types import AgentInfo, PaymentSession
class TestContractQuery:
@pytest.fixture
def mock_client(self):
c = MagicMock(spec=MSGChainClient)
c.query_smart = AsyncMock()
c.query_smart_sync = MagicMock()
return c
@pytest.mark.asyncio
async def test_get_agent(self, mock_client):
mock_client.query_smart.return_value.data = {"agent": {"id": "a1", "name": "test-agent", "description": "desc", "owner": "msg1...", "status": "active", "constitution_id": "c1", "created_at": "2024-01-01"}}
q = ContractQueryClient(mock_client)
agent = await q.get_agent("msg1...", "a1")
assert agent.id == "a1"
assert agent.name == "test-agent"
class TestContractExecute:
@pytest.fixture
def mock_client(self):
c = MagicMock(spec=MSGChainClient)
c.wallet = MagicMock()
c.wallet.address = "msg1..."
c.config.chain_id = "msg-chain-1"
c.config.gas_avg = "1000000000"
c.get_account_info = AsyncMock(return_value={"account_number": 0, "sequence": 0})
c.broadcast_transaction = AsyncMock(return_value={"tx_response": {"txhash": "abc"}})
c.wait_for_tx = AsyncMock(return_value={"code": 0, "txhash": "abc"})
return c
@pytest.mark.asyncio
async def test_register_agent(self, mock_client):
e = ContractExecuteClient(mock_client)
result = await e.register_agent("msg1...", "a1", "name", "desc", "c1")
assert result is not None
test_events.py
# 文件: tests/test_events.py
import pytest
from unittest.mock import AsyncMock, MagicMock, patch
from msg_sdk.events import EventClient
from msg_sdk.types import ChainEvent
@pytest.mark.asyncio
async def test_event_client_connect():
with patch("websockets.connect", new_callable=AsyncMock) as mock_ws:
ec = EventClient("ws://test:26657")
await ec.connect()
assert ec._running is True
await ec.disconnect()
class TestChainEvent:
def test_from_dict(self):
data = {"type": "transfer", "attributes": [{"key": "sender", "value": "msg1..."}], "height": "100"}
event = ChainEvent.from_dict(data)
assert event.event_type == "transfer"
assert event.attributes["sender"] == "msg1..."
assert event.block_height == 100
test_governance.py
# 文件: tests/test_governance.py
import pytest
from unittest.mock import AsyncMock, MagicMock
from msg_sdk.governance import GovernanceClient
from msg_sdk.client import MSGChainClient
class TestGovernance:
@pytest.fixture
def mock_client(self):
c = MagicMock(spec=MSGChainClient)
c._request = AsyncMock(return_value={"proposals": []})
c._request_sync = MagicMock(return_value={"proposals": []})
return c
@pytest.mark.asyncio
async def test_get_proposals(self, mock_client):
g = GovernanceClient(mock_client)
props = await g.get_proposals()
assert isinstance(props, list)
pytest.ini
# 文件: pytest.ini
[pytest]
asyncio_mode = auto
testpaths = tests
python_files = test_*.py
markers =
async: async test
slow: slow integration test
16. init.py 包入口
# 文件: msg_sdk/__init__.py
from msg_sdk.client import MSGChainClient, ChainConfig, GasFee
from msg_sdk.wallet import Wallet, DilithiumWallet
from msg_sdk.bank import BankClient
from msg_sdk.staking import StakingClient
from msg_sdk.contract import ContractQueryClient, ContractExecuteClient
from msg_sdk.agent import AgentAPIClient
from msg_sdk.governance import GovernanceClient
from msg_sdk.events import EventClient
from msg_sdk.errors import *
from msg_sdk.types import *
from msg_sdk.utils import *
__version__ = "1.0.0"
__all__ = [
"MSGChainClient", "ChainConfig", "GasFee",
"Wallet", "DilithiumWallet",
"BankClient", "StakingClient",
"ContractQueryClient", "ContractExecuteClient",
"AgentAPIClient", "GovernanceClient", "EventClient",
]
17. SDK 快速参考
链配置
| 参数 | 值 |
|---|---|
| Chain ID | msg-chain-1 |
| Bech32 前缀 | msg |
| CoinType | 118 |
| MSG 精度 | 18 |
| Gas (低/中/高) | 1,000,000,000 attoMSG/gas |
| 出块时间 | 5s |
| RPC | http://localhost:26657 |
| REST | http://localhost:1317 |
| gRPC | localhost:9090 |
REST API 端点
| 方法 | 路径 | 用途 |
|---|---|---|
| GET | /cosmos/bank/v1beta1/balances/{address} | 查询余额 |
| GET | /cosmos/bank/v1beta1/supply | 总供应量 |
| GET | /cosmos/staking/v1beta1/validators | 验证人列表 |
| GET | /cosmwasm/wasm/v1/code/{code_id} | 合约代码 |
| GET | /cosmwasm/wasm/v1/contract/{address} | 合约信息 |
| POST | /cosmwasm/wasm/v1/contract/{address}/smart | 智能查询 |
Agent API 端点
| 分组 | 方法 | 路径 | 用途 |
|---|---|---|---|
| Query | POST | /api/v1/agent/query | 查询 Agent |
| Query | GET | /api/v1/agent/{id}/status | 状态 |
| Query | GET | /api/v1/agent/{id}/history | 历史 |
| Events | POST | /api/v1/agent/events/subscribe | 订阅 |
| Events | POST | /api/v1/agent/events/unsubscribe | 取消订阅 |
| Events | POST | /api/v1/agent/events/replay | 重放 |
| Oracle | POST | /api/v1/agent/oracle/request | 预言机请求 |
| Oracle | GET | /api/v1/agent/oracle/result/{id} | 结果 |
| Wallet | POST | /api/v1/agent/wallet/create | 创建钱包 |
| Wallet | GET | /api/v1/agent/wallet/{id}/balance | 余额 |
| Wallet | POST | /api/v1/agent/wallet/transfer | 转账 |
| Wallet | GET | /api/v1/agent/wallet/transactions | 交易历史 |
| MPC | POST | /api/v1/agent/mpc/sign | MPC 签名 |
| MPC | POST | /api/v1/agent/mpc/submit | 提交部分签名 |
| MPC | POST | /api/v1/agent/mpc/aggregate | 聚合签名 |
| Payment | POST | /api/v1/agent/payment/create-session | 创建支付会话 |
| Payment | POST | /api/v1/agent/payment/close-session | 关闭支付会话 |
错误码速查
| 代码 | 名称 | Python 异常 |
|---|---|---|
| 0 | OK | OK |
| 1 | STUB | StubError |
| 2 | NOT_IMPLEMENTED | NotImplementedErrorCode |
| 3 | UNAUTHORIZED | UnauthorizedError |
| 4 | INSUFFICIENT_FUNDS | InsufficientFundsError |
| 5 | CONTRACT_FAILED | ContractFailedError |
| 6 | QUERY_TIMEOUT | QueryTimeoutError |
| 7 | INVALID_PARAM | InvalidParamError |
| 8 | DAO_TIMELOCK | DAOTimelockError |
| 9 | HIGH_VALUE | HighValueError |
| 10 | AGENT_NOT_FOUND | AgentNotFoundError |
| 11 | AGENT_EXISTS | AgentExistsError |
| 12 | CONSTITUTION_VIOLATION | ConstitutionViolationError |
| 13 | SESSION_EXPIRED | SessionExpiredError |
| 14 | SESSION_LIMIT | SessionLimitError |
| 15 | INVALID_SIG | InvalidSignatureError |
| 16 | DUPLICATE_NONCE | DuplicateNonceError |
| 17 | RATE_LIMIT | RateLimitError |
AI 合约查询/执行方法
| 合约 | 查询方法 | 执行方法 |
|---|---|---|
| agent_registry_v1 | get_agent, list_agents | register_agent, update_agent, deregister_agent, set_status |
| agent_payment_v1 | get_session | create_session, fund_session, release_payment, dispute_payment |
| aidid_did_registry_v1 | resolve_did | create_did, update_did, deactivate_did |
| ai_agent_constitution_v1 | get_constitution, is_action_allowed | add_rule, remove_rule |
| micropayment_session_v1 | get_channel | open_channel, deposit, claim |
18. 常见问题
1. 连接失败怎么办?
确保节点正在运行且 RPC/REST 端口可访问:
curl http://localhost:26657/status
curl http://localhost:1317/cosmos/base/tendermint/v1beta1/node_info
2. 如何获取测试代币?
使用水龙头 (faucet):
import httpx
resp = httpx.post("http://localhost:8000/faucet", json={"address": "msg1..."})
print(resp.json())
3. 交易广播失败
常见原因:
- 余额不足 (InsufficientFundsError)
- sequence 不匹配 (DuplicateNonceError)
- Gas 不足 (增加 gas_limit)
- 签名无效 (InvalidSignatureError)
4. 合约查询无数据
确认合约地址正确,合约已部署且有数据:
# 查询合约信息
curl http://localhost:1317/cosmwasm/wasm/v1/contract/<address>
5. 如何安全存储私钥?
from msg_sdk.wallet import Wallet
from cryptography.fernet import Fernet
# 加密
key = Fernet.generate_key()
fernet = Fernet(key)
w = Wallet.generate()
encrypted = fernet.encrypt(w.get_private_key_hex().encode())
# 解密
decrypted = fernet.decrypt(encrypted).decode()
wallet = Wallet.from_private_key(decrypted)
19. 从 v0.x 迁移
Breaking Changes
- Python >= 3.11 必需
- Client 重构: 移除旧的 RPCClient,统一使用 MSGChainClient
- 错误处理: 所有错误码映射为 Python 异常
- 异步优先: async/await 为主要 API
- 类型系统: 所有响应使用 dataclass
迁移步骤
- 升级依赖:
pip install -U msg-sdk - 替换 RPCClient -> MSGChainClient
- 添加 await 关键字到所有异步调用
- 更新错误处理: try/except 捕获具体异常类
20. 版本历史
v1.0.0 (2024-06-15)
- 首个稳定版本
- 完整覆盖 MSG Chain 所有模块
- 18 种错误码映射
- 5 个 AI 合约支持
- Agent API 18 个端点
- Dilithium-5 后量子签名支持
- WebSocket 事件订阅
v0.9.0 (2024-05-01)
- Beta 版本
- 核心客户端完成
- 基础银行/质押功能
v0.5.0 (2024-03-01)
- Alpha 版本
- 原型客户端
- 基础钱包功能
