dApp Docs/AI Agent多语言框架集成(LangChain-CrewAI)
Development reference. Not independently verified for production.

AI Agent 多语言框架集成指南:LangChain · CrewAI · AutoGen · LangGraph + MSG Chain

版本: v1.0 | 链: msg-chain-1 | 地址前缀: msg | 签名算法: Dilithium-5 (后量子安全)

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


目录

  1. 概述
  2. LangChain 集成
  3. CrewAI 多 Agent 集成
  4. AutoGen 集成
  5. LangGraph 流程编排
  6. 链下 Agent 与链上合约通信
  7. 完整示例:DeFi 多 Agent 交易系统
  8. 部署指南

1. 概述

1.1 为什么需要将 MSG Chain 与 AI 框架集成?

MSG Chain 是一个高性能、兼容 CosmWasm、采用 Dilithium-5 后量子签名的 Layer 1 区块链。AI Agent 框架(LangChain、CrewAI、AutoGen、LangGraph)为开发者提供了构建智能代理的标准化工具。两者结合,可以实现:

场景 说明
自主 DeFi 代理 AI Agent 自动分析市场数据、执行交易、管理流动性
链上数据分析 用自然语言查询链上状态,Agent 自动翻译为 RPC 调用
智能合约审计 Agent 读取合约代码并执行安全分析
多 Agent 协作 多个专业化 Agent 分工完成复杂链上操作
事件驱动自动化 Agent 监听链上事件并自动触发响应

1.2 架构总览

┌─────────────────────────────────────────────────────────────┐
│                    AI Agent 框架层                           │
│  ┌──────────┐  ┌──────────┐  ┌──────────┐  ┌──────────┐  │
│  │ LangChain │  │  CrewAI  │  │  AutoGen │  │ LangGraph │  │
│  │  Tool     │  │  Agent   │  │ Assistant│  │  State    │  │
│  │  Chain    │  │  Crew    │  │  UserProxy│  │  Graph    │  │
│  └────┬─────┘  └────┬─────┘  └────┬─────┘  └────┬─────┘  │
│       │              │              │              │        │
│  ┌────┴──────────────┴──────────────┴──────────────┴────┐  │
│  │              MSG Chain SDK 抽象层                    │  │
│  │  MsgChainClient · Wallet · TxBuilder · QueryClient │  │
│  └────────────────────────┬────────────────────────────┘  │
│                           │                                │
│  ┌────────────────────────┴────────────────────────────┐  │
│  │              Agent API 网关                           │  │
│  │  REST · WebSocket · gRPC · Event Subscription        │  │
│  └────────────────────────┬────────────────────────────┘  │
├───────────────────────────┼────────────────────────────────┤
│                           │                                │
│  ┌────────────────────────┴────────────────────────────┐  │
│  │              MSG Chain 链层                          │  │
│  │  RPC · ABCI · CosmWasm · IBC · Staking · Gov        │  │
│  │  msg-chain-1 | bech32: msg | Dilithium-5            │  │
│  └─────────────────────────────────────────────────────┘  │
└───────────────────────────────────────────────────────────┘

1.3 框架对比

框架 核心模型 适合场景 MSG Chain 集成方式
LangChain Tool + Chain + Agent 单 Agent 工具调用 自定义 Tool 封装 RPC
CrewAI Agent + Task + Crew 多 Agent 协作 Agent 各司其职,Task 链上操作
AutoGen AssistantAgent + UserProxyAgent 多轮对话任务 Function map 绑定链上方法
LangGraph StateGraph + Node + Edge 复杂状态机流程 Graph node 封装合约交互

1.4 前置条件

# 安装 MSG Chain SDK
pip install msgchain-sdk

# 安装 AI 框架
pip install langchain langchain-community
pip install crewai
pip install pyautogen
pip install langgraph

# 或使用 poetry
poetry add msgchain-sdk langchain crewai pyautogen langgraph

1.5 环境变量

# .env
MSG_CHAIN_RPC=https://rpc.msgchain.org
MSG_CHAIN_REST=https://rest.msgchain.org
MSG_CHAIN_CHAIN_ID=msg-chain-1
MSG_CHAIN_WS=wss://rpc.msgchain.org/websocket

# Agent 钱包(Dilithium-5 私钥)
AGENT_PRIVATE_KEY=your_dilithium5_private_key_hex
AGENT_ADDRESS=msg1agent_address_here

# LLM 配置
OPENAI_API_KEY=sk-...
# 或
ANTHROPIC_API_KEY=sk-ant-...

2. LangChain 集成

2.1 MSG Chain SDK 基础封装

"""
msgchain_sdk_wrapper.py — MSG Chain SDK 封装层
供所有 AI 框架统一调用
"""

import os
import json
from typing import Optional, Any
from dataclasses import dataclass, field
from enum import Enum

import requests
from dotenv import load_dotenv

load_dotenv()


class TxStatus(Enum):
    PENDING = "pending"
    SUCCESS = "success"
    FAILED = "failed"


@dataclass
class MsgChainConfig:
    rpc: str = os.getenv("MSG_CHAIN_RPC", "https://rpc.msgchain.org")
    rest: str = os.getenv("MSG_CHAIN_REST", "https://rest.msgchain.org")
    chain_id: str = os.getenv("MSG_CHAIN_CHAIN_ID", "msg-chain-1")
    ws: str = os.getenv("MSG_CHAIN_WS", "wss://rpc.msgchain.org/websocket")


class MsgChainClient:
    """MSG Chain 统一客户端"""

    def __init__(self, config: Optional[MsgChainConfig] = None):
        self.config = config or MsgChainConfig()
        self.session = requests.Session()
        self.session.headers.update({"Content-Type": "application/json"})

    # ── 账户查询 ──────────────────────────────────────────

    def query_balance(self, address: str, denom: str = "uaimgs") -> dict:
        """查询账户余额"""
        payload = {
            "jsonrpc": "2.0",
            "id": 1,
            "method": "abci_query",
            "params": {
                "path": f"/cosmos.bank.v1beta1/balances/{address}/by_denom",
                "data": denom.encode().hex(),
            },
        }
        resp = self.session.post(self.config.rpc, json=payload)
        resp.raise_for_status()
        return resp.json()

    def query_account(self, address: str) -> dict:
        """查询账户信息(nonce、account number)"""
        payload = {
            "jsonrpc": "2.0",
            "id": 1,
            "method": "abci_query",
            "params": {
                "path": f"/cosmos.auth.v1beta1/accounts/{address}",
            },
        }
        resp = self.session.post(self.config.rpc, json=payload)
        resp.raise_for_status()
        return resp.json()

    # ── 代币与合约查询 ────────────────────────────────────

    def query_contract(self, contract_addr: str, query_msg: dict) -> dict:
        """查询 CosmWasm 合约状态"""
        payload = {
            "jsonrpc": "2.0",
            "id": 1,
            "method": "abci_query",
            "params": {
                "path": f"/cosmwasm.wasm.v1/contract/{contract_addr}/smart",
                "data": json.dumps(query_msg).encode().hex(),
            },
        }
        resp = self.session.post(self.config.rpc, json=payload)
        resp.raise_for_status()
        return resp.json()

    def query_contract_raw(self, contract_addr: str, key: str) -> dict:
        """查询合约原始存储"""
        payload = {
            "jsonrpc": "2.0",
            "id": 1,
            "method": "abci_query",
            "params": {
                "path": (
                    f"/cosmwasm.wasm.v1/contract/{contract_addr}/raw"
                ),
                "data": key.encode().hex(),
            },
        }
        resp = self.session.post(self.config.rpc, json=payload)
        resp.raise_for_status()
        return resp.json()

    def query_all_balances(self, address: str) -> dict:
        """查询所有余额"""
        payload = {
            "jsonrpc": "2.0",
            "id": 1,
            "method": "abci_query",
            "params": {
                "path": f"/cosmos.bank.v1beta1/balances/{address}",
            },
        }
        resp = self.session.post(self.config.rpc, json=payload)
        resp.raise_for_status()
        return resp.json()

    # ── 交易广播 ──────────────────────────────────────────

    def broadcast_tx(self, tx_bytes: bytes, mode: str = "BROADCAST_MODE_SYNC") -> dict:
        """广播交易"""
        payload = {
            "jsonrpc": "2.0",
            "id": 1,
            "method": "broadcast_tx_sync",
            "params": {
                "tx": tx_bytes.hex(),
                "mode": mode,
            },
        }
        resp = self.session.post(self.config.rpc, json=payload)
        resp.raise_for_status()
        return resp.json()

    def broadcast_tx_commit(self, tx_bytes: bytes) -> dict:
        """广播交易并等待 commit"""
        payload = {
            "jsonrpc": "2.0",
            "id": 1,
            "method": "broadcast_tx_commit",
            "params": {"tx": tx_bytes.hex()},
        }
        resp = self.session.post(self.config.rpc, json=payload)
        resp.raise_for_status()
        return resp.json()

    def get_tx(self, tx_hash: str) -> dict:
        """查询交易详情"""
        payload = {
            "jsonrpc": "2.0",
            "id": 1,
            "method": "tx",
            "params": {"hash": tx_hash, "prove": True},
        }
        resp = self.session.post(self.config.rpc, json=payload)
        resp.raise_for_status()
        return resp.json()

    # ── 节点状态 ──────────────────────────────────────────

    def query_validators(self) -> dict:
        """查询验证者集合"""
        payload = {
            "jsonrpc": "2.0",
            "id": 1,
            "method": "validators",
            "params": {"height": 0, "page": 1, "per_page": 100},
        }
        resp = self.session.post(self.config.rpc, json=payload)
        resp.raise_for_status()
        return resp.json()

    def query_block(self, height: Optional[int] = None) -> dict:
        """查询区块"""
        payload = {
            "jsonrpc": "2.0",
            "id": 1,
            "method": "block",
            "params": {"height": height} if height else {},
        }
        resp = self.session.post(self.config.rpc, json=payload)
        resp.raise_for_status()
        return resp.json()

    def query_latest_block_height(self) -> int:
        """获取最新区块高度"""
        result = self.query_block()
        return int(result["result"]["block"]["header"]["height"])

    # ── Staking 查询 ──────────────────────────────────────

    def query_delegations(self, delegator: str) -> dict:
        """查询委托"""
        payload = {
            "jsonrpc": "2.0",
            "id": 1,
            "method": "abci_query",
            "params": {
                "path": f"/cosmos.staking.v1beta1/delegations/{delegator}",
            },
        }
        resp = self.session.post(self.config.rpc, json=payload)
        resp.raise_for_status()
        return resp.json()

    def query_unbonding(self, delegator: str) -> dict:
        """查询解委托"""
        payload = {
            "jsonrpc": "2.0",
            "id": 1,
            "method": "abci_query",
            "params": {
                "path": f"/cosmos.staking.v1beta1/unbonding_delegations/{delegator}",
            },
        }
        resp = self.session.post(self.config.rpc, json=payload)
        resp.raise_for_status()
        return resp.json()

    # ── Governance 查询 ───────────────────────────────────

    def query_proposals(self) -> dict:
        """查询所有提案"""
        payload = {
            "jsonrpc": "2.0",
            "id": 1,
            "method": "abci_query",
            "params": {
                "path": "/cosmos.gov.v1beta1/proposals",
            },
        }
        resp = self.session.post(self.config.rpc, json=payload)
        resp.raise_for_status()
        return resp.json()

    def query_proposal(self, proposal_id: int) -> dict:
        """查询单个提案"""
        payload = {
            "jsonrpc": "2.0",
            "id": 1,
            "method": "abci_query",
            "params": {
                "path": f"/cosmos.gov.v1beta1/proposals/{proposal_id}",
            },
        }
        resp = self.session.post(self.config.rpc, json=payload)
        resp.raise_for_status()
        return resp.json()

    # ── IBC 查询 ──────────────────────────────────────────

    def query_ibc_channels(self) -> dict:
        """查询 IBC 通道"""
        payload = {
            "jsonrpc": "2.0",
            "id": 1,
            "method": "abci_query",
            "params": {
                "path": "/ibc.core.channel.v1/channels",
            },
        }
        resp = self.session.post(self.config.rpc, json=payload)
        resp.raise_for_status()
        return resp.json()

    def query_ibc_denom_trace(self, hash_: str) -> dict:
        """查询 IBC 代币溯源"""
        payload = {
            "jsonrpc": "2.0",
            "id": 1,
            "method": "abci_query",
            "params": {
                "path": f"/ibc.applications.transfer.v1/denom_traces/{hash_}",
            },
        }
        resp = self.session.post(self.config.rpc, json=payload)
        resp.raise_for_status()
        return resp.json()


class MsgChainWallet:
    """Dilithium-5 钱包封装(后量子安全签名)"""

    def __init__(self, private_key_hex: str):
        self.private_key = bytes.fromhex(private_key_hex)
        # 实际 Dilithium-5 实现由底层 SDK 提供
        # self.public_key = dilithium5_get_public_key(self.private_key)
        # self.address = derive_msg_address(self.public_key)

    def sign(self, tx_bytes: bytes) -> bytes:
        """使用 Dilithium-5 签名交易"""
        # from msgchain_sdk.dilithium import sign
        # return sign(self.private_key, tx_bytes)
        raise NotImplementedError("Dilithium-5 signing requires msgchain-sdk")

    def get_address(self) -> str:
        """获取 bech32 地址"""
        # return derive_bech32("msg", self.public_key)
        raise NotImplementedError("Address derivation requires msgchain-sdk")


# ── Cosmos SDK REST API 封装 ─────────────────────────────

class MsgChainRestClient:
    """MSG Chain REST API 客户端"""

    def __init__(self, base_url: str = None):
        self.base_url = base_url or os.getenv(
            "MSG_CHAIN_REST", "https://rest.msgchain.org"
        )
        self.session = requests.Session()

    def get(self, path: str, params: dict = None) -> dict:
        resp = self.session.get(f"{self.base_url}{path}", params=params)
        resp.raise_for_status()
        return resp.json()

    def post(self, path: str, data: dict) -> dict:
        resp = self.session.post(f"{self.base_url}{path}", json=data)
        resp.raise_for_status()
        return resp.json()

    # ── Bank Module ───────────────────────────────────────

    def bank_balance(self, address: str, denom: str = "uaimgs") -> dict:
        return self.get(f"/cosmos/bank/v1beta1/balances/{address}/by_denom", {
            "denom": denom,
        })

    def bank_all_balances(self, address: str) -> dict:
        return self.get(f"/cosmos/bank/v1beta1/balances/{address}")

    # ── Staking Module ────────────────────────────────────

    def staking_delegations(self, delegator: str) -> dict:
        return self.get(f"/cosmos/staking/v1beta1/delegations/{delegator}")

    def staking_validators(self, status: str = "BOND_STATUS_BONDED") -> dict:
        return self.get("/cosmos/staking/v1beta1/validators", {
            "status": status,
        })

    # ── Distribution Module ───────────────────────────────

    def distribution_rewards(self, address: str) -> dict:
        return self.get(
            f"/cosmos/distribution/v1beta1/delegators/{address}/rewards"
        )

    # ── Gov Module ────────────────────────────────────────

    def gov_proposals(self) -> dict:
        return self.get("/cosmos/gov/v1beta1/proposals")

    def gov_proposal(self, proposal_id: int) -> dict:
        return self.get(f"/cosmos/gov/v1beta1/proposals/{proposal_id}")

    # ── CosmWasm Module ───────────────────────────────────

    def wasm_contract_info(self, contract: str) -> dict:
        return self.get(f"/cosmwasm/wasm/v1/contract/{contract}")

    def wasm_contract_query(self, contract: str, query: dict) -> dict:
        return self.get(
            f"/cosmwasm/wasm/v1/contract/{contract}/smart",
            {"query_msg": json.dumps(query)},
        )

    def wasm_code_info(self, code_id: int) -> dict:
        return self.get(f"/cosmwasm/wasm/v1/code/{code_id}")

    # ── IBC Module ────────────────────────────────────────

    def ibc_channels(self) -> dict:
        return self.get("/ibc/core/channel/v1/channels")

    def ibc_connections(self) -> dict:
        return self.get("/ibc/core/connection/v1/connections")

    # ── Mint / Inflation ────────────────────────────────

    def inflation(self) -> dict:
        return self.get("/cosmos/mint/v1beta1/inflation")

    def annual_provisions(self) -> dict:
        return self.get("/cosmos/mint/v1beta1/annual_provisions")

2.2 LangChain Tools 定义

"""
msgchain_langchain_tools.py — MSG Chain LangChain 工具集
"""

from typing import Optional, Type, Any
from pydantic import BaseModel, Field

from langchain.tools import BaseTool
from langchain.callbacks.manager import (
    CallbackManagerForToolRun,
    AsyncCallbackManagerForToolRun,
)

from msgchain_sdk_wrapper import MsgChainClient


# ── 输入 Schema ───────────────────────────────────────────

class BalanceQueryInput(BaseModel):
    address: str = Field(description="MSG Chain bech32 address (msg...)")
    denom: str = Field(
        default="uaimgs",
        description="Denom to query (default: uaimgs)",
    )


class ContractQueryInput(BaseModel):
    contract_addr: str = Field(description="CosmWasm contract address")
    query_msg: str = Field(
        description="JSON query message as string (e.g. '{\\"balance\\": {\\"address\\": \\"msg1...\\"}}')"
    )


class TokenInfoInput(BaseModel):
    contract_addr: str = Field(description="CW20 token contract address")


class StakingInfoInput(BaseModel):
    delegator: str = Field(description="Delegator address")
    validator: Optional[str] = Field(default=None, description="Optional validator address")


class ValidatorListInput(BaseModel):
    status: str = Field(default="BOND_STATUS_BONDED", description="Validator status")


class ProposalQueryInput(BaseModel):
    proposal_id: Optional[int] = Field(default=None, description="Proposal ID")


class BlockQueryInput(BaseModel):
    height: Optional[int] = Field(default=None, description="Block height")


# ── 工具: 查询余额 ───────────────────────────────────────

class MsgChainBalanceTool(BaseTool):
    name: str = "msgchain_query_balance"
    description: str = (
        "Query MSG Chain account balance for a given address and denom. "
        "Returns the current balance in the smallest unit (uaimgs)."
    )
    args_schema: Type[BaseModel] = BalanceQueryInput

    def _run(
        self,
        address: str,
        denom: str = "uaimgs",
        run_manager: Optional[CallbackManagerForToolRun] = None,
    ) -> str:
        client = MsgChainClient()
        try:
            result = client.query_balance(address, denom)
            return json.dumps(result, indent=2, ensure_ascii=False)
        except Exception as e:
            return f"Error querying balance: {str(e)}"


# ── 工具: 查询所有余额 ──────────────────────────────────

class MsgChainAllBalancesTool(BaseTool):
    name: str = "msgchain_query_all_balances"
    description: str = (
        "Query all balances for an MSG Chain address. "
        "Returns all denominations with their balances."
    )
    args_schema: Type[BaseModel] = BalanceQueryInput

    def _run(
        self,
        address: str,
        denom: str = "",
        run_manager: Optional[CallbackManagerForToolRun] = None,
    ) -> str:
        client = MsgChainClient()
        try:
            result = client.query_all_balances(address)
            return json.dumps(result, indent=2, ensure_ascii=False)
        except Exception as e:
            return f"Error querying all balances: {str(e)}"


# ── 工具: 查询合约 ──────────────────────────────────────

class MsgChainContractQueryTool(BaseTool):
    name: str = "msgchain_contract_query"
    description: str = (
        "Query a CosmWasm smart contract on MSG Chain. "
        "Provide the contract address and a JSON query message. "
        "Example query: '{\\\"balance\\\": {\\\"address\\\": \\\"msg1...\\\"}}'"
    )
    args_schema: Type[BaseModel] = ContractQueryInput

    def _run(
        self,
        contract_addr: str,
        query_msg: str,
        run_manager: Optional[CallbackManagerForToolRun] = None,
    ) -> str:
        client = MsgChainClient()
        try:
            msg = json.loads(query_msg)
            result = client.query_contract(contract_addr, msg)
            return json.dumps(result, indent=2, ensure_ascii=False)
        except json.JSONDecodeError as e:
            return f"Invalid JSON query_msg: {str(e)}"
        except Exception as e:
            return f"Error querying contract: {str(e)}"


# ── 工具: 查询 CW20 代币信息 ───────────────────────────

class MsgChainTokenInfoTool(BaseTool):
    name: str = "msgchain_token_info"
    description: str = (
        "Query CW20 token information including name, symbol, decimals, "
        "total supply on MSG Chain."
    )
    args_schema: Type[BaseModel] = TokenInfoInput

    def _run(
        self,
        contract_addr: str,
        run_manager: Optional[CallbackManagerForToolRun] = None,
    ) -> str:
        client = MsgChainClient()
        try:
            info = client.query_contract(contract_addr, {"token_info": {}})
            supply = client.query_contract(contract_addr, {"total_supply": {}})
            return json.dumps({
                "token_info": info,
                "total_supply": supply,
            }, indent=2, ensure_ascii=False)
        except Exception as e:
            return f"Error querying token info: {str(e)}"


# ── 工具: 查询区块 ──────────────────────────────────────

class MsgChainBlockQueryTool(BaseTool):
    name: str = "msgchain_block_query"
    description: str = (
        "Query block information from MSG Chain. "
        "Leave height empty for latest block."
    )
    args_schema: Type[BaseModel] = BlockQueryInput

    def _run(
        self,
        height: Optional[int] = None,
        run_manager: Optional[CallbackManagerForToolRun] = None,
    ) -> str:
        client = MsgChainClient()
        try:
            result = client.query_block(height)
            return json.dumps(result, indent=2, ensure_ascii=False)
        except Exception as e:
            return f"Error querying block: {str(e)}"


# ── 工具: 查询验证者 ──────────────────────────────────

class MsgChainValidatorTool(BaseTool):
    name: str = "msgchain_validators"
    description: str = "Query MSG Chain validator set."
    args_schema: Type[BaseModel] = ValidatorListInput

    def _run(
        self,
        status: str = "BOND_STATUS_BONDED",
        run_manager: Optional[CallbackManagerForToolRun] = None,
    ) -> str:
        client = MsgChainClient()
        try:
            result = client.query_validators()
            return json.dumps(result, indent=2, ensure_ascii=False)
        except Exception as e:
            return f"Error querying validators: {str(e)}"


# ── 工具: 查询委托 ────────────────────────────────────

class MsgChainDelegationTool(BaseTool):
    name: str = "msgchain_delegations"
    description: str = (
        "Query staking delegations for an MSG Chain address. "
        "Returns all active delegations and unbonding information."
    )
    args_schema: Type[BaseModel] = StakingInfoInput

    def _run(
        self,
        delegator: str,
        validator: Optional[str] = None,
        run_manager: Optional[CallbackManagerForToolRun] = None,
    ) -> str:
        client = MsgChainClient()
        try:
            delegations = client.query_delegations(delegator)
            unbonding = client.query_unbonding(delegator)
            return json.dumps({
                "delegations": delegations,
                "unbonding": unbonding,
            }, indent=2, ensure_ascii=False)
        except Exception as e:
            return f"Error querying delegations: {str(e)}"


# ── 工具: 查询治理提案 ─────────────────────────────────

class MsgChainProposalTool(BaseTool):
    name: str = "msgchain_proposals"
    description: str = (
        "Query governance proposals on MSG Chain. "
        "Leave proposal_id empty to list all proposals."
    )
    args_schema: Type[BaseModel] = ProposalQueryInput

    def _run(
        self,
        proposal_id: Optional[int] = None,
        run_manager: Optional[CallbackManagerForToolRun] = None,
    ) -> str:
        client = MsgChainClient()
        try:
            if proposal_id:
                result = client.query_proposal(proposal_id)
            else:
                result = client.query_proposals()
            return json.dumps(result, indent=2, ensure_ascii=False)
        except Exception as e:
            return f"Error querying proposals: {str(e)}"


# ── 工具: 获取最新区块高度 ────────────────────────────

class MsgChainLatestHeightTool(BaseTool):
    name: str = "msgchain_latest_height"
    description: str = "Get the latest block height on MSG Chain."
    args_schema: Type[BaseModel] = BlockQueryInput

    def _run(
        self,
        height: Optional[int] = None,
        run_manager: Optional[CallbackManagerForToolRun] = None,
    ) -> str:
        client = MsgChainClient()
        try:
            height = client.query_latest_block_height()
            return f"Latest block height: {height}"
        except Exception as e:
            return f"Error: {str(e)}"


# ── 工具: 查询交易 ──────────────────────────────────────

class MsgChainTxQueryTool(BaseTool):
    name: str = "msgchain_tx_query"
    description: str = "Query a transaction by hash on MSG Chain."

    class TxQueryInput(BaseModel):
        tx_hash: str = Field(description="Transaction hash (hex)")

    args_schema: Type[BaseModel] = TxQueryInput

    def _run(
        self,
        tx_hash: str,
        run_manager: Optional[CallbackManagerForToolRun] = None,
    ) -> str:
        client = MsgChainClient()
        try:
            result = client.get_tx(tx_hash)
            return json.dumps(result, indent=2, ensure_ascii=False)
        except Exception as e:
            return f"Error querying tx: {str(e)}"


# ── 工具: 查询通胀参数 ────────────────────────────────

class MsgChainInflationTool(BaseTool):
    name: str = "msgchain_inflation"
    description: str = "Query MSG Chain inflation rate and annual provisions."
    args_schema: Type[BaseModel] = BlockQueryInput

    def _run(
        self,
        height: Optional[int] = None,
        run_manager: Optional[CallbackManagerForToolRun] = None,
    ) -> str:
        rest = MsgChainRestClient()
        try:
            inflation = rest.inflation()
            provisions = rest.annual_provisions()
            return json.dumps({
                "inflation": inflation,
                "annual_provisions": provisions,
            }, indent=2, ensure_ascii=False)
        except Exception as e:
            return f"Error querying inflation: {str(e)}"

2.3 转账与合约执行工具(含签名)

"""
msgchain_langchain_tx_tools.py — MSG Chain 交易执行工具
需要 Dilithium-5 钱包支持
"""

import json
import os
from typing import Optional, Type

from pydantic import BaseModel, Field
from langchain.tools import BaseTool
from langchain.callbacks.manager import CallbackManagerForToolRun

from msgchain_sdk_wrapper import MsgChainClient, MsgChainConfig


# ── 转账输入 ─────────────────────────────────────────────

class TransferInput(BaseModel):
    recipient: str = Field(description="Recipient bech32 address")
    amount: str = Field(description="Amount in smallest unit (uaimgs)")
    memo: str = Field(default="", description="Optional transaction memo")
    denom: str = Field(default="uaimgs", description="Denom to transfer")


class ContractExecuteInput(BaseModel):
    contract_addr: str = Field(description="Contract address")
    execute_msg: str = Field(
        description="Execute message as JSON string. "
        "Example: '{\\"transfer\\": {\\"recipient\\": \\"msg1...\\", "
        "\\"amount\\": \\"1000000\\"}}'"
    )
    funds: str = Field(
        default="",
        description="Funds to send with execute, e.g. '1000000uaimgs'",
    )


class DelegateInput(BaseModel):
    validator: str = Field(description="Validator address (msgvaloper...)")
    amount: str = Field(description="Amount to delegate in uaimgs")


class VoteInput(BaseModel):
    proposal_id: int = Field(description="Proposal ID")
    option: str = Field(
        description="Vote option: yes/no/abstain/no_with_veto"
    )


# ── 转账工具 ─────────────────────────────────────────────

class MsgChainTransferTool(BaseTool):
    name: str = "msgchain_transfer"
    description: str = (
        "Transfer MSG tokens (uaimgs) to another address. "
        "Requires the agent wallet to be configured with private key. "
        "Returns transaction hash on success."
    )
    args_schema: Type[BaseModel] = TransferInput

    def _run(
        self,
        recipient: str,
        amount: str,
        memo: str = "",
        denom: str = "uaimgs",
        run_manager: Optional[CallbackManagerForToolRun] = None,
    ) -> str:
        try:
            return self._build_and_broadcast(recipient, amount, memo, denom)
        except Exception as e:
            return f"Transfer failed: {str(e)}"

    def _build_and_broadcast(
        self, recipient: str, amount: str, memo: str, denom: str
    ) -> str:
        """构建并广播转账交易"""
        private_key = os.getenv("AGENT_PRIVATE_KEY")
        if not private_key:
            return "Error: AGENT_PRIVATE_KEY not set"

        client = MsgChainClient()
        config = MsgChainConfig()

        # 1. 查询账户信息
        sender = os.getenv("AGENT_ADDRESS")
        account_info = client.query_account(sender)

        # 2. 构建 StdTx
        tx = {
            "body": {
                "messages": [
                    {
                        "@type": "/cosmos.bank.v1beta1.MsgSend",
                        "from_address": sender,
                        "to_address": recipient,
                        "amount": [{"denom": denom, "amount": amount}],
                    }
                ],
                "memo": memo,
                "timeout_height": "0",
                "extension_options": [],
                "non_critical_extension_options": [],
            },
            "auth_info": {
                "signer_infos": [],
                "fee": {
                    "amount": [{"denom": "uaimgs", "amount": "5000"}],
                    "gas_limit": "200000",
                    "payer": "",
                    "granter": "",
                },
            },
            "signatures": [],
        }

        # 3. 签名(Dilithium-5)
        # tx_bytes = proto_encode(tx)
        # signature = sign_dilithium5(private_key, tx_bytes)
        # tx["signatures"] = [signature]
        # signed_bytes = proto_encode(tx)

        # 4. 广播
        # result = client.broadcast_tx(signed_bytes)
        # tx_hash = result["result"]["hash"]

        return json.dumps(
            {
                "status": "simulated",
                "sender": sender,
                "recipient": recipient,
                "amount": f"{amount}{denom}",
                "tx_hash": "TX_PENDING_AFTER_SIGNING",
                "chain_id": config.chain_id,
            },
            indent=2,
            ensure_ascii=False,
        )


# ── 合约执行工具 ─────────────────────────────────────────

class MsgChainContractExecuteTool(BaseTool):
    name: str = "msgchain_contract_execute"
    description: str = (
        "Execute a CosmWasm smart contract on MSG Chain. "
        "Can include funds to send with the execution. "
        "Common actions: transfer CW20 tokens, swap on DEX, "
        "deposit to lending protocol, etc."
    )
    args_schema: Type[BaseModel] = ContractExecuteInput

    def _run(
        self,
        contract_addr: str,
        execute_msg: str,
        funds: str = "",
        run_manager: Optional[CallbackManagerForToolRun] = None,
    ) -> str:
        try:
            msg = json.loads(execute_msg)
            return self._build_and_broadcast_contract(
                contract_addr, msg, funds
            )
        except json.JSONDecodeError as e:
            return f"Invalid execute_msg JSON: {str(e)}"
        except Exception as e:
            return f"Contract execute failed: {str(e)}"

    def _build_and_broadcast_contract(
        self, contract_addr: str, msg: dict, funds: str
    ) -> str:
        private_key = os.getenv("AGENT_PRIVATE_KEY")
        if not private_key:
            return "Error: AGENT_PRIVATE_KEY not set"

        sender = os.getenv("AGENT_ADDRESS")
        config = MsgChainConfig()

        funds_list = []
        if funds:
            parts = funds.split(",")
            for part in parts:
                part = part.strip()
                if part:
                    import re
                    m = re.match(r"(\\d+)(\\D+)", part)
                    if m:
                        funds_list.append({
                            "denom": m.group(2),
                            "amount": m.group(1),
                        })

        execute_msg = {
            "@type": "/cosmwasm.wasm.v1.MsgExecuteContract",
            "sender": sender,
            "contract": contract_addr,
            "msg": msg,
            "funds": funds_list,
        }

        tx = {
            "body": {
                "messages": [execute_msg],
                "memo": "",
                "timeout_height": "0",
                "extension_options": [],
                "non_critical_extension_options": [],
            },
            "auth_info": {
                "signer_infos": [],
                "fee": {
                    "amount": [{"denom": "uaimgs", "amount": "5000"}],
                    "gas_limit": "500000",
                    "payer": "",
                    "granter": "",
                },
            },
            "signatures": [],
        }

        return json.dumps(
            {
                "status": "simulated",
                "sender": sender,
                "contract": contract_addr,
                "msg": msg,
                "funds": funds_list,
                "chain_id": config.chain_id,
                "note": "Replace with actual signing in production",
            },
            indent=2,
            ensure_ascii=False,
        )


# ── 委托工具 ─────────────────────────────────────────────

class MsgChainDelegateTool(BaseTool):
    name: str = "msgchain_delegate"
    description: str = (
        "Delegate MSG tokens to a validator. "
        "Requires agent wallet configuration."
    )
    args_schema: Type[BaseModel] = DelegateInput

    def _run(
        self,
        validator: str,
        amount: str,
        run_manager: Optional[CallbackManagerForToolRun] = None,
    ) -> str:
        sender = os.getenv("AGENT_ADDRESS")
        delegate_msg = {
            "@type": "/cosmos.staking.v1beta1.MsgDelegate",
            "delegator_address": sender,
            "validator_address": validator,
            "amount": {"denom": "uaimgs", "amount": amount},
        }

        return json.dumps({
            "status": "simulated",
            "action": "delegate",
            "sender": sender,
            "validator": validator,
            "amount": amount,
            "note": "Requires Dilithium-5 signing in production",
        }, indent=2, ensure_ascii=False)


# ── 投票工具 ─────────────────────────────────────────────

class MsgChainVoteTool(BaseTool):
    name: str = "msgchain_vote"
    description: str = (
        "Vote on an MSG Chain governance proposal. "
        "Options: yes, no, abstain, no_with_veto"
    )
    args_schema: Type[BaseModel] = VoteInput

    def _run(
        self,
        proposal_id: int,
        option: str,
        run_manager: Optional[CallbackManagerForToolRun] = None,
    ) -> str:
        valid_options = ["yes", "no", "abstain", "no_with_veto"]
        if option not in valid_options:
            return f"Invalid option: {option}. Must be one of {valid_options}"

        sender = os.getenv("AGENT_ADDRESS")
        return json.dumps({
            "status": "simulated",
            "action": "vote",
            "voter": sender,
            "proposal_id": proposal_id,
            "option": option,
        }, indent=2, ensure_ascii=False)

2.4 LangChain Agent 集成

"""
msgchain_langchain_agent.py — 将 MSG Chain 工具注入 LangChain Agent
"""

from langchain.agents import (
    AgentExecutor,
    create_openai_functions_agent,
    create_react_agent,
)
from langchain.prompts import ChatPromptTemplate, PromptTemplate
from langchain.memory import ConversationBufferMemory
from langchain_openai import ChatOpenAI

from msgchain_langchain_tools import (
    MsgChainBalanceTool,
    MsgChainAllBalancesTool,
    MsgChainContractQueryTool,
    MsgChainTokenInfoTool,
    MsgChainBlockQueryTool,
    MsgChainValidatorTool,
    MsgChainDelegationTool,
    MsgChainProposalTool,
    MsgChainInflationTool,
    MsgChainLatestHeightTool,
    MsgChainTxQueryTool,
)
from msgchain_langchain_tx_tools import (
    MsgChainTransferTool,
    MsgChainContractExecuteTool,
    MsgChainDelegateTool,
    MsgChainVoteTool,
)


def create_msgchain_agent(
    model: str = "gpt-4o",
    temperature: float = 0.1,
    verbose: bool = True,
    include_tx_tools: bool = False,
) -> AgentExecutor:
    """创建 MSG Chain LangChain Agent"""

    tools = [
        MsgChainBalanceTool(),
        MsgChainAllBalancesTool(),
        MsgChainContractQueryTool(),
        MsgChainTokenInfoTool(),
        MsgChainBlockQueryTool(),
        MsgChainValidatorTool(),
        MsgChainDelegationTool(),
        MsgChainProposalTool(),
        MsgChainInflationTool(),
        MsgChainLatestHeightTool(),
        MsgChainTxQueryTool(),
    ]

    if include_tx_tools:
        tools.extend([
            MsgChainTransferTool(),
            MsgChainContractExecuteTool(),
            MsgChainDelegateTool(),
            MsgChainVoteTool(),
        ])

    llm = ChatOpenAI(model=model, temperature=temperature)

    system_prompt = (
        "You are an AI assistant specialized in MSG Chain blockchain operations. "
        "You can query on-chain data, analyze blockchain state, and execute transactions. "
        "\\n\\n"
        "MSG Chain is a CosmWasm-enabled Layer 1 blockchain with Dilithium-5 post-quantum signatures. "
        "The native token is uaimgs (1 MSG = 1,000,000 uaimgs). "
        "Addresses use the 'msg' bech32 prefix. "
        "\\n\\n"
        "Available operation[未公开路径]"
        "- Query balances, token info, contract state\\n"
        "- Query blocks, validators, delegations\\n"
        "- Query governance proposals\\n"
        "- Transfer tokens, execute contracts, delegate, vote\\n"
        "\\n"
        "Always use the available tools to answer questions. "
        "For transactions, clearly show the constructed message before execution. "
        "Report results in a clear, structured format."
    )

    prompt = ChatPromptTemplate.from_messages([
        ("system", system_prompt),
        ("human", "{input}"),
        ("placeholder", "{agent_scratchpad}"),
    ])

    agent = create_openai_functions_agent(
        llm=llm,
        tools=tools,
        prompt=prompt,
    )

    executor = AgentExecutor(
        agent=agent,
        tools=tools,
        verbose=verbose,
        handle_parsing_errors=True,
        max_iterations=15,
        early_stopping_method="generate",
    )

    return executor


def create_msgchain_agent_with_memory(
    model: str = "gpt-4o",
    verbose: bool = True,
) -> AgentExecutor:
    """创建带记忆的 MSG Chain Agent"""

    tools = [
        MsgChainBalanceTool(),
        MsgChainAllBalancesTool(),
        MsgChainContractQueryTool(),
        MsgChainTokenInfoTool(),
        MsgChainBlockQueryTool(),
        MsgChainValidatorTool(),
        MsgChainDelegationTool(),
        MsgChainProposalTool(),
        MsgChainInflationTool(),
        MsgChainLatestHeightTool(),
        MsgChainTxQueryTool(),
        MsgChainTransferTool(),
        MsgChainContractExecuteTool(),
    ]

    llm = ChatOpenAI(model=model, temperature=0.1)

    memory = ConversationBufferMemory(
        memory_key="chat_history",
        return_messages=True,
    )

    system_prompt = (
        "You are an MSG Chain blockchain agent with memory. "
        "You remember previous queries and transactions. "
        "Use the tools to answer questions about MSG Chain state."
    )

    prompt = ChatPromptTemplate.from_messages([
        ("system", system_prompt),
        ("placeholder", "{chat_history}"),
        ("human", "{input}"),
        ("placeholder", "{agent_scratchpad}"),
    ])

    agent = create_openai_functions_agent(
        llm=llm,
        tools=tools,
        prompt=prompt,
    )

    executor = AgentExecutor(
        agent=agent,
        tools=tools,
        memory=memory,
        verbose=verbose,
        handle_parsing_errors=True,
        max_iterations=15,
    )

    return executor


# ── ReAct Agent ──────────────────────────────────────────

def create_msgchain_react_agent(
    model: str = "gpt-4o",
    verbose: bool = True,
) -> AgentExecutor:
    """创建基于 ReAct 的 MSG Chain Agent"""

    tools = [
        MsgChainBalanceTool(),
        MsgChainContractQueryTool(),
        MsgChainTokenInfoTool(),
        MsgChainBlockQueryTool(),
        MsgChainLatestHeightTool(),
        MsgChainTransferTool(),
    ]

    llm = ChatOpenAI(model=model, temperature=0.1)

    template = (
        "You are an MSG Chain blockchain agent. Answer the user's questions "
        "using the available tools.\\n\\n"
        "Available tool[未公开路径]}\\n\\n"
        "Use the following forma[未公开路径]"
        "Question: the input question\\n"
        "Thought: what you need to do\\n"
        "Action: the tool name\\n"
        "Action Input: the tool input\\n"
        "Observation: the tool result\\n"
        "... (repeat as needed)\\n"
        "Thought: I now have the answer\\n"
        "Final Answer: the final answer\\n\\n"
        "Previous conversatio[未公开路径]}\\n\\n"
        "Question: {input}\\n"
        "{agent_scratchpad}"
    )

    prompt = PromptTemplate.from_template(template)

    agent = create_react_agent(
        llm=llm,
        tools=tools,
        prompt=prompt,
    )

    executor = AgentExecutor(
        agent=agent,
        tools=tools,
        memory=ConversationBufferMemory(
            memory_key="chat_history",
            return_messages=True,
        ),
        verbose=verbose,
        handle_parsing_errors=True,
        max_iterations=10,
    )

    return executor


# ── 运行示例 ─────────────────────────────────────────────

if __name__ == "__main__":
    agent = create_msgchain_agent(
        model="gpt-4o",
        include_tx_tools=True,
        verbose=True,
    )

    queries = [
        "What is the latest block height on MSG Chain?",
        "Query the balance of address msg1agent...",
        "What are the current validators?",
        "Show me all governance proposals",
        "What is the current inflation rate?",
    ]

    for query in queries:
        print(f"\\n{'='*60}")
        print(f"Query: {query}")
        result = agent.invoke({"input": query})
        print(f"Result: {result['output']}")

2.5 LangChain 顺序链(Sequential Chain)

"""
msgchain_langchain_chain.py — LangChain 顺序链与路由
"""

from langchain.chains import LLMChain, SimpleSequentialChain
from langchain.prompts import PromptTemplate
from langchain_openai import ChatOpenAI

from msgchain_sdk_wrapper import MsgChainClient


class MsgChainAnalysisChain:
    """MSG Chain 数据分析流水线"""

    def __init__(self, model: str = "gpt-4o"):
        self.llm = ChatOpenAI(model=model, temperature=0.1)
        self.client = MsgChainClient()

    def analyze_validator_set(self) -> str:
        """分析验证者集合"""
        validators = self.client.query_validators()

        prompt = PromptTemplate.from_template(
            "Analyze the following MSG Chain validator se[未公开路径]"
            "{validators}\\n\\n"
            "Provide: total validators, voting power distribution, "
            "top 5 validators by voting power, and any observations."
        )

        chain = LLMChain(llm=self.llm, prompt=prompt)
        return chain.run(validators=validators)

    def analyze_governance(self) -> str:
        """分析治理提案"""
        proposals = self.client.query_proposals()

        prompt = PromptTemplate.from_template(
            "Analyze the following MSG Chain governance proposal[未公开路径]"
            "{proposals}\\n\\n"
            "Summarize each active proposal, its status, "
            "voting deadline, and recommend how to vote."
        )

        chain = LLMChain(llm=self.llm, prompt=prompt)
        return chain.run(proposals=proposals)

    def full_chain_report(self) -> str:
        """生成完整链状态报告"""
        chain = SimpleSequentialChain(
            chains=[
                LLMChain(
                    llm=self.llm,
                    prompt=PromptTemplate.from_template(
                        "Get current MSG Chain state summary including "
                        "latest block, validator count, and recent activity."
                    ),
                ),
                LLMChain(
                    llm=self.llm,
                    prompt=PromptTemplate.from_template(
                        "Based on this MSG Chain state: {input}\\n"
                        "Generate a comprehensive status report with "
                        "recommendations for validators and delegators."
                    ),
                ),
            ],
            verbose=True,
        )

        return chain.run("start")
"""

# ── 路由链(Router Chain) ─────────────────────────────

from langchain.chains.router import MultiPromptChain
from langchain.chains.router.llm_router import LLMRouterChain
from langchain.chains.router.multi_prompt_prompt import MULTI_PROMPT_ROUTER_TEMPLATE


def create_msgchain_router() -> MultiPromptChain:
    """创建路由链,根据用户意图路由到不同链"""

    llm = ChatOpenAI(model="gpt-4o", temperature=0)

    prompt_infos = [
        {
            "name": "balance",
            "description": "Good for answering questions about token balances",
            "prompt_template": (
                "You are a balance query specialist for MSG Chain. "
                "Given the user's balance query, use the balance tool "
                "and provide detailed token holdings information.\\n"
                "Query: {input}"
            ),
        },
        {
            "name": "validators",
            "description": "Good for answering questions about validators and staking",
            "prompt_template": (
                "You are a staking specialist for MSG Chain. "
                "Answer questions about validators, delegations, "
                "and staking rewards.\\n"
                "Query: {input}"
            ),
        },
        {
            "name": "governance",
            "description": "Good for answering about proposals and voting",
            "prompt_template": (
                "You are a governance specialist for MSG Chain. "
                "Answer questions about proposals, voting, "
                "and governance parameters.\\n"
                "Query: {input}"
            ),
        },
        {
            "name": "general",
            "description": "Good for general MSG Chain questions",
            "prompt_template": (
                "You are a general MSG Chain assistant. "
                "Answer general questions about the blockchain.\\n"
                "Query: {input}"
            ),
        },
    ]

    destination_chains = {}
    for info in prompt_infos:
        name = info["name"]
        prompt = PromptTemplate.from_template(info["prompt_template"])
        chain = LLMChain(llm=llm, prompt=prompt)
        destination_chains[name] = chain

    default_chain = LLMChain(
        llm=llm,
        prompt=PromptTemplate.from_template(
            "Answer the following MSG Chain question: {input}"
        ),
    )

    router = MultiPromptChain(
        router_chain=LLMRouterChain.from_instance(
            llm=llm,
            prompt=PromptTemplate.from_template(
                MULTI_PROMPT_ROUTER_TEMPLATE
            ),
        ),
        destination_chains=destination_chains,
        default_chain=default_chain,
        verbose=True,
    )

    return router


# ── 带工具的 LLM Chain ─────────────────────────────────

class MsgChainToolChain:
    """将 MSG Chain 工具与 LLM Chain 结合使用"""

    def __init__(self, model: str = "gpt-4o"):
        self.llm = ChatOpenAI(model=model, temperature=0)
        self.client = MsgChainClient()
        self.balance_tool = MsgChainBalanceTool()

    def explain_balance(self, address: str) -> str:
        """查询余额并用自然语言解释"""
        balance = self.balance_tool.run(address)

        prompt = PromptTemplate.from_template(
            "The following is the MSG Chain balance data for address {address}:\\n\\n"
            "{balance}\\n\\n"
            "Explain this balance in simple terms. Include the equivalent in MSG "
            "(1 MSG = 1,000,000 uaimgs). Suggest what the user could do with this balance "
            "(stake, transfer, provide liquidity, etc.)."
        )

        chain = LLMChain(llm=self.llm, prompt=prompt)
        return chain.run(address=address, balance=balance)

    def analyze_token(self, contract_addr: str) -> str:
        """分析 CW20 代币"""
        from msgchain_langchain_tools import MsgChainTokenInfoTool
        token_info = MsgChainTokenInfoTool().run(contract_addr)

        prompt = PromptTemplate.from_template(
            "Analyze this CW20 token on MSG Chai[未公开路径]}\\n\\n"
            "Provide: token name, symbol, total supply, "
            "and thoughts on its utility and distribution."
        )

        chain = LLMChain(llm=self.llm, prompt=prompt)
        return chain.run(info=token_info)


if __name__ == "__main__":
    analysis = MsgChainAnalysisChain()
    print(analysis.analyze_validator_set())

2.6 LangChain 回调与日志

"""
msgchain_langchain_callbacks.py — MSG Chain 操作的 LangChain 回调
"""

from langchain.callbacks.base import BaseCallbackHandler
from langchain.schema import LLMResult
from typing import Any, Dict, List


class MsgChainLoggingCallback(BaseCallbackHandler):
    """记录所有 MSG Chain 相关操作的回调处理器"""

    def __init__(self):
        self.logs: List[Dict[str, Any]] = []

    def on_tool_start(
        self,
        serialized: Dict[str, Any],
        input_str: str,
        **kwargs: Any,
    ) -> None:
        if "msgchain" in serialized.get("name", "").lower():
            self.logs.append({
                "event": "tool_start",
                "tool": serialized["name"],
                "input": input_str,
            })
            print(f"[MSG Chain] Tool Start: {serialized['name']}")
            print(f"  Input: {input_str[:200]}")

    def on_tool_end(
        self,
        output: str,
        **kwargs: Any,
    ) -> None:
        self.logs.append({
            "event": "tool_end",
            "output": output[:500],
        })
        print(f"[MSG Chain] Tool End")

    def on_tool_error(
        self,
        error: Exception,
        **kwargs: Any,
    ) -> None:
        self.logs.append({
            "event": "tool_error",
            "error": str(error),
        })
        print(f"[MSG Chain] Tool Error: {str(error)}")

    def get_logs(self) -> List[Dict[str, Any]]:
        return self.logs


class MsgChainTxAuditCallback(BaseCallbackHandler):
    """交易审计回调 — 记录所有链上交易操作"""

    def __init__(self):
        self.transactions: List[Dict[str, Any]] = []

    def on_tool_end(
        self,
        output: str,
        **kwargs: Any,
    ) -> None:
        import json
        try:
            data = json.loads(output)
            if isinstance(data, dict) and data.get("tx_hash"):
                self.transactions.append(data)
                print(f"[MSG Chain AUDIT] Transaction recorded: {data['tx_hash']}")
        except (json.JSONDecodeError, AttributeError):
            pass

    def get_tx_history(self) -> List[Dict[str, Any]]:
        return self.transactions


if __name__ == "__main__":
    from msgchain_langchain_agent import create_msgchain_agent

    callback = MsgChainLoggingCallback()
    agent = create_msgchain_agent(
        model="gpt-4o",
        include_tx_tools=True,
        verbose=False,
    )

    result = agent.invoke(
        {"input": "Show me the latest block height on MSG Chain"},
        config={"callbacks": [callback]},
    )

    print("\\n--- Operation Logs ---")
    for log in callback.get_logs():
        print(f"[{log['event']}] {log.get('tool', '')}")

2.7 带 Retry 的 Tool 封装

"""
msgchain_langchain_retry.py — MSG Chain Tool 重试与容错封装
"""

import time
import json
from functools import wraps
from typing import Optional, Type, Any, Callable

from pydantic import BaseModel, Field
from langchain.tools import BaseTool
from langchain.callbacks.manager import CallbackManagerForToolRun


def retry_on_failure(
    max_retries: int = 3,
    delay: float = 1.0,
    backoff: float = 2.0,
    exceptions: tuple = (Exception,),
):
    """重试装饰器"""
    def decorator(func: Callable):
        @wraps(func)
        def wrapper(*args, **kwargs):
            last_exception = None
            current_delay = delay
            for attempt in range(max_retries):
                try:
                    return func(*args, **kwargs)
                except exceptions as e:
                    last_exception = e
                    if attempt < max_retries - 1:
                        time.sleep(current_delay)
                        current_delay *= backoff
            raise last_exception
        return wrapper
    return decorator


class RetryableBalanceTool(BaseTool):
    """带重试机制的余额查询工具"""

    name: str = "msgchain_retry_balance"
    description: str = (
        "Query MSG Chain balance with automatic retry on failure. "
        "Useful when the RPC endpoint may be temporarily unavailable."
    )

    class BalanceInput(BaseModel):
        address: str = Field(description="MSG Chain address")
        denom: str = Field(default="uaimgs", description="Denom")

    args_schema: Type[BaseModel] = BalanceInput

    @retry_on_failure(max_retries=3, delay=0.5, backoff=2.0)
    def _run(
        self,
        address: str,
        denom: str = "uaimgs",
        run_manager: Optional[CallbackManagerForToolRun] = None,
    ) -> str:
        from msgchain_sdk_wrapper import MsgChainClient
        client = MsgChainClient()
        result = client.query_balance(address, denom)
        return json.dumps(result, indent=2, ensure_ascii=False)

2.8 TypeScript LangChain 集成

// msgchain-langchain.ts — MSG Chain LangChain TypeScript 集成

import { DynamicTool, Tool } from "langchain/tools";
import { AgentExecutor, createOpenAIFunctionsAgent } from "langchain/agents";
import { ChatOpenAI } from "@langchain/openai";
import { ChatPromptTemplate } from "@langchain/core/prompts";

class MsgChainClient {
  private rpc: string;
  private rest: string;

  constructor(rpc?: string, rest?: string) {
    this.rpc = rpc || "https://rpc.msgchain.org";
    this.rest = rest || "https://rest.msgchain.org";
  }

  async queryBalance(address: string, denom: string = "uaimgs"): Promise<any> {
    const resp = await fetch(this.rpc, {
      method: "POST",
      headers: { "Content-Type": "application/json" },
      body: JSON.stringify({
        jsonrpc: "2.0",
        id: 1,
        method: "abci_query",
        params: {
          path: `/cosmos.bank.v1beta1/balances/${address}/by_denom`,
          data: Buffer.from(denom).toString("hex"),
        },
      }),
    });
    return resp.json();
  }

  async queryContract(
    contract: string,
    msg: Record<string, any>
  ): Promise<any> {
    const resp = await fetch(this.rpc, {
      method: "POST",
      headers: { "Content-Type": "application/json" },
      body: JSON.stringify({
        jsonrpc: "2.0",
        id: 1,
        method: "abci_query",
        params: {
          path: `/cosmwasm.wasm.v1/contract/${contract}/smart`,
          data: Buffer.from(JSON.stringify(msg)).toString("hex"),
        },
      }),
    });
    return resp.json();
  }

  async queryLatestBlock(): Promise<any> {
    const resp = await fetch(this.rpc, {
      method: "POST",
      headers: { "Content-Type": "application/json" },
      body: JSON.stringify({
        jsonrpc: "2.0",
        id: 1,
        method: "block",
        params: {},
      }),
    });
    return resp.json();
  }

  async queryValidators(): Promise<any> {
    const resp = await fetch(this.rpc, {
      method: "POST",
      headers: { "Content-Type": "application/json" },
      body: JSON.stringify({
        jsonrpc: "2.0",
        id: 1,
        method: "validators",
        params: { height: 0, page: 1, per_page: 100 },
      }),
    });
    return resp.json();
  }
}

export function createMsgChainTools(): Tool[] {
  const client = new MsgChainClient();

  return [
    new DynamicTool({
      name: "msgchain_balance",
      description: "Query MSG Chain account balance. Input: address",
      func: async (address: string) => {
        try {
          const result = await client.queryBalance(address.trim());
          return JSON.stringify(result, null, 2);
        } catch (e: any) {
          return `Error: ${e.message}`;
        }
      },
    }),

    new DynamicTool({
      name: "msgchain_latest_block",
      description: "Get the latest block height on MSG Chain",
      func: async () => {
        try {
          const result = await client.queryLatestBlock();
          const height = result?.result?.block?.header?.height || "unknown";
          return `Latest block height: ${height}`;
        } catch (e: any) {
          return `Error: ${e.message}`;
        }
      },
    }),

    new DynamicTool({
      name: "msgchain_validators",
      description: "Query MSG Chain validator set",
      func: async () => {
        try {
          const result = await client.queryValidators();
          return JSON.stringify(result, null, 2);
        } catch (e: any) {
          return `Error: ${e.message}`;
        }
      },
    }),

    new DynamicTool({
      name: "msgchain_contract_query",
      description:
        "Query a CosmWasm contract on MSG Chain. "
        + "Input format: contract_address | json_query_msg",
      func: async (input: string) => {
        try {
          const parts = input.split("|").map((s) => s.trim());
          if (parts.length !== 2) {
            return "Invalid format. Use: contract_addr | json_msg";
          }
          const result = await client.queryContract(
            parts[0], JSON.parse(parts[1])
          );
          return JSON.stringify(result, null, 2);
        } catch (e: any) {
          return `Error: ${e.message}`;
        }
      },
    }),
  ];
}

export async function createMsgChainAgent() {
  const tools = createMsgChainTools();
  const llm = new ChatOpenAI({ model: "gpt-4o", temperature: 0.1 });

  const prompt = ChatPromptTemplate.fromMessages([
    ["system",
      `You are an MSG Chain blockchain assistant.
       MSG Chain is a CosmWasm Layer 1 with Dilithium-5 signatures.
       Addresses use the "msg" prefix.
       Use the available tools to query chain state.`],
    ["human", "{input}"],
    ["placeholder", "{agent_scratchpad}"],
  ]);

  const agent = await createOpenAIFunctionsAgent({ llm, tools, prompt });
  return AgentExecutor.fromAgentAndTools({ agent, tools, verbose: true });
}

async function main() {
  const agent = await createMsgChainAgent();
  const queries = [
    "What is the latest block height on MSG Chain?",
    "Show me the current validators",
  ];
  for (const query of queries) {
    console.log(`\\nQuery: ${query}`);
    const result = await agent.invoke({ input: query });
    console.log(`Result: ${result.output}`);
  }
}

3. CrewAI 多 Agent 集成

3.1 CrewAI 简介

CrewAI 是一个多 Agent 协作框架,允许定义具有特定角色、目标、工具的 Agent,并将其组织成 Crew(团队)来完成任务。核心概念:

3.2 MSG Chain Agent 定义

"""
crewai_msgchain_agents.py — MSG Chain CrewAI Agent 定义
"""

import json
from typing import List, Optional

from crewai import Agent, Task, Crew, Process
from crewai.tools import BaseTool as CrewAITool
from langchain.tools import BaseTool as LangChainTool

from msgchain_langchain_tools import (
    MsgChainBalanceTool,
    MsgChainAllBalancesTool,
    MsgChainContractQueryTool,
    MsgChainTokenInfoTool,
    MsgChainBlockQueryTool,
    MsgChainValidatorTool,
    MsgChainDelegationTool,
    MsgChainProposalTool,
    MsgChainInflationTool,
    MsgChainLatestHeightTool,
    MsgChainTxQueryTool,
)
from msgchain_langchain_tx_tools import (
    MsgChainTransferTool,
    MsgChainContractExecuteTool,
    MsgChainDelegateTool,
    MsgChainVoteTool,
)


# ── LangChain → CrewAI Tool 适配器 ──────────────────────

class LangChainToCrewAITool(CrewAITool):
    """将 LangChain BaseTool 适配为 CrewAI BaseTool"""

    name: str = ""
    description: str = ""

    def __init__(self, langchain_tool: LangChainTool):
        super().__init__(
            name=langchain_tool.name,
            description=langchain_tool.description,
        )
        self._langchain_tool = langchain_tool

    def _run(self, **kwargs) -> str:
        return self._langchain_tool._run(**kwargs)


# ── MSG Chain Agent 工厂 ────────────────────────────────

def create_data_analyst_agent(
    model: str = "gpt-4o",
    temperature: float = 0.1,
) -> Agent:
    """数据分析师 Agent — 负责链上数据查询与分析"""

    tools = [
        LangChainToCrewAITool(MsgChainBalanceTool()),
        LangChainToCrewAITool(MsgChainAllBalancesTool()),
        LangChainToCrewAITool(MsgChainContractQueryTool()),
        LangChainToCrewAITool(MsgChainTokenInfoTool()),
        LangChainToCrewAITool(MsgChainBlockQueryTool()),
        LangChainToCrewAITool(MsgChainValidatorTool()),
        LangChainToCrewAITool(MsgChainDelegationTool()),
        LangChainToCrewAITool(MsgChainProposalTool()),
        LangChainToCrewAITool(MsgChainInflationTool()),
        LangChainToCrewAITool(MsgChainLatestHeightTool()),
        LangChainToCrewAITool(MsgChainTxQueryTool()),
    ]

    return Agent(
        role="MSG Chain Data Analyst",
        goal=(
            "Query and analyze MSG Chain on-chain data accurately. "
            "Provide insights on balances, token holdings, validator sets, "
            "governance proposals, and blockchain state."
        ),
        backstory=(
            "You are an expert blockchain data analyst specializing in MSG Chain. "
            "You have deep knowledge of CosmWasm smart contracts, CW20 tokens, "
            "staking mechanics, and Cosmos SDK modules."
        ),
        tools=tools,
        allow_delegation=True,
        verbose=True,
        llm_config={"model": model, "temperature": temperature},
    )


def create_defi_trader_agent(
    model: str = "gpt-4o",
    temperature: float = 0.2,
) -> Agent:
    """DeFi 交易 Agent — 负责执行链上交易操作"""

    tools = [
        LangChainToCrewAITool(MsgChainBalanceTool()),
        LangChainToCrewAITool(MsgChainTransferTool()),
        LangChainToCrewAITool(MsgChainContractExecuteTool()),
        LangChainToCrewAITool(MsgChainTokenInfoTool()),
    ]

    return Agent(
        role="DeFi Trader",
        goal=(
            "Execute profitable trades and manage positions on MSG Chain DeFi protocols. "
            "Monitor market conditions, execute swaps, provide liquidity, "
            "and optimize returns."
        ),
        backstory=(
            "You are an experienced DeFi trader operating on MSG Chain. "
            "You understand AMM mechanics, liquidity pools, yield farming, "
            "and arbitrage opportunities."
        ),
        tools=tools,
        allow_delegation=True,
        verbose=True,
        llm_config={"model": model, "temperature": temperature},
    )


def create_staking_agent(
    model: str = "gpt-4o",
    temperature: float = 0.1,
) -> Agent:
    """质押管理 Agent — 负责验证者分析和委托管理"""

    tools = [
        LangChainToCrewAITool(MsgChainBalanceTool()),
        LangChainToCrewAITool(MsgChainValidatorTool()),
        LangChainToCrewAITool(MsgChainDelegationTool()),
        LangChainToCrewAITool(MsgChainInflationTool()),
        LangChainToCrewAITool(MsgChainDelegateTool()),
    ]

    return Agent(
        role="Staking Manager",
        goal=(
            "Optimize staking returns on MSG Chain. "
            "Analyze validator performance, commission rates, "
            "and recommend optimal delegation strategies."
        ),
        backstory=(
            "You are a proof-of-stake optimization specialist. "
            "You track validator uptime, commission rates, "
            "and voting power distribution."
        ),
        tools=tools,
        allow_delegation=True,
        verbose=True,
        llm_config={"model": model, "temperature": temperature},
    )


def create_governance_agent(
    model: str = "gpt-4o",
    temperature: float = 0.2,
) -> Agent:
    """治理 Agent — 负责提案分析和投票"""

    tools = [
        LangChainToCrewAITool(MsgChainProposalTool()),
        LangChainToCrewAITool(MsgChainVoteTool()),
        LangChainToCrewAITool(MsgChainBalanceTool()),
    ]

    return Agent(
        role="Governance Analyst",
        goal=(
            "Analyze MSG Chain governance proposals and "
            "recommend voting decisions based on protocol health "
            "and community benefit."
        ),
        backstory=(
            "You are a governance participant on MSG Chain. "
            "You carefully analyze each proposal and make informed "
            "voting recommendations."
        ),
        tools=tools,
        allow_delegation=True,
        verbose=True,
        llm_config={"model": model, "temperature": temperature},
    )


def create_security_agent(
    model: str = "gpt-4o",
    temperature: float = 0.1,
) -> Agent:
    """安全审计 Agent — 负责合约安全分析"""

    tools = [
        LangChainToCrewAITool(MsgChainContractQueryTool()),
        LangChainToCrewAITool(MsgChainTxQueryTool()),
    ]

    return Agent(
        role="Security Auditor",
        goal=(
            "Audit CosmWasm smart contracts on MSG Chain for vulnerabilities. "
            "Analyze contract state and transaction patterns."
        ),
        backstory=(
            "You are a smart contract security expert. "
            "You specialize in CosmWasm security and on-chain forensic analysis."
        ),
        tools=tools,
        allow_delegation=True,
        verbose=True,
        llm_config={"model": model, "temperature": temperature},
    )

3.3 Task 定义

"""
crewai_msgchain_tasks.py — MSG Chain CrewAI Task 定义
"""

from crewai import Task
from crewai import Agent


def create_balance_analysis_task(address: str, agent: Agent) -> Task:
    """分析指定地址的余额"""
    return Task(
        description=(
            f"Query and analyze the balance of address {address} on MSG Chain. "
            "Include: all token balances, equivalent MSG value, "
            "and recommendations for asset allocation."
        ),
        expected_output=(
            "A comprehensive balance report with token holdings, "
            "values, and actionable recommendations."
        ),
        agent=agent,
    )


def create_validator_analysis_task(agent: Agent) -> Task:
    """分析验证者集合"""
    return Task(
        description=(
            "Query the current validator set on MSG Chain. "
            "Analyze voting power distribution, commission rates, "
            "and identify the top validators. "
            "Provide staking recommendations."
        ),
        expected_output="Validator analysis report with delegation recommendations.",
        agent=agent,
    )


def create_governance_analysis_task(agent: Agent) -> Task:
    """分析治理提案"""
    return Task(
        description=(
            "Query all active governance proposals on MSG Chain. "
            "For each proposal, analyze its content, potential impact, "
            "and provide voting recommendations."
        ),
        expected_output="Governance summary with recommended voting positions.",
        agent=agent,
    )


def create_market_analysis_task(agent: Agent) -> Task:
    """分析市场状态"""
    return Task(
        description=(
            "Query top DEX pools on MSG Chain, analyze trading volumes, "
            "liquidity depth, and price trends. "
            "Identify potential arbitrage opportunities."
        ),
        expected_output="Market analysis with pool statistics and opportunities.",
        agent=agent,
    )


def create_swap_execution_task(
    token_in: str, token_out: str, amount: str, agent: Agent,
) -> Task:
    """执行代币交换"""
    return Task(
        description=(
            f"Execute a swap of {amount} of {token_in} for {token_out} "
            "on the best available DEX on MSG Chain. "
            "Consider slippage and fees."
        ),
        expected_output="Swap result with transaction hash and execution price.",
        agent=agent,
    )


def create_delegation_task(validator: str, amount: str, agent: Agent) -> Task:
    """执行委托操作"""
    return Task(
        description=(
            f"Delegate {amount} uaimgs to validator {validator} on MSG Chain. "
            "First verify the validator is active and check current delegation state."
        ),
        expected_output="Delegation result with transaction hash.",
        agent=agent,
    )


def create_contract_audit_task(contract_addr: str, agent: Agent) -> Task:
    """审计合约安全性"""
    return Task(
        description=(
            f"Audit CosmWasm contract {contract_addr} on MSG Chain. "
            "Query contract state and check for known vulnerability patterns."
        ),
        expected_output="Security audit report with risk assessment.",
        agent=agent,
    )

3.4 Crew 编排

"""
crewai_msgchain_crew.py — MSG Chain CrewAI Crew 编排
"""

import json
from typing import List, Optional

from crewai import Agent, Task, Crew, Process

from crewai_msgchain_agents import (
    create_data_analyst_agent,
    create_defi_trader_agent,
    create_staking_agent,
    create_governance_agent,
    create_security_agent,
)
from crewai_msgchain_tasks import (
    create_balance_analysis_task,
    create_validator_analysis_task,
    create_governance_analysis_task,
    create_market_analysis_task,
    create_swap_execution_task,
    create_delegation_task,
    create_contract_audit_task,
)


class MsgChainCrewManager:
    """MSG Chain Crew 管理器"""

    def __init__(self, model: str = "gpt-4o", verbose: bool = True):
        self.model = model
        self.verbose = verbose

    # ── 数据分析 Crew ──────────────────────────────────

    def create_analysis_crew(self, addresses: Optional[List[str]] = None) -> Crew:
        """数据分析 Crew"""
        analyst = create_data_analyst_agent(model=self.model)
        tasks = []

        if addresses:
            for addr in addresses:
                tasks.append(create_balance_analysis_task(addr, analyst))

        tasks.append(create_validator_analysis_task(analyst))
        tasks.append(create_governance_analysis_task(analyst))

        return Crew(
            agents=[analyst],
            tasks=tasks,
            process=Process.sequential,
            verbose=self.verbose,
        )

    # ── DeFi 交易 Crew ─────────────────────────────────

    def create_defi_crew(self, token_in: str, token_out: str, amount: str) -> Crew:
        """DeFi 交易 Crew"""
        analyst = create_data_analyst_agent(model=self.model)
        trader = create_defi_trader_agent(model=self.model)

        return Crew(
            agents=[analyst, trader],
            tasks=[
                create_market_analysis_task(analyst),
                create_swap_execution_task(token_in, token_out, amount, trader),
            ],
            process=Process.sequential,
            verbose=self.verbose,
        )

    # ── 全功能 Crew ────────────────────────────────────

    def create_full_crew(self) -> Crew:
        """全功能 Crew"""
        analyst = create_data_analyst_agent(model=self.model)
        trader = create_defi_trader_agent(model=self.model)
        staking_mgr = create_staking_agent(model=self.model)
        gov_analyst = create_governance_agent(model=self.model)
        security = create_security_agent(model=self.model)

        tasks = [
            Task(
                description="Query MSG Chain latest block height and active validators.",
                expected_output="Current chain state overview.",
                agent=analyst,
            ),
            Task(
                description="Analyze top DEX pools and market conditions.",
                expected_output="Market analysis report.",
                agent=analyst,
            ),
            Task(
                description="Review active governance proposals.",
                expected_output="Governance voting recommendations.",
                agent=gov_analyst,
            ),
            Task(
                description="Recommend optimal delegation strategy.",
                expected_output="Staking optimization strategy.",
                agent=staking_mgr,
            ),
            Task(
                description="Audit recently deployed contracts for risks.",
                expected_output="Security audit report.",
                agent=security,
            ),
            Task(
                description="Compile all findings into a comprehensive MSG Chain report.",
                expected_output="Final ecosystem report.",
                agent=analyst,
            ),
        ]

        return Crew(
            agents=[analyst, trader, staking_mgr, gov_analyst, security],
            tasks=tasks,
            process=Process.sequential,
            verbose=self.verbose,
        )

    # ── 分层流程(Hierarchical) ───────────────────────

    def create_hierarchical_crew(self) -> Crew:
        """使用分层流程的 Crew"""
        analyst = create_data_analyst_agent(model=self.model)
        trader = create_defi_trader_agent(model=self.model)
        staking_mgr = create_staking_agent(model=self.model)

        return Crew(
            agents=[analyst, trader, staking_mgr],
            tasks=[Task(
                description="Gather MSG Chain market data and execute optimal strategies.",
                expected_output="Completed operations with transaction hashes.",
            )],
            process=Process.hierarchical,
            verbose=self.verbose,
            manager_llm_config={"model": self.model, "temperature": 0.1},
        )


if __name__ == "__main__":
    manager = MsgChainCrewManager(verbose=True)
    print("MSG Chain CrewAI Multi-Agent System")

    analysis_crew = manager.create_analysis_crew(
        addresses=["msg1qgjw2qfc5v82j5c8dl5qen8wxm3r3v4e9k4cfk"]
    )
    analysis_result = analysis_crew.kickoff()
    print(f"\\nAnalysis Resul[未公开路径]}")

    defi_result = manager.create_defi_crew("uusdc", "uaimgs", "1000000").kickoff()
    print(f"\\nDeFi Resul[未公开路径]}")

    full_result = manager.create_full_crew().kickoff()
    print(f"\\nFull Repor[未公开路径]}")

3.5 高级 Agent 协作模式

"""
crewai_msgchain_advanced.py — MSG Chain CrewAI 高级协作模式
"""

from crewai import Agent, Task, Crew, Process
from crewai_msgchain_agents import LangChainToCrewAITool
from msgchain_langchain_tools import MsgChainContractQueryTool
from msgchain_langchain_tx_tools import MsgChainContractExecuteTool


def create_yield_farmer_agent(model: str = "gpt-4o") -> Agent:
    """收益农场 Agent — 寻找最优收益率"""
    return Agent(
        role="Yield Farmer",
        goal="Identify and execute the highest yield farming opportunities.",
        backstory="You are a yield optimization expert monitoring liquidity pools.",
        tools=[
            LangChainToCrewAITool(MsgChainContractQueryTool()),
            LangChainToCrewAITool(MsgChainContractExecuteTool()),
        ],
        allow_delegation=True,
        verbose=True,
        llm_config={"model": model, "temperature": 0.1},
    )


def create_arbitrage_agent(model: str = "gpt-4o") -> Agent:
    """套利 Agent — 发现并执行套利机会"""
    return Agent(
        role="Arbitrage Hunter",
        goal="Find and execute profitable arbitrage opportunities across MSG Chain DEXes.",
        backstory="You monitor prices across all MSG Chain DEXes simultaneously.",
        tools=[
            LangChainToCrewAITool(MsgChainContractQueryTool()),
            LangChainToCrewAITool(MsgChainContractExecuteTool()),
        ],
        allow_delegation=True,
        verbose=True,
        llm_config={"model": model, "temperature": 0.1},
    )


class MsgChainCollaborativeCrew:
    """高级协作 Crew 模式"""

    def __init__(self, model: str = "gpt-4o", verbose: bool = True):
        self.model = model
        self.verbose = verbose

    def defi_squad(self) -> Crew:
        """DeFi 特遣队"""
        from crewai_msgchain_agents import create_data_analyst_agent, create_defi_trader_agent

        analyst = create_data_analyst_agent(model=self.model)
        trader = create_defi_trader_agent(model=self.model)
        yield_farmer = create_yield_farmer_agent(model=self.model)
        arb = create_arbitrage_agent(model=self.model)

        tasks = [
            Task(
                description="Survey MSG Chain DeFi ecosystem and collect key metrics.",
                expected_output="DeFi landscape report with key metrics.",
                agent=analyst,
            ),
            Task(
                description="Identify top opportunities including best swap routes and APY.",
                expected_output="Ranked list of opportunities.",
                agent=analyst,
            ),
            Task(
                description="Execute identified opportunities and arbitrage trades.",
                expected_output="Execution results with transaction hashes.",
                agent=trader,
            ),
            Task(
                description="Manage all active positions and provide P&L dashboard.",
                expected_output="Position management report.",
                agent=yield_farmer,
            ),
        ]

        return Crew(
            agents=[analyst, trader, yield_farmer, arb],
            tasks=tasks,
            process=Process.sequential,
            verbose=self.verbose,
        )

    def governance_squad(self) -> Crew:
        """治理特遣队"""
        from crewai_msgchain_agents import (
            create_data_analyst_agent, create_governance_agent, create_security_agent,
        )

        analyst = create_data_analyst_agent(model=self.model)
        gov = create_governance_agent(model=self.model)
        security = create_security_agent(model=self.model)

        tasks = [
            Task(
                description="Collect all active governance proposals and their details.",
                expected_output="Complete proposal inventory.",
                agent=analyst,
            ),
            Task(
                description="Audit each proposal's smart contract changes for vulnerabilities.",
                expected_output="Security assessment per proposal.",
                agent=security,
            ),
            Task(
                description="Recommend voting positions based on technical and community analysis.",
                expected_output="Voting guide with recommendations.",
                agent=gov,
            ),
        ]

        return Crew(
            agents=[analyst, security, gov],
            tasks=tasks,
            process=Process.sequential,
            verbose=self.verbose,
        )


if __name__ == "__main__":
    collab = MsgChainCollaborativeCrew(verbose=True)
    print("MSG Chain Collaborative Crew System")

    defi_result = collab.defi_squad().kickoff()
    print(f"\\nDeFi Squad Resul[未公开路径]}")

    gov_result = collab.governance_squad().kickoff()
    print(f"\\nGovernance Squad Resul[未公开路径]}")

3.6 CrewAI 独立工具封装

"""
crewai_msgchain_tools.py — MSG Chain 原生 CrewAI 工具
"""

from crewai.tools import BaseTool
from pydantic import BaseModel, Field
from typing import Type, Optional


class BalanceInput(BaseModel):
    address: str = Field(description="MSG Chain bech32 address")


class ContractQueryInput(BaseModel):
    contract: str = Field(description="CosmWasm contract address")
    query: str = Field(description="JSON query string")


class CrewAIMsgChainBalanceTool(BaseTool):
    name: str = "Query MSG Chain Balance"
    description: str = "Query the token balance of an MSG Chain address."
    args_schema: Type[BaseModel] = BalanceInput

    def _run(self, address: str) -> str:
        from msgchain_sdk_wrapper import MsgChainClient
        client = MsgChainClient()
        result = client.query_balance(address)
        return str(result)


class CrewAIMsgChainContractTool(BaseTool):
    name: str = "Query MSG Chain Contract"
    description: str = "Query a CosmWasm smart contract on MSG Chain."
    args_schema: Type[BaseModel] = ContractQueryInput

    def _run(self, contract: str, query: str) -> str:
        import json
        from msgchain_sdk_wrapper import MsgChainClient
        client = MsgChainClient()
        result = client.query_contract(contract, json.loads(query))
        return str(result)

4. AutoGen 集成

4.1 AutoGen 简介

AutoGen(微软)是一个多 Agent 对话框架,允许创建可对话的 AI Agent,支持工具调用和人类参与。核心概念:

4.2 MSG Chain AutoGen Assistant

"""
autogen_msgchain.py — MSG Chain AutoGen 集成
"""

import json
import os
from typing import Dict, List, Optional, Any

from dotenv import load_dotenv
load_dotenv()

try:
    import autogen
    from autogen import (
        AssistantAgent, UserProxyAgent,
        GroupChat, GroupChatManager, ConversableAgent,
    )
except ImportError:
    print("Install pyautogen: pip install pyautogen")
    raise


# ── MSG Chain 函数定义 ─────────────────────────────────

def query_balance(address: str, denom: str = "uaimgs") -> str:
    from msgchain_sdk_wrapper import MsgChainClient
    client = MsgChainClient()
    result = client.query_balance(address, denom)
    return json.dumps(result, indent=2, ensure_ascii=False)


def query_contract(contract_addr: str, query_msg: str) -> str:
    from msgchain_sdk_wrapper import MsgChainClient
    client = MsgChainClient()
    result = client.query_contract(contract_addr, json.loads(query_msg))
    return json.dumps(result, indent=2, ensure_ascii=False)


def get_latest_block() -> str:
    from msgchain_sdk_wrapper import MsgChainClient
    client = MsgChainClient()
    height = client.query_latest_block_height()
    return f"Latest block height: {height}"


def get_validators() -> str:
    from msgchain_sdk_wrapper import MsgChainClient
    client = MsgChainClient()
    result = client.query_validators()
    return json.dumps(result, indent=2, ensure_ascii=False)


def get_token_info(contract_addr: str) -> str:
    from msgchain_sdk_wrapper import MsgChainClient
    client = MsgChainClient()
    info = client.query_contract(contract_addr, {"token_info": {}})
    supply = client.query_contract(contract_addr, {"total_supply": {}})
    return json.dumps({"token_info": info, "total_supply": supply},
                       indent=2, ensure_ascii=False)


def get_proposals() -> str:
    from msgchain_sdk_wrapper import MsgChainClient
    client = MsgChainClient()
    result = client.query_proposals()
    return json.dumps(result, indent=2, ensure_ascii=False)


def simulate_transfer(recipient: str, amount: str, memo: str = "") -> str:
    return json.dumps({
        "action": "simulate_transfer",
        "recipient": recipient,
        "amount": amount,
        "note": "In production, sign with Dilithium-5 and broadcast",
    }, indent=2, ensure_ascii=False)


# 函数映射
MSG_CHAIN_FUNCTIONS: Dict[str, callable] = {
    "query_balance": query_balance,
    "query_contract": query_contract,
    "get_latest_block": get_latest_block,
    "get_validators": get_validators,
    "get_token_info": get_token_info,
    "get_proposals": get_proposals,
    "simulate_transfer": simulate_transfer,
}

# 函数 Schema
MSG_CHAIN_FUNCTION_SCHEMAS = [
    {"function": {
        "name": "query_balance",
        "description": "Query MSG Chain token balance",
        "parameters": {
            "type": "object",
            "properties": {
                "address": {"type": "string", "description": "MSG Chain address"},
                "denom": {"type": "string", "description": "Denom (default: uaimgs)"},
            },
            "required": ["address"],
        },
    }},
    {"function": {
        "name": "query_contract",
        "description": "Query CosmWasm smart contract",
        "parameters": {
            "type": "object",
            "properties": {
                "contract_addr": {"type": "string"},
                "query_msg": {"type": "string"},
            },
            "required": ["contract_addr", "query_msg"],
        },
    }},
    {"function": {
        "name": "get_latest_block",
        "description": "Get latest MSG Chain block height",
        "parameters": {"type": "object", "properties": {}},
    }},
    {"function": {
        "name": "get_validators",
        "description": "Get MSG Chain validator set",
        "parameters": {"type": "object", "properties": {}},
    }},
    {"function": {
        "name": "get_token_info",
        "description": "Get CW20 token information",
        "parameters": {
            "type": "object",
            "properties": {"contract_addr": {"type": "string"}},
            "required": ["contract_addr"],
        },
    }},
    {"function": {
        "name": "get_proposals",
        "description": "Get governance proposals",
        "parameters": {"type": "object", "properties": {}},
    }},
]


def create_msgchain_assistant(
    name: str = "MSG_Assistant",
    llm_config: Optional[Dict] = None,
) -> AssistantAgent:
    """创建 MSG Chain 助手 Agent"""

    default_llm_config = {
        "config_list": [{"model": "gpt-4o", "api_key": os.getenv("OPENAI_API_KEY")}],
        "temperature": 0.1,
        "functions": MSG_CHAIN_FUNCTION_SCHEMAS,
    }

    return AssistantAgent(
        name=name,
        system_message=(
            "You are an MSG Chain blockchain assistant with access to "
            "on-chain query and simulation tools.\\n"
            "MSG Chain is a CosmWasm Layer 1 with Dilithium-5 signatures.\\n"
            "Addresses: msg... prefix | Native token: uaimgs\\n\\n"
            "Available functions: query_balance, query_contract, "
            "get_latest_block, get_validators, get_token_info, "
            "get_proposals, simulate_transfer\\n\\n"
            "Always explain what you find and suggest next actions."
        ),
        llm_config=llm_config or default_llm_config,
        code_execution_config=False,
    )


def create_user_proxy(
    name: str = "User",
    human_input_mode: str = "NEVER",
) -> UserProxyAgent:
    """创建用户代理"""
    return UserProxyAgent(
        name=name,
        human_input_mode=human_input_mode,
        max_consecutive_auto_reply=10,
        is_termination_msg=lambda x: x.get("content", "").rstrip().endswith("TERMINATE"),
        code_execution_config={"work_dir": "autogen_workspace", "use_docker": False},
    )


# ── 单轮对话 ──────────────────────────────────────────

def run_single_query(query: str, llm_config: Optional[Dict] = None) -> List[Dict]:
    """执行单轮 MSG Chain 查询"""
    assistant = create_msgchain_assistant(llm_config=llm_config)
    user = create_user_proxy()

    for func_name, func in MSG_CHAIN_FUNCTIONS.items():
        user.register_function(function_map={func_name: func})

    user.initiate_chat(assistant, message=query)
    return assistant.chat_messages


# ── 多轮对话 ──────────────────────────────────────────

def run_multi_turn_conversation(tasks: List[str], llm_config: Optional[Dict] = None) -> None:
    """执行多轮 MSG Chain 对话任务"""
    assistant = create_msgchain_assistant(llm_config=llm_config)
    user = create_user_proxy(human_input_mode="TERMINATE")

    for func_name, func in MSG_CHAIN_FUNCTIONS.items():
        user.register_function(function_map={func_name: func})

    for i, task in enumerate(tasks):
        print(f"\\n--- Task {i+1}/{len(tasks)}: {task} ---")
        user.initiate_chat(assistant, message=task)


# ── 多 Agent 对话 ─────────────────────────────────────

def create_debate_agents(llm_config: Optional[Dict] = None) -> Dict[str, ConversableAgent]:
    """创建辩论型多 Agent"""
    config = llm_config or {
        "config_list": [{"model": "gpt-4o", "api_key": os.getenv("OPENAI_API_KEY")}],
        "temperature": 0.3,
    }

    return {
        "analyst": AssistantAgent(
            name="DataAnalyst",
            system_message="You query on-chain data and provide factual analysis.",
            llm_config=config,
            function_map={"query_balance": query_balance, "get_latest_block": get_latest_block},
        ),
        "trader": AssistantAgent(
            name="DeFiTrader",
            system_message="You recommend trading strategies based on data.",
            llm_config=config,
            function_map={"query_balance": query_balance, "query_contract": query_contract},
        ),
        "critic": AssistantAgent(
            name="Critic",
            system_message="You review recommendations and point out risks.",
            llm_config=config,
        ),
    }


# ── GroupChat 群聊 ────────────────────────────────────

def run_group_chat(topic: str, llm_config: Optional[Dict] = None) -> None:
    """运行多 Agent 群聊讨论 MSG Chain 话题"""
    agents = create_debate_agents(llm_config)
    config = llm_config or {
        "config_list": [{"model": "gpt-4o", "api_key": os.getenv("OPENAI_API_KEY")}],
        "temperature": 0.3,
    }

    groupchat = GroupChat(agents=list(agents.values()), messages=[], max_round=12)
    manager = GroupChatManager(groupchat=groupchat, llm_config=config)
    user = create_user_proxy()
    user.initiate_chat(manager, message=f"Discuss: {topic}")


# ── 主入口 ─────────────────────────────────────────────

if __name__ == "__main__":
    print("MSG Chain AutoGen Integration")

    result = run_single_query(
        "What is the latest block height on MSG Chain? "
        "Also check who the top validators are."
    )
    print(f"\\nResult messages: {len(result)} exchanges")

    run_multi_turn_conversation([
        "Query the balance of address msg1qgjw2qfc5v82j5c8dl5qen8wxm3r3v4e9k4cfk",
        "What governance proposals are currently active?",
    ])

    run_group_chat(
        "Should we delegate more tokens given current inflation rate?"
    )

4.3 两步验证工作流

"""
autogen_msgchain_2fa.py — 两步验证 AutoGen 工作流
"""

import json
from autogen import AssistantAgent, UserProxyAgent


def query_balance(address: str) -> str:
    from msgchain_sdk_wrapper import MsgChainClient
    client = MsgChainClient()
    return json.dumps(client.query_balance(address), indent=2, ensure_ascii=False)


def execute_if_approved(action: dict) -> str:
    """仅当人类批准后才执行"""
    return json.dumps({
        "status": "awaiting_human_approval",
        "action": action,
        "message": "Please confirm this transaction.",
    }, indent=2)


class TwoFactorMSGChainWorkflow:
    """两步验证: Propose -> Approve -> Execute"""

    def __init__(self):
        self.assistant = AssistantAgent(
            name="MSG_Assistant",
            system_message=(
                "Step 1: Propose action with full details. "
                "Step 2: Wait for human approval. "
                "Never execute without approval."
            ),
            llm_config={
                "config_list": [{"model": "gpt-4o", "api_key": "..."}],
                "temperature": 0.1,
                "functions": [
                    {"function": {
                        "name": "query_balance",
                        "description": "Query balance",
                        "parameters": {
                            "type": "object",
                            "properties": {"address": {"type": "string"}},
                            "required": ["address"],
                        },
                    }},
                    {"function": {
                        "name": "execute_if_approved",
                        "description": "Execute after approval",
                        "parameters": {
                            "type": "object",
                            "properties": {"action": {"type": "object"}},
                            "required": ["action"],
                        },
                    }},
                ],
            },
            function_map={
                "query_balance": query_balance,
                "execute_if_approved": execute_if_approved,
            },
        )
        self.user = UserProxyAgent(
            name="User",
            human_input_mode="ALWAYS",
            code_execution_config=False,
        )

    def run(self, task: str):
        self.user.initiate_chat(self.assistant, message=task)


if __name__ == "__main__":
    TwoFactorMSGChainWorkflow().run(
        "Transfer 5000000 uaimgs to msg1recipient..."
    )

5. LangGraph 流程编排

5.1 LangGraph 简介

LangGraph 是一个构建状态化、多参与者 AI 应用的框架。核心概念:

5.2 MSG Chain DeFi 状态机

"""
langgraph_msgchain.py — MSG Chain LangGraph 流程编排
"""

import json
from typing import TypedDict, Dict, List, Optional, Any, Literal

from langgraph.graph import StateGraph, END
from langgraph.checkpoint import MemorySaver
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage

from msgchain_sdk_wrapper import MsgChainClient, MsgChainRestClient


# ── 状态定义 ───────────────────────────────────────────

class DeFiState(TypedDict):
    address: str
    action: str
    params: Dict[str, Any]
    balance: Optional[Dict[str, Any]]
    all_balances: Optional[Dict[str, Any]]
    market_data: Optional[Dict[str, Any]]
    positions: Optional[List[Dict[str, Any]]]
    analysis: Optional[str]
    recommendation: Optional[str]
    risk_assessment: Optional[str]
    tx_hash: Optional[str]
    tx_result: Optional[Dict[str, Any]]
    execution_status: Literal["pending", "simulated", "success", "failed"]
    error: Optional[str]
    logs: List[str]


def initial_state(address: str, action: str = "") -> DeFiState:
    return {
        "address": address,
        "action": action,
        "params": {},
        "balance": None,
        "all_balances": None,
        "market_data": None,
        "positions": None,
        "analysis": None,
        "recommendation": None,
        "risk_assessment": None,
        "tx_hash": None,
        "tx_result": None,
        "execution_status": "pending",
        "error": None,
        "logs": [],
    }


# ── 节点函数 ───────────────────────────────────────────

def check_balance(state: DeFiState) -> DeFiState:
    client = MsgChainClient()
    logs = list(state.get("logs", []))
    try:
        logs.append("[1/6] Querying balance...")
        balance = client.query_balance(state["address"])
        all_balances = client.query_all_balances(state["address"])
        logs.append(f"  Balance retrieved")
        return {**state, "balance": balance, "all_balances": all_balances, "logs": logs}
    except Exception as e:
        logs.append(f"  Error: {str(e)}")
        return {**state, "error": str(e), "logs": logs, "execution_status": "failed"}


def analyze_market(state: DeFiState) -> DeFiState:
    client = MsgChainClient()
    rest = MsgChainRestClient()
    logs = list(state.get("logs", []))
    try:
        logs.append("[2/6] Analyzing market conditions...")
        block = client.query_block()
        validators = client.query_validators()
        inflation = rest.inflation()
        market_data = {
            "latest_block": block.get("result", {}).get("block", {}).get("header", {}),
            "validator_count": len(validators.get("result", {}).get("validators", [])),
            "inflation": inflation,
        }
        logs.append("  Market data collected")
        return {**state, "market_data": market_data, "logs": logs}
    except Exception as e:
        logs.append(f"  Error: {str(e)}")
        return {**state, "error": str(e), "logs": logs}


def assess_risk(state: DeFiState) -> DeFiState:
    llm = ChatOpenAI(model="gpt-4o", temperature=0.1)
    logs = list(state.get("logs", []))
    try:
        logs.append("[3/6] Assessing risk...")
        prompt = (
            f"Assess risk for address {state['address']} on MSG Chain.\\n"
            f"Balance: {json.dumps(state.get('all_balances', {}), indent=2)[:300]}\\n"
            f"Provide a risk score (1-10) and explanation."
        )
        response = llm.invoke([HumanMessage(content=prompt)])
        logs.append("  Risk assessed")
        return {**state, "risk_assessment": response.content, "logs": logs}
    except Exception as e:
        return {**state, "error": str(e), "logs": logs}


def generate_recommendation(state: DeFiState) -> DeFiState:
    llm = ChatOpenAI(model="gpt-4o", temperature=0.3)
    logs = list(state.get("logs", []))
    try:
        logs.append("[4/6] Generating recommendation...")
        prompt = (
            f"Based on MSG Chain data for {state['address']}:\\n"
            f"Risk: {state.get('risk_assessment', 'N/A')}\\n"
            f"Action: {state.get('action', 'analysis')}\\n"
            "Recommend specific actions: stake, trade, or hold."
        )
        response = llm.invoke([HumanMessage(content=prompt)])
        logs.append("  Recommendation generated")
        return {**state, "recommendation": response.content, "logs": logs}
    except Exception as e:
        return {**state, "error": str(e), "logs": logs}


def simulate_execution(state: DeFiState) -> DeFiState:
    logs = list(state.get("logs", []))
    if state.get("action") != "trade" or state.get("error"):
        logs.append("[5/6] Skipping execution")
        return {**state, "logs": logs}
    try:
        logs.append("[5/6] Simulating trade...")
        params = state.get("params", {})
        simulated = {
            "from": state["address"],
            "to": params.get("contract", ""),
            "value": params.get("amount", "0"),
            "status": "simulated",
        }
        logs.append(f"  Simulated TX prepared")
        return {**state, "tx_result": simulated, "execution_status": "simulated", "logs": logs}
    except Exception as e:
        return {**state, "error": str(e), "execution_status": "failed", "logs": logs}


def compile_report(state: DeFiState) -> DeFiState:
    logs = list(state.get("logs", []))
    logs.append("[6/6] Compiling report...")
    report = {
        "address": state["address"],
        "execution_status": state["execution_status"],
        "error": state.get("error"),
        "logs": logs,
    }
    return {**state, "analysis": json.dumps(report, indent=2, ensure_ascii=False), "logs": logs}


# ── 条件边 ─────────────────────────────────────────────

def decide_execution(state: DeFiState) -> str:
    if state.get("error"):
        return "skip"
    if state.get("action") == "trade" and state.get("params"):
        return "execute"
    return "skip"


# ── 构建图 ─────────────────────────────────────────────

def build_defi_workflow() -> StateGraph:
    workflow = StateGraph(DeFiState)

    workflow.add_node("check_balance", check_balance)
    workflow.add_node("analyze_market", analyze_market)
    workflow.add_node("assess_risk", assess_risk)
    workflow.add_node("generate_recommendation", generate_recommendation)
    workflow.add_node("simulate_execution", simulate_execution)
    workflow.add_node("compile_report", compile_report)

    workflow.set_entry_point("check_balance")

    workflow.add_edge("check_balance", "analyze_market")
    workflow.add_edge("analyze_market", "assess_risk")
    workflow.add_edge("assess_risk", "generate_recommendation")

    workflow.add_conditional_edges(
        "generate_recommendation",
        decide_execution,
        {"execute": "simulate_execution", "skip": "compile_report"},
    )

    workflow.add_edge("simulate_execution", "compile_report")
    workflow.add_edge("compile_report", END)

    return workflow


# ── 运行器 ─────────────────────────────────────────────

class MSGChainDeFiWorkflow:
    def __init__(self):
        self.graph = build_defi_workflow()
        self.memory = MemorySaver()
        self.app = self.graph.compile(checkpointer=self.memory)

    def run(self, address: str, action: str = "", params: Optional[Dict] = None,
            thread_id: str = "1") -> Dict[str, Any]:
        state = initial_state(address, action)
        if params:
            state["params"] = params
        return self.app.invoke(state, {"configurable": {"thread_id": thread_id}})

    def stream(self, address: str, action: str = "", params: Optional[Dict] = None,
               thread_id: str = "1"):
        state = initial_state(address, action)
        if params:
            state["params"] = params
        for event in self.app.stream(state, {"configurable": {"thread_id": thread_id}}):
            for node_name, node_state in event.items():
                yield node_name, node_state


# ── 质押专用工作流 ─────────────────────────────────────

def build_staking_workflow() -> StateGraph:
    workflow = StateGraph(DeFiState)

    def check_staking(state: DeFiState) -> DeFiState:
        client = MsgChainClient()
        logs = list(state.get("logs", []))
        logs.append("[Staking] Checking delegations...")
        delegations = client.query_delegations(state["address"])
        unbonding = client.query_unbonding(state["address"])
        return {**state, "positions": [
            {"type": "delegation", "data": delegations},
            {"type": "unbonding", "data": unbonding},
        ], "logs": logs}

    def analyze_validators(state: DeFiState) -> DeFiState:
        client = MsgChainClient()
        logs = list(state.get("logs", []))
        logs.append("[Staking] Analyzing validators...")
        return {**state, "market_data": client.query_validators(), "logs": logs}

    def recommend(state: DeFiState) -> DeFiState:
        llm = ChatOpenAI(model="gpt-4o", temperature=0.1)
        logs = list(state.get("logs", []))
        logs.append("[Staking] Generating recommendation...")
        prompt = f"Current positions: {json.dumps(state.get('positions', [{}])[0], indent=2)[:300]}\\nRecommend delegation strategy."
        response = llm.invoke([HumanMessage(content=prompt)])
        return {**state, "recommendation": response.content, "logs": logs}

    workflow.add_node("check_staking", check_staking)
    workflow.add_node("analyze_validators", analyze_validators)
    workflow.add_node("recommend", recommend)
    workflow.set_entry_point("check_staking")
    workflow.add_edge("check_staking", "analyze_validators")
    workflow.add_edge("analyze_validators", "recommend")
    workflow.add_edge("recommend", END)
    return workflow


# ── 主入口 ─────────────────────────────────────────────

if __name__ == "__main__":
    print("MSG Chain LangGraph Workflow Engine")

    defi = MSGChainDeFiWorkflow()
    result = defi.run(
        address="msg1qgjw2qfc5v82j5c8dl5qen8wxm3r3v4e9k4cfk",
        action="analyze",
    )

    print("\\n=== DeFi Workflow Result ===")
    for log in result.get("logs", []):
        print(log)
    print(f"\\nRecommendation: {result.get('recommendation', 'N/A')}")

    print("\\n=== Streaming ===")
    for node_name, node_state in defi.stream("msg1...", "trade", {"contract": "msg1dex", "amount": "1000000"}):
        print(f"Node: {node_name} -> Status: {node_state.get('execution_status', 'running')}")

6. 链下 Agent 与链上合约通信

6.1 架构概览

链下 AI Agent 与 MSG Chain 链上合约之间通过以下方式进行通信:


┌─────────────────┐     WebSocket       ┌─────────────────┐

│  链下 AI Agent   │ ◄──── 事件 ──────►  │  MSG Chain 节点  │

│  (LangChain/     │                    │  (RPC/WS)       │

│   CrewAI/AutoGen)│ ──── 交易 ────────► │                 │

│                 │                    │  ┌───────────┐  │

│  Dilithium-5    │                    │  │ CosmWasm  │  │

│  签名 + 广播     │                    │  │  合约     │  │

└─────────────────┘                    └──┴───────────┴──┘

6.2 WebSocket 事件监听

"""
onchain_events.py — MSG Chain 链上事件监听
"""

import json
import os
import asyncio
import websockets
from typing import Callable, Dict, Any, Optional
from datetime import datetime


class MsgChainEventMonitor:
    """MSG Chain 链上事件监控器"""

    def __init__(self, ws_url: Optional[str] = None):
        self.ws_url = ws_url or os.getenv(
            "MSG_CHAIN_WS", "wss://rpc.msgchain.org/websocket"
        )
        self.handlers: Dict[str, list[Callable]] = {}
        self.connection = None
        self.running = False

    def on(self, event_type: str, handler: Callable):
        """注册事件处理器"""
        if event_type not in self.handlers:
            self.handlers[event_type] = []
        self.handlers[event_type].append(handler)

    async def subscribe(self, query: str):
        """订阅事件查询"""
        async with websockets.connect(self.ws_url) as ws:
            self.connection = ws
            self.running = True

            # 发送订阅请求
            subscribe_msg = {
                "jsonrpc": "2.0",
                "method": "subscribe",
                "id": 1,
                "params": {"query": query},
            }
            await ws.send(json.dumps(subscribe_msg))

            # 监听事件
            while self.running:
                try:
                    response = await asyncio.wait_for(ws.recv(), timeout=30)
                    data = json.loads(response)
                    await self._dispatch("message", data)

                    # 解析事件
                    result = data.get("result", {})
                    if "events" in result:
                        for event_type, event_data in result["events"].items():
                            await self._dispatch(event_type, event_data)

                    # 解析交易事件
                    tx_result = result.get("data", {}).get("value", {}).get("TxResult", {})
                    if tx_result:
                        await self._dispatch("tx", tx_result)

                except asyncio.TimeoutError:
                    # 发送心跳
                    await ws.send(json.dumps({
                        "jsonrpc": "2.0",
                        "method": "ping",
                        "id": 1,
                    }))
                except websockets.exceptions.ConnectionClosed:
                    print("[WS] Connection closed, reconnecting...")
                    break

    async def _dispatch(self, event_type: str, data: Any):
        """分发事件到已注册的处理器"""
        handlers = self.handlers.get(event_type, [])
        handlers_all = self.handlers.get("*", [])  # 通配符处理器
        for handler in handlers + handlers_all:
            try:
                if asyncio.iscoroutinefunction(handler):
                    await handler(data)
                else:
                    handler(data)
            except Exception as e:
                print(f"[Event Handler Error] {e}")

    async def listen_tx_events(self):
        """监听所有交易事件"""
        await self.subscribe("tm.event = 'Tx'")

    async def listen_new_block(self):
        """监听新区块"""
        await self.subscribe("tm.event = 'NewBlock'")

    async def listen_contract_events(self, contract_addr: str):
        """监听特定合约的事件"""
        await self.subscribe(
            f"wasm.contract_address = '{contract_addr}'"
        )

    def stop(self):
        """停止监听"""
        self.running = False
        if self.connection:
            asyncio.create_task(self.connection.close())


# ── 事件处理器示例 ──────────────────────────────────────

def on_tx_event(tx_data: dict):
    """处理交易事件的回调"""
    height = tx_data.get("height", "unknown")
    tx_hash = tx_data.get("hash", "unknown")
    print(f"[TX Event] Height: {height}, Hash: {tx_hash}")


async def on_swap_event(event_data: dict):
    """处理 DEX 交换事件的回调"""
    print(f"[SWAP] Event received: {json.dumps(event_data, indent=2)[:200]}")


# ── Agent 触发器 ────────────────────────────────────────

class AgentTrigger:
    """基于链上事件的 Agent 触发器"""

    def __init__(self, agent_executor):
        self.monitor = MsgChainEventMonitor()
        self.agent = agent_executor

    def setup_triggers(self):
        """设置事件触发规则"""

        @self.monitor.on("tx")
        def on_transaction(tx_data):
            """每笔交易后检查是否需要 Agent 介入"""
            print(f"[Trigger] Transaction detected, checking conditions...")
            # Agent 可以分析交易并决定是否响应

        @self.monitor.on("transfer")
        def on_transfer(event_data):
            """检测到转账后触发分析"""
            amount = event_data.get("amount", "")
            if "1000000" in str(amount):  # 大额转账
                print(f"[Trigger] Large transfer detected! Alerting agent...")
                # agent.run("Analyze this large transfer...")

    async def start(self):
        """启动事件监听"""
        self.setup_triggers()
        await self.monitor.listen_tx_events()


# ── 运行 ─────────────────────────────────────────────────

async def main():
    monitor = MsgChainEventMonitor()

    # 注册处理器
    monitor.on("tx", on_tx_event)
    monitor.on("*", lambda d: print(f"[Event] {type(d).__name__}"))

    print("Starting MSG Chain event monitor...")
    await monitor.listen_tx_events()


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

6.3 Agent 签名并广播交易

"""
onchain_tx_signer.py — Agent 签名并广播 MSG Chain 交易
"""

import json
import os
import hashlib
from typing import Optional, Dict, Any


class Dilithium5Signer:
    """Dilithium-5 签名器(封装底层 SDK)"""

    def __init__(self, private_key_hex: Optional[str] = None):
        self.private_key = bytes.fromhex(
            private_key_hex or os.getenv("AGENT_PRIVATE_KEY", "")
        )
        self.public_key = None
        self.address = None

    def sign(self, tx_bytes: bytes) -> bytes:
        """签名交易"""
        # 生产环境使用:
        # from msgchain_sdk.dilithium import dilithium5_sign
        # return dilithium5_sign(self.private_key, tx_bytes)
        raise NotImplementedError(
            "Install msgchain-sdk for Dilithium-5 support"
        )

    def get_address(self) -> str:
        """派生 bech32 地址"""
        # from msgchain_sdk.bech32 import encode_bech32
        # return encode_bech32("msg", self.public_key)
        return os.getenv("AGENT_ADDRESS", "msg1agent_address")


class AgentTxBuilder:
    """Agent 交易构建器"""

    def __init__(self, signer: Dilithium5Signer):
        self.signer = signer
        self.chain_id = os.getenv("MSG_CHAIN_CHAIN_ID", "msg-chain-1")

    def build_send_tx(
        self,
        recipient: str,
        amount: str,
        denom: str = "uaimgs",
        memo: str = "",
    ) -> bytes:
        """构建转账交易"""
        from msgchain_sdk_wrapper import MsgChainClient

        sender = self.signer.get_address()
        client = MsgChainClient()
        account = client.query_account(sender)

        # 提取 account_number 和 sequence
        acc_data = account.get("result", {}).get("response", {}).get("value", {})
        if isinstance(acc_data, str):
            acc_data = json.loads(acc_data)
        account_number = acc_data.get("account_number", "0")
        sequence = acc_data.get("sequence", "0")

        tx_body = {
            "messages": [{
                "@type": "/cosmos.bank.v1beta1.MsgSend",
                "from_address": sender,
                "to_address": recipient,
                "amount": [{"denom": denom, "amount": amount}],
            }],
            "memo": memo,
            "timeout_height": "0",
            "extension_options": [],
            "non_critical_extension_options": [],
        }

        auth_info = {
            "signer_infos": [],
            "fee": {
                "amount": [{"denom": "uaimgs", "amount": "5000"}],
                "gas_limit": "200000",
                "payer": "",
                "granter": "",
            },
        }

        sign_doc = {
            "body_bytes": json.dumps(tx_body).encode().hex(),
            "auth_info_bytes": json.dumps(auth_info).encode().hex(),
            "chain_id": self.chain_id,
            "account_number": account_number,
        }

        # 序列化为 protobuf 并签名
        # tx_bytes = proto_encode(sign_doc)
        # signature = self.signer.sign(tx_bytes)

        return json.dumps(sign_doc).encode()

    def build_contract_execute_tx(
        self,
        contract_addr: str,
        msg: dict,
        funds: str = "",
    ) -> bytes:
        """构建合约执行交易"""
        sender = self.signer.get_address()
        funds_list = []
        if funds:
            import re
            for part in funds.split(","):
                m = re.match(r"(\d+)(\D+)", part.strip())
                if m:
                    funds_list.append({"denom": m.group(2), "amount": m.group(1)})

        execute_msg = {
            "@type": "/cosmwasm.wasm.v1.MsgExecuteContract",
            "sender": sender,
            "contract": contract_addr,
            "msg": msg,
            "funds": funds_list,
        }

        return json.dumps(execute_msg).encode()


class AgentTxBroadcaster:
    """Agent 交易广播器"""

    def __init__(self):
        from msgchain_sdk_wrapper import MsgChainClient
        self.client = MsgChainClient()

    def broadcast(self, tx_bytes: bytes, mode: str = "BROADCAST_MODE_SYNC") -> Dict:
        """广播交易"""
        result = self.client.broadcast_tx(tx_bytes, mode)
        return result

    def broadcast_and_wait(self, tx_bytes: bytes) -> Dict:
        """广播并等待确认"""
        result = self.client.broadcast_tx_commit(tx_bytes)
        return result

    def wait_for_tx(self, tx_hash: str, max_retries: int = 10) -> Optional[Dict]:
        """轮询等待交易确认"""
        import time
        for i in range(max_retries):
            try:
                result = self.client.get_tx(tx_hash)
                if result.get("result", {}).get("tx_response", {}).get("code") == 0:
                    return result
            except Exception:
                pass
            time.sleep(1)
        return None


# ── Agent 条件触发引擎 ──────────────────────────────────

class ConditionEngine:
    """事件条件触发引擎"""

    def __init__(self):
        self.rules = []

    def add_rule(
        self,
        condition: callable,
        action: callable,
        description: str = "",
    ):
        """添加触发规则"""
        self.rules.append({
            "condition": condition,
            "action": action,
            "description": description,
        })

    def evaluate(self, event_data: dict) -> list:
        """评估事件并返回匹配的动作"""
        triggered = []
        for rule in self.rules:
            try:
                if rule["condition"](event_data):
                    triggered.append(rule)
            except Exception as e:
                print(f"Rule evaluation error: {e}")
        return triggered


# ── 条件规则示例 ────────────────────────────────────────

def price_drop_condition(event: dict) -> bool:
    """价格下跌条件"""
    price = float(event.get("price", 1))
    return price < 0.95  # 下跌超过 5%


def large_tx_condition(event: dict) -> bool:
    """大额交易条件"""
    amount = int(event.get("amount", "0"))
    return amount > 1_000_000_000  # > 1000 MSG


def governance_alert_condition(event: dict) -> bool:
    """新提案条件"""
    return event.get("type") == "new_proposal"

6.4 自动做市(AMM)Agent 示例

"""
onchain_amm_agent.py — 链上 AMM 交互 Agent
"""

import json
import os
from typing import Optional


class AMMAgent:
    """与 MSG Chain 上 AMM 合约交互的 Agent"""

    def __init__(self, router_contract: str):
        self.router = router_contract
        from msgchain_sdk_wrapper import MsgChainClient
        self.client = MsgChainClient()

    def get_pool_info(self, pool_id: str) -> dict:
        """查询 AMM 池信息"""
        return self.client.query_contract(
            self.router,
            {"pool": {"pool_id": pool_id}},
        )

    def get_all_pools(self) -> list:
        """查询所有池"""
        result = self.client.query_contract(
            self.router,
            {"pools": {}},
        )
        return result.get("pools", [])

    def estimate_swap(
        self,
        pool_id: str,
        token_in: str,
        amount_in: str,
        token_out: str,
    ) -> dict:
        """估算交换输出"""
        return self.client.query_contract(
            self.router,
            {
                "simulate_swap": {
                    "pool_id": pool_id,
                    "token_in": token_in,
                    "amount_in": amount_in,
                    "token_out": token_out,
                }
            },
        )

    def find_best_route(
        self,
        token_in: str,
        token_out: str,
        amount_in: str,
    ) -> Optional[dict]:
        """寻找最优交换路径"""
        pools = self.get_all_pools()
        best = None
        best_amount = 0

        for pool in pools:
            pid = pool.get("pool_id", "")
            if not pid:
                continue
            estimate = self.estimate_swap(pid, token_in, amount_in, token_out)
            amount_out = int(estimate.get("amount_out", "0"))
            if amount_out > best_amount:
                best_amount = amount_out
                best = {
                    "pool_id": pid,
                    "amount_out": amount_out,
                    **estimate,
                }

        return best

    def execute_swap(
        self,
        pool_id: str,
        token_in: str,
        amount_in: str,
        token_out: str,
        min_amount_out: str,
    ) -> str:
        """构建交换执行交易"""
        swap_msg = {
            "execute_swap": {
                "pool_id": pool_id,
                "token_in": token_in,
                "amount_in": amount_in,
                "token_out": token_out,
                "min_amount_out": min_amount_out,
            }
        }

        # 返回待签名的交易消息
        return json.dumps(swap_msg, indent=2)


# ── 条件触发的自动交易 Agent ───────────────────────────

class AutoTraderAgent:
    """自动交易 Agent(事件驱动)"""

    def __init__(self, amm: AMMAgent, condition_engine: ConditionEngine):
        self.amm = amm
        self.conditions = condition_engine
        self.signer = Dilithium5Signer()
        self.broadcaster = AgentTxBroadcaster()
        self.monitor = MsgChainEventMonitor()

    async def start(self):
        """启动自动交易"""

        @self.monitor.on("price_update")
        async def on_price_update(data):
            rules = self.conditions.evaluate(data)
            for rule in rules:
                print(f"[AutoTrader] Rule triggered: {rule.get('description', '')}")
                await rule["action"](data)

        async def arbitrage_action(data):
            """套利动作"""
            token_in = data.get("token_in", "uusdc")
            token_out = data.get("token_out", "uaimgs")
            amount = data.get("amount", "1000000")

            best = self.amm.find_best_route(token_in, token_out, amount)
            if best:
                tx_bytes = self.amm.execute_swap(
                    best["pool_id"], token_in, amount, token_out, "0"
                )
                # result = self.broadcaster.broadcast(tx_bytes)
                print(f"[AutoTrader] Would execute swap: {json.dumps(best, indent=2)}")

        self.conditions.add_rule(
            condition=lambda d: float(d.get("price_diff", 0)) > 0.02,
            action=arbitrage_action,
            description="Price difference > 2% - arbitrage opportunity",
        )

        await self.monitor.listen_tx_events()

6.5 Agent 链上事件处理循环

"""
onchain_agent_loop.py — Agent 链上事件处理主循环
"""

import asyncio
import json
from datetime import datetime


class MSGChainAgentLoop:
    """Agent 链上事件处理主循环"""

    def __init__(self):
        self.monitor = MsgChainEventMonitor()
        self.running = False
        self.task_queue = asyncio.Queue()

    async def event_handler(self, event_data: dict):
        """事件处理入口"""
        await self.task_queue.put({
            "timestamp": datetime.now().isoformat(),
            "data": event_data,
        })

    async def worker(self):
        """工作线程:处理事件队列"""
        while self.running:
            try:
                task = await asyncio.wait_for(
                    self.task_queue.get(), timeout=1.0
                )
                print(f"[Worker] Processing event from {task['timestamp']}")
                await self.process_event(task["data"])
            except asyncio.TimeoutError:
                continue

    async def process_event(self, data: dict):
        """处理单个事件"""
        # 解析事件类型
        events = data.get("events", {})
        for event_type, attributes in events.items():
            print(f"[Process] Event: {event_type}")
            # 根据事件类型分派到不同的 Agent 逻辑

    async def run(self):
        """运行事件循环"""
        self.running = True
        self.monitor.on("*", self.event_handler)

        # 并发运行监听和工作者
        await asyncio.gather(
            self.monitor.listen_tx_events(),
            self.worker(),
        )

    def stop(self):
        self.running = False
        self.monitor.stop()

7. 完整示例:DeFi 多 Agent 交易系统

7.1 系统概述

本节构建一个完整的 DeFi 多 Agent 交易系统,包含 4 个专业化 Agent:

Agent 角色 工具
Data Agent 数据采集 链上查询工具
Strategy Agent 策略分析 LLM + 数据分析
Trader Agent 交易执行 合约执行工具
Settlement Agent 结算确认 交易查询工具

7.2 完整代码

"""
complete_defi_system.py — MSG Chain 完整 DeFi 多 Agent 交易系统
"""

import json
import os
import asyncio
from typing import Dict, List, Optional, Any
from enum import Enum
from dataclasses import dataclass
from datetime import datetime

from dotenv import load_dotenv
load_dotenv()

# ── 导入 SDK ───────────────────────────────────────────
from msgchain_sdk_wrapper import MsgChainClient, MsgChainRestClient
from msgchain_langchain_tools import (
    MsgChainBalanceTool,
    MsgChainAllBalancesTool,
    MsgChainContractQueryTool,
    MsgChainTokenInfoTool,
    MsgChainBlockQueryTool,
    MsgChainValidatorTool,
    MsgChainLatestHeightTool,
    MsgChainTxQueryTool,
)
from msgchain_langchain_tx_tools import (
    MsgChainTransferTool,
    MsgChainContractExecuteTool,
)


# ── 数据模型 ───────────────────────────────────────────

class AgentRole(Enum):
    DATA = "data"
    STRATEGY = "strategy"
    TRADER = "trader"
    SETTLEMENT = "settlement"


@dataclass
class MarketData:
    """市场数据结构"""
    block_height: int
    token_prices: Dict[str, float]
    pool_liquidity: Dict[str, str]
    gas_prices: Dict[str, int]
    timestamp: str


@dataclass
class TradeSignal:
    """交易信号"""
    action: str  # buy, sell, swap
    token_in: str
    token_out: str
    amount: str
    min_out: str
    confidence: float
    reasoning: str


@dataclass
class ExecutionResult:
    """执行结果"""
    success: bool
    tx_hash: Optional[str]
    amount_in: str
    amount_out: str
    fee: str
    error: Optional[str]


# ── Agent 基类 ──────────────────────────────────────────

class BaseAgent:
    """Agent 基类"""

    def __init__(self, name: str, role: AgentRole):
        self.name = name
        self.role = role
        self.client = MsgChainClient()
        self.rest = MsgChainRestClient()
        self.log: List[str] = []

    def log_action(self, message: str):
        entry = f"[{datetime.now().isoformat()}] [{self.name}] {message}"
        self.log.append(entry)
        print(entry)

    def get_log(self) -> str:
        return "\n".join(self.log)


# ── Data Agent ──────────────────────────────────────────

class DataAgent(BaseAgent):
    """数据采集 Agent"""

    def __init__(self):
        super().__init__("DataAgent", AgentRole.DATA)
        self.tools = {
            "balance": MsgChainBalanceTool(),
            "all_balances": MsgChainAllBalancesTool(),
            "contract_query": MsgChainContractQueryTool(),
            "token_info": MsgChainTokenInfoTool(),
            "block": MsgChainBlockQueryTool(),
            "validators": MsgChainValidatorTool(),
            "latest_height": MsgChainLatestHeightTool(),
        }

    def collect_market_data(self, pools: List[str]) -> MarketData:
        """收集市场数据"""
        self.log_action("Collecting market data...")

        # 区块高度
        block_data = self.client.query_block()
        height = int(
            block_data.get("result", {}).get("block", {})
            .get("header", {}).get("height", 0)
        )

        # 池数据
        pool_data = {}
        for pool in pools:
            try:
                info = self.client.query_contract(
                    pool, {"info": {}}
                )
                pool_data[pool] = str(info)[:200]
            except Exception as e:
                pool_data[pool] = f"Error: {e}"

        # Gas 价格
        gas_prices = {
            "low": 5000,
            "average": 7500,
            "high": 10000,
        }

        self.log_action(f"Market data collected at height {height}")

        return MarketData(
            block_height=height,
            token_prices={},
            pool_liquidity=pool_data,
            gas_prices=gas_prices,
            timestamp=datetime.now().isoformat(),
        )

    def query_wallet(self, address: str) -> dict:
        """查询钱包信息"""
        self.log_action(f"Querying wallet: {address}")
        result = self.client.query_all_balances(address)
        return result


# ── Strategy Agent ──────────────────────────────────────

class StrategyAgent(BaseAgent):
    """策略分析 Agent"""

    def __init__(self):
        super().__init__("StrategyAgent", AgentRole.STRATEGY)
        self.pending_signals: List[TradeSignal] = []

    def analyze(self, market: MarketData, portfolio: dict) -> List[TradeSignal]:
        """分析市场并生成交易信号"""
        self.log_action("Analyzing market data...")

        signals = []

        # 模拟策略: 检查流动性池
        for pool, liquidity in market.pool_liquidity.items():
            if "uusdc" in str(liquidity).lower():
                signals.append(TradeSignal(
                    action="swap",
                    token_in="uusdc",
                    token_out="uaimgs",
                    amount="1000000",
                    min_out="0",
                    confidence=0.75,
                    reasoning="Standard swap based on liquidity availability",
                ))
                break

        self.log_action(f"Generated {len(signals)} trade signals")
        self.pending_signals = signals
        return signals

    def get_best_signal(self) -> Optional[TradeSignal]:
        """获取最佳信号"""
        if not self.pending_signals:
            return None
        return max(self.pending_signals, key=lambda s: s.confidence)


# ── Trader Agent ────────────────────────────────────────

class TraderAgent(BaseAgent):
    """交易执行 Agent"""

    def __init__(self):
        super().__init__("TraderAgent", AgentRole.TRADER)
        self.tx_tools = {
            "transfer": MsgChainTransferTool(),
            "contract_execute": MsgChainContractExecuteTool(),
        }

    def execute(self, signal: TradeSignal) -> ExecutionResult:
        """执行交易信号"""
        self.log_action(
            f"Executing {signal.action}: {signal.amount} {signal.token_in} -> {signal.token_out}"
        )

        try:
            # 构建合约执行交易
            execute_msg = json.dumps({
                "swap": {
                    "token_in": signal.token_in,
                    "token_out": signal.token_out,
                    "amount": signal.amount,
                    "min_out": signal.min_out,
                    "recipient": os.getenv("AGENT_ADDRESS", ""),
                }
            })

            # 模拟执行
            result = ExecutionResult(
                success=True,
                tx_hash="TX_SIMULATED_" + hashlib.sha256(
                    execute_msg.encode()
                ).hexdigest()[:16],
                amount_in=signal.amount,
                amount_out="950000",  # 模拟 5% 滑点
                fee="5000",
                error=None,
            )

            self.log_action(f"Trade executed: {result.tx_hash}")
            return result

        except Exception as e:
            self.log_action(f"Trade failed: {e}")
            return ExecutionResult(
                success=False,
                tx_hash=None,
                amount_in=signal.amount,
                amount_out="0",
                fee="5000",
                error=str(e),
            )


# ── Settlement Agent ───────────────────────────────────

class SettlementAgent(BaseAgent):
    """结算确认 Agent"""

    def __init__(self):
        super().__init__("SettlementAgent", AgentRole.SETTLEMENT)
        self.tx_query = MsgChainTxQueryTool()

    def confirm(self, result: ExecutionResult) -> bool:
        """确认交易结算"""
        self.log_action(f"Confirming settlement for tx: {result.tx_hash}")

        if not result.success:
            self.log_action("Transaction failed, no confirmation needed")
            return False

        if not result.tx_hash:
            self.log_action("No transaction hash to confirm")
            return False

        try:
            # 查询交易状态
            tx_info = self.client.get_tx(result.tx_hash)
            code = tx_info.get("result", {}).get("tx_response", {}).get("code", None)

            if code == 0:
                self.log_action("Transaction confirmed successfully")
                return True
            else:
                self.log_action(f"Transaction failed with code: {code}")
                return False

        except Exception as e:
            self.log_action(f"Confirmation error: {e}")
            return False

    def generate_report(self, results: List[ExecutionResult]) -> str:
        """生成结算报告"""
        successful = sum(1 for r in results if r.success)
        failed = sum(1 for r in results if not r.success)

        report = {
            "timestamp": datetime.now().isoformat(),
            "total_trades": len(results),
            "successful": successful,
            "failed": failed,
            "details": [
                {
                    "tx_hash": r.tx_hash,
                    "status": "success" if r.success else "failed",
                    "amount_in": r.amount_in,
                    "amount_out": r.amount_out,
                    "fee": r.fee,
                    "error": r.error,
                }
                for r in results
            ],
        }

        self.log_action(f"Report generated: {successful}/{len(results)} successful")
        return json.dumps(report, indent=2, ensure_ascii=False)


# ── 协调器 ──────────────────────────────────────────────

class DeFiCoordinator:
    """DeFi 多 Agent 协调器"""

    def __init__(self):
        self.data_agent = DataAgent()
        self.strategy_agent = StrategyAgent()
        self.trader_agent = TraderAgent()
        self.settlement_agent = SettlementAgent()
        self.execution_history: List[ExecutionResult] = []

    def run_pipeline(self, pools: List[str], address: str):
        """运行完整 DeFi 交易流水线"""
        print("=" * 60)
        print("MSG Chain DeFi Multi-Agent Trading Pipeline")
        print("=" * 60)

        # Phase 1: 数据采集
        print("\n[Phase 1] Data Collection")
        market = self.data_agent.collect_market_data(pools)
        portfolio = self.data_agent.query_wallet(address)
        print(f"  Block Height: {market.block_height}")
        print(f"  Pools Queried: {len(pools)}")

        # Phase 2: 策略分析
        print("\n[Phase 2] Strategy Analysis")
        signals = self.strategy_agent.analyze(market, portfolio)
        print(f"  Signals Generated: {len(signals)}")

        if not signals:
            print("  No trading signals - pipeline complete")
            return

        # Phase 3: 交易执行
        print("\n[Phase 3] Trade Execution")
        for signal in signals:
            result = self.trader_agent.execute(signal)
            self.execution_history.append(result)
            print(f"  {'✓' if result.success else '✗'} {signal.action}: {result.tx_hash or 'FAILED'}")

        # Phase 4: 结算确认
        print("\n[Phase 4] Settlement Confirmation")
        for result in self.execution_history:
            confirmed = self.settlement_agent.confirm(result)
            print(f"  {'✓' if confirmed else '✗'} {result.tx_hash}: {'Confirmed' if confirmed else 'Failed'}")

        # Phase 5: 报告生成
        print("\n[Phase 5] Report Generation")
        report = self.settlement_agent.generate_report(self.execution_history)
        print(f"\n{report}")

        # 输出日志
        print("\n=== Agent Logs ===")
        for agent in [self.data_agent, self.strategy_agent,
                      self.trader_agent, self.settlement_agent]:
            print(f"\n--- {agent.name} Log ---")
            print(agent.get_log())

    def get_summary(self) -> dict:
        """获取执行摘要"""
        return {
            "total_trades": len(self.execution_history),
            "successful": sum(1 for r in self.execution_history if r.success),
            "total_amount_in": sum(int(r.amount_in) for r in self.execution_history),
            "total_amount_out": sum(int(r.amount_out) for r in self.execution_history if r.success),
        }


# ── 入口 ───────────────────────────────────────────────

if __name__ == "__main__":
    import hashlib

    coordinator = DeFiCoordinator()

    coordinator.run_pipeline(
        pools=[
            "msg1dex_pool_uusdc_uaimgs",
            "msg1dex_pool_uusdc_uusdt",
            "msg1dex_pool_uaimgs_uusdt",
        ],
        address=os.getenv("AGENT_ADDRESS", "msg1agent_address"),
    )

    summary = coordinator.get_summary()
    print(f"\n=== Final Summary ===")
    print(json.dumps(summary, indent=2))

7.3 使用 CrewAI 实现的等效系统

"""
complete_defi_crew.py — 使用 CrewAI 实现的 DeFi 交易系统
"""

from crewai import Agent, Task, Crew, Process

from crewai_msgchain_agents import (
    create_data_analyst_agent,
    create_defi_trader_agent,
)


def create_defi_trading_crew(
    address: str,
    pools: list,
    model: str = "gpt-4o",
) -> Crew:
    """创建 DeFi 交易 Crew"""

    data_agent = create_data_analyst_agent(model=model)
    trader = create_defi_trader_agent(model=model)

    tasks = [
        Task(
            description=(
                f"Query MSG Chain market data for address {address}. "
                f"Analyze these pools: {pools}. "
                "Get current prices, liquidity, and trading volumes."
            ),
            expected_output="Market data report with prices and liquidity.",
            agent=data_agent,
        ),
        Task(
            description=(
                "Based on the market data, identify and rank "
                "trading opportunities. Consider arbitrage, "
                "yield farming, and optimal swap routes."
            ),
            expected_output="Ranked list of trading opportunities.",
            agent=data_agent,
        ),
        Task(
            description=(
                "Execute the best trading opportunity. "
                "Simulate the swap and report expected outcome. "
                "Consider gas costs and slippage."
            ),
            expected_output="Trade execution result with simulated tx hash.",
            agent=trader,
        ),
        Task(
            description=(
                "Confirm the settlement of executed trades. "
                "Query transaction status and generate a report "
                "summarizing all operations."
            ),
            expected_output="Settlement confirmation and final report.",
            agent=trader,
        ),
    ]

    return Crew(
        agents=[data_agent, trader],
        tasks=tasks,
        process=Process.sequential,
        verbose=True,
    )


if __name__ == "__main__":
    crew = create_defi_trading_crew(
        address="msg1agent_address",
        pools=["msg1dex_pool_1", "msg1dex_pool_2"],
    )
    result = crew.kickoff()
    print(f"\nCrew Resul[未公开路径]}")

7.4 使用 LangGraph 实现的等效系统

"""
complete_defi_graph.py — 使用 LangGraph 实现的 DeFi 交易状态机
"""

from typing import TypedDict, Dict, List, Optional, Any
from langgraph.graph import StateGraph, END


class DeFiTradeState(TypedDict):
    address: str
    pools: List[str]
    market_data: Optional[Dict]
    portfolio: Optional[Dict]
    signals: List[Dict]
    execution_results: List[Dict]
    confirmed: List[bool]
    report: Optional[str]
    error: Optional[str]
    phase: str


# 节点实现
def collect_data(state: DeFiTradeState) -> DeFiTradeState:
    print("[Phase 1] Collecting market data...")
    from msgchain_sdk_wrapper import MsgChainClient
    client = MsgChainClient()
    block = client.query_block()
    return {**state, "market_data": {"block": block}, "phase": "collect"}


def analyze_and_generate_signals(state: DeFiTradeState) -> DeFiTradeState:
    print("[Phase 2] Analyzing and generating signals...")
    signals = [{
        "action": "swap",
        "token_in": "uusdc",
        "token_out": "uaimgs",
        "amount": "1000000",
        "confidence": 0.8,
    }]
    return {**state, "signals": signals, "phase": "analyze"}


def execute_trades(state: DeFiTradeState) -> DeFiTradeState:
    print("[Phase 3] Executing trades...")
    results = []
    for signal in state["signals"]:
        results.append({
            "success": True,
            "tx_hash": f"tx_{hash(str(signal))}",
            **signal,
        })
    return {**state, "execution_results": results, "phase": "execute"}


def confirm_settlements(state: DeFiTradeState) -> DeFiTradeState:
    print("[Phase 4] Confirming settlements...")
    confirmed = [r["success"] for r in state["execution_results"]]
    return {**state, "confirmed": confirmed, "phase": "confirm"}


def generate_report(state: DeFiTradeState) -> DeFiTradeState:
    print("[Phase 5] Generating report...")
    import json, datetime
    report = json.dumps({
        "timestamp": datetime.datetime.now().isoformat(),
        "total": len(state["execution_results"]),
        "successful": sum(state["confirmed"]),
        "results": state["execution_results"],
    }, indent=2)
    return {**state, "report": report, "phase": "done"}


def check_errors(state: DeFiTradeState) -> str:
    if state.get("error"):
        return "error"
    return "continue"


# 构建图
workflow = StateGraph(DeFiTradeState)

workflow.add_node("collect_data", collect_data)
workflow.add_node("analyze", analyze_and_generate_signals)
workflow.add_node("execute", execute_trades)
workflow.add_node("confirm", confirm_settlements)
workflow.add_node("report", generate_report)

workflow.set_entry_point("collect_data")
workflow.add_edge("collect_data", "analyze")
workflow.add_edge("analyze", "execute")
workflow.add_edge("execute", "confirm")
workflow.add_edge("confirm", "report")
workflow.add_edge("report", END)


if __name__ == "__main__":
    app = workflow.compile()
    result = app.invoke({
        "address": "msg1agent_address",
        "pools": ["pool1", "pool2"],
        "market_data": None,
        "portfolio": None,
        "signals": [],
        "execution_results": [],
        "confirmed": [],
        "report": None,
        "error": None,
        "phase": "start",
    })
    print(f"\nFinal Repor[未公开路径]'report', 'N/A')}")

8. 部署指南

8.1 Docker 部署

# Dockerfile — MSG Chain AI Agent 运行环境
FROM python:3.11-slim

WORKDIR /app

# 安装系统依赖
RUN apt-get update && apt-get install -y \\n    build-essential \\n    curl \\n    git \\n    && rm -rf /var/lib/apt/lists/*

# 安装 Python 依赖
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt

# 复制应用代码
COPY . .

# 环境变量
ENV MSG_CHAIN_RPC=https://rpc.msgchain.org
ENV MSG_CHAIN_CHAIN_ID=msg-chain-1
ENV PYTHONUNBUFFERED=1

# 启动命令
CMD ["python", "main.py"]
# requirements.txt
msgchain-sdk>=0.1.0
langchain>=0.2.0
langchain-community>=0.2.0
langchain-openai>=0.1.0
langgraph>=0.1.0
crewai>=0.30.0
pyautogen>=0.2.0
openai>=1.0.0
websockets>=12.0
requests>=2.31.0
python-dotenv>=1.0.0
pydantic>=2.0.0

8.2 Docker Compose 编排

# docker-compose.yml
version: "3.8"

services:
  msgchain-agent:
    build: .
    container_name: msgchain-ai-agent
    env_file:
      - .env
    volumes:
      - ./data:/app/data
      - ./logs:/app/logs
    restart: unless-stopped
    networks:
      - msgchain-net
    healthcheck:
      test: ["CMD", "python", "-c", "import msgchain_sdk; print('ok')"]
      interval: 30s
      timeout: 10s
      retries: 3

  msgchain-event-monitor:
    build: .
    container_name: msgchain-event-monitor
    command: python -m onchain_events
    env_file:
      - .env
    volumes:
      - ./logs:/app/logs
    restart: unless-stopped
    depends_on:
      - msgchain-agent
    networks:
      - msgchain-net

networks:
  msgchain-net:
    driver: bridge

8.3 环境配置示例

# .env 完整示例

# MSG Chain 节点配置
MSG_CHAIN_RPC=https://rpc.msgchain.org
MSG_CHAIN_REST=https://rest.msgchain.org
MSG_CHAIN_CHAIN_ID=msg-chain-1
MSG_CHAIN_WS=wss://rpc.msgchain.org/websocket

# Agent 密钥(Dilithium-5)
AGENT_PRIVATE_KEY=0123456789abcdef0123456789abcdef...
AGENT_ADDRESS=msg1agent_wallet_address

# LLM API 密钥(至少配置一个)
OPENAI_API_KEY=sk-your-openai-api-key
# ANTHROPIC_API_KEY=sk-ant-your-anthropic-api-key

# 日志级别
LOG_LEVEL=INFO

# 监控配置
ENABLE_EVENT_MONITORING=true
MONITOR_CONTRACTS=msg1contract1,msg1contract2

# Agent 参数
AGENT_MAX_ITERATIONS=15
AGENT_TEMPERATURE=0.1
DEFAULT_MODEL=gpt-4o

8.4 Kubernetes 部署

# k8s-deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
  name: msgchain-agent
  labels:
    app: msgchain-agent
spec:
  replicas: 2
  selector:
    matchLabels:
      app: msgchain-agent
  template:
    metadata:
      labels:
        app: msgchain-agent
    spec:
      containers:
      - name: agent
        image: msgchain/ai-agent:latest
        ports:
        - containerPort: 8080
        envFrom:
        - secretRef:
            name: msgchain-agent-secrets
        env:
        - name: MSG_CHAIN_RPC
          value: "https://rpc.msgchain.org"
        - name: MSG_CHAIN_CHAIN_ID
          value: "msg-chain-1"
        resources:
          requests:
            memory: "512Mi"
            cpu: "250m"
          limits:
            memory: "1Gi"
            cpu: "500m"
        livenessProbe:
          httpGet:
            path: /health
            port: 8080
          initialDelaySeconds: 30
          periodSeconds: 10
---
apiVersion: v1
kind: Secret
metadata:
  name: msgchain-agent-secrets
type: Opaque
data:
  AGENT_PRIVATE_KEY: <base64-encoded-key>
  OPENAI_API_KEY: <base64-encoded-key>

8.5 安全性注意事项

## 安全最佳实践

1. **密钥管理**
   - 永远不要将私钥硬编码在代码中
   - 使用环境变量或密钥管理服务(AWS KMS / Azure Key Vault / HashiCorp Vault)
   - 定期轮换 Agent 私钥
   - Dilithium-5 私钥需安全存储(HSM 推荐)

2. **API 密钥**
   - 使用 .gitignore 排除 .env 文件
   - 为每个环境使用不同的 API Key
   - 限制 LLM API 调用频率和预算

3. **交易安全**
   - 始终模拟交易后再广播
   - 设置合理的 Gas 限制以防止资金耗尽
   - 实现交易限额(例如单笔最大转账量)
   - 使用两步确认机制

4. **RPC 安全**
   - 使用可信的 RPC 端点
   - 实施请求速率限制
   - 监控异常 RPC 调用模式

5. **Agent 权限**
   - 最小权限原则:Agent 只应有执行所需的最少权限
   - 使用专用 Agent 账户,不与主账户共用
   - 定期审计 Agent 交易历史

6. **日志与监控**
   - 记录所有 Agent 操作
   - 设置告警:异常交易、大额转账、频繁失败
   - 使用集中式日志系统(ELK / Grafana Loki)

8.6 监控与告警配置

"""
monitoring.py — MSG Chain Agent 监控系统
"""

import json
import time
import logging
from datetime import datetime, timedelta
from typing import Dict, List, Optional


class AgentMonitor:
    """Agent 运行监控器"""

    def __init__(self):
        self.metrics = {
            "total_queries": 0,
            "successful_txs": 0,
            "failed_txs": 0,
            "total_fees_paid": 0,
            "start_time": datetime.now().isoformat(),
            "last_activity": None,
        }
        self.alerts: List[Dict] = []
        self.logger = logging.getLogger("msgchain.monitor")

    def record_query(self):
        self.metrics["total_queries"] += 1
        self.metrics["last_activity"] = datetime.now().isoformat()

    def record_tx(self, success: bool, fee: str = "0"):
        if success:
            self.metrics["successful_txs"] += 1
        else:
            self.metrics["failed_txs"] += 1
        self.metrics["total_fees_paid"] += int(fee)
        self.metrics["last_activity"] = datetime.now().isoformat()

    def check_health(self) -> Dict:
        """检查 Agent 健康状态"""
        now = datetime.now()
        last = datetime.fromisoformat(
            self.metrics["last_activity"] or self.metrics["start_time"]
        )

        status = "healthy"
        alerts = []

        # 检查是否长时间未活动
        if now - last > timedelta(minutes=30):
            status = "warning"
            alerts.append("No activity for 30+ minutes")

        # 检查高失败率
        total = self.metrics["successful_txs"] + self.metrics["failed_txs"]
        if total > 10 and self.metrics["failed_txs"] / total > 0.3:
            status = "degraded"
            alerts.append(f"High failure rate: {self.metrics['failed_txs']}/{total}")

        return {
            "status": status,
            "metrics": self.metrics,
            "alerts": alerts,
            "timestamp": now.isoformat(),
        }

    def raise_alert(self, severity: str, message: str, data: Optional[Dict] = None):
        """触发告警"""
        alert = {
            "severity": severity,
            "message": message,
            "data": data,
            "timestamp": datetime.now().isoformat(),
        }
        self.alerts.append(alert)
        self.logger.warning(f"[{severity}] {message}")
        return alert

    def get_stats(self) -> str:
        """获取统计摘要"""
        return json.dumps(self.metrics, indent=2, ensure_ascii=False)

8.7 快速启动脚本

#!/bin/bash
# setup_msgchain_agent.sh — MSG Chain AI Agent 快速部署脚本

set -e

echo "=== MSG Chain AI Agent Setup ==="

# 1. 创建项目目录
PROJECT_DIR="msgchain-agent-$(date +%Y%m%d)"
mkdir -p $PROJECT_DIR/{src,config,logs,data}
cd $PROJECT_DIR

# 2. 创建虚拟环境
python3 -m venv venv
source venv/bin/activate

# 3. 创建 requirements.txt
cat > requirements.txt << 'REQ'
msgchain-sdk>=0.1.0
langchain>=0.2.0
langgraph>=0.1.0
crewai>=0.30.0
pyautogen>=0.2.0
openai>=1.0.0
websockets>=12.0
requests>=2.31.0
python-dotenv>=1.0.0
pydantic>=2.0.0
REQ

# 4. 安装依赖
pip install -r requirements.txt

# 5. 创建 .env 模板
cat > .env << 'ENV'
MSG_CHAIN_RPC=https://rpc.msgchain.org
MSG_CHAIN_CHAIN_ID=msg-chain-1
AGENT_PRIVATE_KEY=your_key_here
AGENT_ADDRESS=msg1your_address_here
OPENAI_API_KEY=sk-your_key_here
LOG_LEVEL=INFO
ENV

# 6. 创建主入口文件
cat > main.py << 'MAIN'
"""MSG Chain AI Agent 主入口"""
import os
from dotenv import load_dotenv

load_dotenv()


def main():
    print("MSG Chain AI Agent starting...")
    print(f"Chain ID: {os.getenv('MSG_CHAIN_CHAIN_ID')}")
    print(f"RPC: {os.getenv('MSG_CHAIN_RPC')}")
    print("Agent initialized successfully")


if __name__ == "__main__":
    main()
MAIN

echo ""
echo "=== Setup Complete ==="
echo "Project: $PWD"
echo "To start: source venv/bin/activate && python main.py"
echo "Edit .env to configure your keys"

8.8 CI/CD 示例

# .github/workflows/deploy-agent.yml
name: Deploy MSG Chain Agent

on:
  push:
    branches: [main]
  pull_request:
    branches: [main]

jobs:
  test:
    runs-on: ubuntu-latest
    steps:
    - uses: actions/checkout@v4
    - name: Set up Python
      uses: actions/setup-python@v5
      with:
        python-version: "3.11"
    - name: Install dependencies
      run: |
        pip install -r requirements.txt
    - name: Lint
      run: |
        pip install flake8
        flake8 src/ --max-line-length=100
    - name: Test
      run: |
        pip install pytest
        pytest tests/

  deploy:
    needs: test
    runs-on: ubuntu-latest
    if: github.ref == 'refs/heads/main'
    steps:
    - uses: actions/checkout@v4
    - name: Build Docker image
      run: |
        docker build -t msgchain/ai-agent:${{ github.sha }} .
    - name: Deploy to Kubernetes
      env:
        KUBE_CONFIG: ${{ secrets.KUBE_CONFIG }}
      run: |
        echo "$KUBE_CONFIG" | base64 --decode > kubeconfig
        kubectl --kubeconfig=kubeconfig apply -f k8s/

附录

A. 常见问题

Q: Agent 如何获取 Dilithium-5 私钥?

A: 私钥应在安全环境中生成。可以使用 msgchain-sdk 的密钥生成工具:

python -c "from msgchain_sdk.dilithium import generate_keypair; \\n  pk, sk = generate_keypair(); \\n  print(f'Private: {sk.hex()}')"

Q: 不同框架如何选择?

Q: 如何测试 Agent 而不花费实际 Gas?

A: 使用测试网。MSG Chain 提供测试网 faucet:

curl -X POST https://faucet.testnet.msgchain.org/claim \\n  -H "Content-Type: application/json" \\n  -d '{"address": "msg1test_address"}'

B. MSG Chain 测试网信息

参数 值
Chain ID msg-chain-test-1
RPC https://rpc.testnet.msgchain.org
REST https://rest.testnet.msgchain.org
Faucet https://faucet.testnet.msgchain.org
Explorer https://explorer.testnet.msgchain.org
Denom uaimgs

C. 参考资源