dApp Docs/Python SDK开发指南
Development reference. Not independently verified for production.

MSG Chain Python SDK 开发指南

数据来源:MSG Chain 代码库核实

主网状态: No-Go — 当前 MSGChain 主网裁决为 No-Go,以下内容反映代码实际状态,不代表生产可用。

1. 概述

MSG Chain Python SDK 是一个用于与 MSG Chain 区块链交互的完整 Python 工具包。该 SDK 的设计参考了 web3.py 和 cosmpy 的架构模式,提供了面向对象、类型安全的接口,支持同步和异步两种编程模型。

核心特性

依赖

包名 版本 用途
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. 交易广播失败

常见原因:

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

  1. Python >= 3.11 必需
  2. Client 重构: 移除旧的 RPCClient,统一使用 MSGChainClient
  3. 错误处理: 所有错误码映射为 Python 异常
  4. 异步优先: async/await 为主要 API
  5. 类型系统: 所有响应使用 dataclass

迁移步骤

  1. 升级依赖: pip install -U msg-sdk
  2. 替换 RPCClient -> MSGChainClient
  3. 添加 await 关键字到所有异步调用
  4. 更新错误处理: try/except 捕获具体异常类

20. 版本历史

v1.0.0 (2024-06-15)

v0.9.0 (2024-05-01)

v0.5.0 (2024-03-01)