dApp Docs/AI Agent 引导启动流程指南
Development reference. Not independently verified for production.

AI Agent 引导启动流程指南

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

版本:v1.0
适用链:msg-chain-1 | 编写语言:Python 3.10+ / TypeScript 5.x


目录

  1. 概述
  2. 第零步:环境准备与工具链安装
  3. 第一步:创建去中心化身份(DID)
  4. 第二步:为 Agent 充值
  5. 第三步:注册 AI Agent
  6. 第四步:设定宪章与策略引擎
  7. 第五步:连接 Agent API 网关
  8. 第六步:订阅链上事件
  9. 第七步:集成 MPC 钱包与 AIPAY 支付
  10. 第八步:Agent 间通信(A2A)
  11. 第九步:监控与日志
  12. 第十步:一键启动脚本
  13. 附录

1. 概述

本文档提供在 MSG Chain 上从零启动一个 AI Agent 的完整引导流程。无论你是构建自动化交易 Agent、客服 Agent、数据索引 Agent 还是自治治理 Agent,以下 10 个步骤将覆盖从环境搭建到生产上线的全部环节。

1.1 什么是 MSG Chain AI Agent?

MSG Chain 上的 AI Agent 是一个拥有 去中心化身份(DID)、注册在链上的智能合约账户,并受 AI Agent 宪章 约束的自主程序。Agent 可以:

1.2 架构概览

┌─────────────────────────────────────────────────────────────┐
│                    AI Agent 运行时架构                        │
├─────────────────────────────────────────────────────────────┤
│                                                             │
│  ┌──────────────┐    ┌──────────────┐    ┌──────────────┐  │
│  │  Agent API   │    │ 策略引擎     │    │ 事件处理器   │  │
│  │  Gateway     │◄──►│ (Constitution)│◄──►│ (Subscriber) │  │
│  └──────┬───────┘    └──────┬───────┘    └──────┬───────┘  │
│         │                   │                   │           │
│  ┌──────▼───────────────────▼───────────────────▼───────┐  │
│  │              Agent 核心层 (Core Runtime)              │  │
│  │  任务调度 | 状态管理 | 记忆存储 | 工具调用            │  │
│  └────────────────────────┬──────────────────────────────┘  │
│                           │                                  │
│  ┌────────────────────────▼──────────────────────────────┐  │
│  │               MSG Chain 交互层                         │  │
│  │  DID | Registry | Constitution | MPC | AIPAY | A2A    │  │
│  └─────────────────────────────────────────────────────────┘  │
│                                                             │
│  ┌─────────────────────────────────────────────────────────┐  │
│  │                  MSG Chain 主网                          │  │
│  │        msg-chain-1 | CosmWasm | IBC | CosmTrez         │  │
│  └─────────────────────────────────────────────────────────┘  │
│                                                             │
└─────────────────────────────────────────────────────────────┘

1.3 前置条件

1.4 核心合约地址

合约 地址
DID Registry msg14hj2tavq8fpesdwxxcu44rty3hh90vhujrvcmstl4zr3txmfvw9s4hmal
Agent Registry msg1qypqxpq9kcrn2c9afea5lq35ef37c5x7jqylz3
AI Agent Constitution msg1s9xu5h2nkl6dcncxr48e7q07gqy4n6p0x7vq9k
AIPAY Module msg1aipayxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx
MPC Wallet Factory msg1mpcxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx

2. 第零步:环境准备与工具链安装

2.1 安装 MSG Chain 节点与 CLI

方法 A:预编译二进制(推荐)

# 下载最新版本(请从官方 Release 页面获取最新版本号)
VERSION="v1.0.0"
wget https://github.com/msgchain/msg-chain/releases/download/${VERSION}/msg-chain-linux-amd64.tar.gz
tar -xzf msg-chain-linux-amd64.tar.gz
sudo mv msgd /usr/local/bin/
msgd version

方法 B:从源码编译

# 前置依赖
sudo apt-get update && sudo apt-get install -y git golang-go make gcc libssl-dev

# 克隆仓库
git clone https://github.com/msgchain/msg-chain.git
cd msg-chain
git checkout v1.0.0

# 编译
make install

# 验证
msgd version

2.2 配置节点连接

# 初始化配置(不初始化 validator)
msgd config chain-id msg-chain-1
msgd config node https://rpc.msgchain.org:443
msgd config output json
msgd config keyring-backend os

# 测试连接
msgd status

预期输出:显示节点同步状态、最新区块高度等信息。如返回 "sync_info": {"catching_up": false} 则表示节点已同步。

2.3 安装 Python 依赖

# Python 3.10+ 推荐
python3 -m venv agent-env
source agent-env/bin/activate

# 安装核心依赖
pip install msg-chain-sdk==1.0.0
pip install cosmwasm-cli==1.0.0
pip install httpx==0.27.0
pip install pydantic==2.5.0
pip install asyncio==3.4.3
pip install websockets==12.0
pip install cryptography==41.0.0

# 验证安装
python3 -c "from msg_chain_sdk import Client; print('SDK OK')"

2.4 安装 TypeScript 依赖

# 初始化项目
mkdir my-agent && cd my-agent
npm init -y

# 安装核心依赖
npm install @msgchain/sdk@^1.0.0
npm install @msgchain/cosmwasm@^1.0.0
npm install axios@^1.6.0
npm install ws@^8.16.0
npm install ethers@^6.9.0
npm install typescript@^5.3.0 ts-node@^10.9.0

# TypeScript 配置
cat > tsconfig.json << 'EOF'
{
  "compilerOptions": {
    "target": "ES2022",
    "module": "commonjs",
    "lib": ["ES2022"],
    "outDir": "./dist",
    "rootDir": "./src",
    "strict": true,
    "esModuleInterop": true,
    "skipLibCheck": true,
    "forceConsistentCasingInFileNames": true,
    "resolveJsonModule": true,
    "declaration": true
  },
  "include": ["src/**/*"],
  "exclude": ["node_modules", "dist"]
}
EOF

# 创建源码目录
mkdir -p src

2.5 安装 CosmWasm CLI 工具

# 安装 Rust(如果需要从源码编译)
curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh
source "$HOME/.cargo/env"

# 安装 cosmwasm-cli
cargo install cosmwasm-cli

# 或者使用预编译版本
wget https://github.com/CosmWasm/cosmwasm-cli/releases/download/v1.0.0/cosmwasm-cli-linux-amd64.tar.gz
tar -xzf cosmwasm-cli-linux-amd64.tar.gz
sudo mv cosmwasm-cli /usr/local/bin/

# 验证
cosmwasm-cli --version

2.6 创建钱包账户

# 创建新钱包(请安全保存助记词!)
msgd keys add agent_owner

# 输出示例:
# - name: agent_owner
#   type: local
#   address: msg1abcdef...
#   pubkey: '{"@type":"/cosmos.crypto.secp256k1.PubKey","key":"..."}'
#   mnemonic: "abandon abandon abandon ..."

# 查看所有钱包
msgd keys list

# 查看钱包地址
AGENT_OWNER=$(msgd keys show agent_owner -a)
echo "Owner Address: $AGENT_OWNER"

安全提醒:助记词是恢复钱包的唯一方式。请离线存储在安全位置。永远不要将助记词提交到代码仓库或通过网络传输。

2.7 环境变量配置

创建 .env 文件集中管理配置:

cat > .env << 'EOF'
# MSG Chain 配置
MSG_CHAIN_ID=msg-chain-1
MSG_RPC_URL=https://rpc.msgchain.org:443
MSG_REST_URL=https://rest.msgchain.org:443
MSG_WEBSOCKET_URL=wss://rpc.msgchain.org:443/websocket

# 合约地址
DID_REGISTRY_ADDRESS=msg14hj2tavq8fpesdwxxcu44rty3hh90vhujrvcmstl4zr3txmfvw9s4hmal
AGENT_REGISTRY_ADDRESS=msg1qypqxpq9kcrn2c9afea5lq35ef37c5x7jqylz3
CONSTITUTION_ADDRESS=msg1s9xu5h2nkl6dcncxr48e7q07gqy4n6p0x7vq9k
AIPAY_MODULE_ADDRESS=msg1aipayxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx
MPC_FACTORY_ADDRESS=msg1mpcxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx

# Agent 配置
AGENT_OWNER=$(msgd keys show agent_owner -a 2>/dev/null)
AGENT_MNEMONIC=
AGENT_NAME=my_first_agent
AGENT_DESCRIPTION="My first AI Agent on MSG Chain"
AGENT_API_ENDPOINT=http://localhost:8080

# Dilithium-5 密钥路径
DILITHIUM_PRIV_KEY_PATH=./keys/dilithium_priv.bin
DILITHIUM_PUB_KEY_PATH=./keys/dilithium_pub.bin

# 日志配置
LOG_LEVEL=INFO
LOG_FORMAT=json
EOF
# Python: 加载环境变量
from dotenv import load_dotenv
import os

load_dotenv()

MSG_RPC_URL = os.getenv("MSG_RPC_URL")
DID_REGISTRY = os.getenv("DID_REGISTRY_ADDRESS")
AGENT_REGISTRY = os.getenv("AGENT_REGISTRY_ADDRESS")
AGENT_OWNER = os.getenv("AGENT_OWNER")
// TypeScript: 加载环境变量
import dotenv from 'dotenv';
dotenv.config();

const MSG_RPC_URL = process.env.MSG_RPC_URL!;
const DID_REGISTRY = process.env.DID_REGISTRY_ADDRESS!;
const AGENT_REGISTRY = process.env.AGENT_REGISTRY_ADDRESS!;
const AGENT_OWNER = process.env.AGENT_OWNER!;

2.8 验证环境

# Python: 验证环境连接
import asyncio
from msg_chain_sdk import Client

async def verify_environment():
    client = Client(
        rpc_url=os.getenv("MSG_RPC_URL"),
        chain_id=os.getenv("MSG_CHAIN_ID")
    )

    # 检查节点状态
    status = await client.get_status()
    print(f"✓ 节点连接成功")
    print(f"  最新区块: {status.sync_info.latest_block_height}")
    print(f"  同步状态: {'同步中' if status.sync_info.catching_up else '已同步'}")

    # 检查账户余额
    owner = os.getenv("AGENT_OWNER")
    balance = await client.get_balance(owner, "umsg")
    print(f"✓ 账户余额: {balance} umsg")

    # 检查合约是否存在
    code_id = await client.get_contract_code_id(DID_REGISTRY)
    print(f"✓ DID Registry 合约存在 (code_id: {code_id})")

    await client.close()
    return True

asyncio.run(verify_environment())
// TypeScript: 验证环境连接
import { MsgChainClient } from '@msgchain/sdk';

async function verifyEnvironment(): Promise<boolean> {
  const client = new MsgChainClient({
    rpcUrl: process.env.MSG_RPC_URL!,
    chainId: process.env.MSG_CHAIN_ID!,
  });

  try {
    // 检查节点状态
    const status = await client.getStatus();
    console.log(`✓ 节点连接成功`);
    console.log(`  最新区块: ${status.syncInfo.latestBlockHeight}`);
    console.log(`  同步状态: ${status.syncInfo.catchingUp ? '同步中' : '已同步'}`);

    // 检查账户余额
    const balance = await client.getBalance(process.env.AGENT_OWNER!, 'umsg');
    console.log(`✓ 账户余额: ${balance} umsg`);

    // 检查合约是否存在
    const codeId = await client.getContractCodeId(process.env.DID_REGISTRY_ADDRESS!);
    console.log(`✓ DID Registry 合约存在 (code_id: ${codeId})`);

    return true;
  } finally {
    await client.close();
  }
}

verifyEnvironment().catch(console.error);

3. 第一步:创建去中心化身份(DID)

每个 AI Agent 需要一个去中心化身份(Decentralized Identifier, DID)才能在 MSG Chain 上被唯一识别和信任。

3.1 DID 文档结构

{
  "@context": [
    "https://www.w3.org/ns/did/v1",
    "https://msgchain.org/ns/did-config/v1"
  ],
  "id": "did:msg:msg1agentaddress...",
  "alsoKnownAs": ["my_first_agent"],
  "controller": ["did:msg:msg1owneraddress..."],
  "verificationMethod": [
    {
      "id": "did:msg:msg1agentaddress...#keys-1",
      "type": "EcdsaSecp256k1RecoveryMethod2020",
      "controller": "did:msg:msg1agentaddress...",
      "blockchainAccountId": "msg1agentaddress..."
    },
    {
      "id": "did:msg:msg1agentaddress...#dilithium-5-keys-1",
      "type": "Dilithium5VerificationKey2026",
      "controller": "did:msg:msg1agentaddress...",
      "publicKeyMultibase": "zDilithium5PubKeyBase64..."
    }
  ],
  "authentication": [
    "did:msg:msg1agentaddress...#keys-1",
    "did:msg:msg1agentaddress...#dilithium-5-keys-1"
  ],
  "assertionMethod": [
    "did:msg:msg1agentaddress...#keys-1"
  ],
  "keyAgreement": [
    {
      "id": "did:msg:msg1agentaddress...#x25519-keys-1",
      "type": "X25519KeyAgreementKey2019",
      "controller": "did:msg:msg1agentaddress...",
      "publicKeyMultibase": "zX25519PubKeyBase64..."
    }
  ],
  "service": [
    {
      "id": "did:msg:msg1agentaddress...#agent-endpoint",
      "type": "AIAgentEndpoint",
      "serviceEndpoint": "http://localhost:8080"
    },
    {
      "id": "did:msg:msg1agentaddress...#a2a-endpoint",
      "type": "A2AEndpoint",
      "serviceEndpoint": "http://localhost:8080/a2a"
    }
  ],
  "created": "2026-07-06T00:00:00Z",
  "updated": "2026-07-06T00:00:00Z"
}

3.2 生成 Dilithium-5 密钥对

# 使用 msgd CLI 生成后量子密钥
msgd keys add-dilithium --output-file ./keys/dilithium_keys.json

# 或使用独立的 Dilithium 工具
mkdir -p keys
python3 -c "
from msg_chain_sdk.crypto import Dilithium5

# 生成密钥对
dk = Dilithium5.generate()
dk.save_private_key('keys/dilithium_priv.bin')
dk.save_public_key('keys/dilithium_pub.bin')

print(f'私钥已保存: keys/dilithium_priv.bin')
print(f'公钥已保存: keys/dilithium_pub.bin')
print(f'公钥 (hex): {dk.public_key.hex()[:64]}...')
"
# Python: 生成 Dilithium-5 密钥对
from msg_chain_sdk.crypto import Dilithium5
from pathlib import Path

def generate_dilithium_keys(key_dir: str = "./keys"):
    Path(key_dir).mkdir(parents=True, exist_ok=True)

    # 生成密钥对
    dk = Dilithium5.generate()

    # 保存
    priv_path = Path(key_dir) / "dilithium_priv.bin"
    pub_path = Path(key_dir) / "dilithium_pub.bin"
    dk.save_private_key(str(priv_path))
    dk.save_public_key(str(pub_path))

    print(f"✓ Dilithium-5 密钥对已生成")
    print(f"  私钥: {priv_path}")
    print(f"  公钥: {pub_path}")
    print(f"  公钥 (hex): {dk.public_key.hex()[:64]}...")

    return dk

dk = generate_dilithium_keys()
// TypeScript: 生成 Dilithium-5 密钥对
import { Dilithium5 } from '@msgchain/sdk';
import * as fs from 'fs';
import * as path from 'path';

function generateDilithiumKeys(keyDir: string = './keys'): Dilithium5 {
  fs.mkdirSync(keyDir, { recursive: true });

  const dk = Dilithium5.generate();

  fs.writeFileSync(path.join(keyDir, 'dilithium_priv.bin'), dk.privateKeyBytes);
  fs.writeFileSync(path.join(keyDir, 'dilithium_pub.bin'), dk.publicKeyBytes);

  console.log(`✓ Dilithium-5 密钥对已生成`);
  console.log(`  公钥 (hex): ${Buffer.from(dk.publicKeyBytes).toString('hex').slice(0, 64)}...`);

  return dk;
}

const dk = generateDilithiumKeys();

3.3 在链上创建 DID

# 使用 msgd CLI 创建 DID
msgd tx wasm execute $DID_REGISTRY_ADDRESS \
  '{"create_did":{
    "did_document": {
      "id": "did:msg:'$(msgd keys show agent_owner -a)'",
      "verification_method": [
        {
          "id": "did:msg:'$(msgd keys show agent_owner -a)'#keys-1",
          "type": "EcdsaSecp256k1RecoveryMethod2020",
          "controller": "did:msg:'$(msgd keys show agent_owner -a)'",
          "blockchain_account_id": "'$(msgd keys show agent_owner -a)'"
        }
      ],
      "authentication": ["did:msg:'$(msgd keys show agent_owner -a)'#keys-1"]
    }
  }}' \
  --from agent_owner \
  --gas auto \
  --gas-adjustment 1.5 \
  --fees 500umsg \
  --output json
# Python: 在链上创建 DID
import asyncio
import json
from msg_chain_sdk import Client, Wallet
from msg_chain_sdk.crypto import Dilithium5

async def create_did():
    # 初始化客户端和钱包
    client = Client(rpc_url=os.getenv("MSG_RPC_URL"))
    wallet = Wallet.from_mnemonic(os.getenv("AGENT_MNEMONIC"))
    agent_address = wallet.address

    # 加载 Dilithium-5 密钥
    dk = Dilithium5.load_private_key("keys/dilithium_priv.bin")

    # 构建 DID 文档
    did_document = {
        "@context": [
            "https://www.w3.org/ns/did/v1",
            "https://msgchain.org/ns/did-config/v1"
        ],
        "id": f"did:msg:{agent_address}",
        "alsoKnownAs": [os.getenv("AGENT_NAME")],
        "controller": [f"did:msg:{agent_address}"],
        "verificationMethod": [
            {
                "id": f"did:msg:{agent_address}#keys-1",
                "type": "EcdsaSecp256k1RecoveryMethod2020",
                "controller": f"did:msg:{agent_address}",
                "blockchainAccountId": agent_address
            },
            {
                "id": f"did:msg:{agent_address}#dilithium-5-keys-1",
                "type": "Dilithium5VerificationKey2026",
                "controller": f"did:msg:{agent_address}",
                "publicKeyMultibase": dk.public_key_multibase()
            }
        ],
        "authentication": [
            f"did:msg:{agent_address}#keys-1",
            f"did:msg:{agent_address}#dilithium-5-keys-1"
        ],
        "assertionMethod": [f"did:msg:{agent_address}#keys-1"],
        "service": [
            {
                "id": f"did:msg:{agent_address}#agent-endpoint",
                "type": "AIAgentEndpoint",
                "serviceEndpoint": os.getenv("AGENT_API_ENDPOINT")
            }
        ],
        "created": "2026-07-06T00:00:00Z",
        "updated": "2026-07-06T00:00:00Z"
    }

    # 构造交易消息
    msg = {
        "create_did": {
            "did_document": did_document
        }
    }

    # 签名并发送交易
    result = await client.execute_contract(
        sender=wallet,
        contract_address=os.getenv("DID_REGISTRY_ADDRESS"),
        msg=msg,
        funds=[],
        gas_limit=500000
    )

    tx_hash = result.transaction_hash
    print(f"✓ DID 创建成功")
    print(f"  DID: did:msg:{agent_address}")
    print(f"  交易哈希: {tx_hash}")

    # 保存 DID 文档
    with open("keys/did_document.json", "w") as f:
        json.dump(did_document, f, indent=2)
    print(f"  DID 文档已保存: keys/did_document.json")

    await client.close()
    return did_document

asyncio.run(create_did())
// TypeScript: 在链上创建 DID
import { MsgChainClient, Wallet } from '@msgchain/sdk';
import { Dilithium5 } from '@msgchain/sdk/crypto';
import * as fs from 'fs';

async function createDID() {
  const client = new MsgChainClient({ rpcUrl: process.env.MSG_RPC_URL! });
  const wallet = await Wallet.fromMnemonic(process.env.AGENT_MNEMONIC!);
  const agentAddress = wallet.address;

  // 加载 Dilithium-5 密钥
  const dk = Dilithium5.loadPrivateKey('keys/dilithium_priv.bin');

  // 构建 DID 文档
  const didDocument = {
    '@context': [
      'https://www.w3.org/ns/did/v1',
      'https://msgchain.org/ns/did-config/v1',
    ],
    id: `did:msg:${agentAddress}`,
    alsoKnownAs: [process.env.AGENT_NAME!],
    controller: [`did:msg:${agentAddress}`],
    verificationMethod: [
      {
        id: `did:msg:${agentAddress}#keys-1`,
        type: 'EcdsaSecp256k1RecoveryMethod2020',
        controller: `did:msg:${agentAddress}`,
        blockchainAccountId: agentAddress,
      },
      {
        id: `did:msg:${agentAddress}#dilithium-5-keys-1`,
        type: 'Dilithium5VerificationKey2026',
        controller: `did:msg:${agentAddress}`,
        publicKeyMultibase: dk.publicKeyMultibase(),
      },
    ],
    authentication: [
      `did:msg:${agentAddress}#keys-1`,
      `did:msg:${agentAddress}#dilithium-5-keys-1`,
    ],
    assertionMethod: [`did:msg:${agentAddress}#keys-1`],
    service: [
      {
        id: `did:msg:${agentAddress}#agent-endpoint`,
        type: 'AIAgentEndpoint',
        serviceEndpoint: process.env.AGENT_API_ENDPOINT!,
      },
    ],
    created: '2026-07-06T00:00:00Z',
    updated: '2026-07-06T00:00:00Z',
  };

  // 执行合约调用
  const result = await client.executeContract(
    wallet,
    process.env.DID_REGISTRY_ADDRESS!,
    { create_did: { did_document: didDocument } },
    [],
    500_000
  );

  console.log(`✓ DID 创建成功`);
  console.log(`  DID: did:msg:${agentAddress}`);
  console.log(`  交易哈希: ${result.transactionHash}`);

  // 保存 DID 文档
  fs.writeFileSync('keys/did_document.json', JSON.stringify(didDocument, null, 2));
  console.log(`  DID 文档已保存: keys/did_document.json`);

  await client.close();
  return didDocument;
}

createDID().catch(console.error);

3.4 查询 DID

# 查询 DID
msgd query wasm contract-state smart $DID_REGISTRY_ADDRESS \
  '{"query_did":{"did":"did:msg:'$(msgd keys show agent_owner -a)'"}}' \
  --output json
# Python: 查询 DID
async def query_did(did: str):
    client = Client(rpc_url=os.getenv("MSG_RPC_URL"))
    result = await client.query_contract(
        os.getenv("DID_REGISTRY_ADDRESS"),
        {"query_did": {"did": did}}
    )
    print(json.dumps(result, indent=2))
    await client.close()
    return result

asyncio.run(query_did("did:msg:msg1..."))
// TypeScript: 查询 DID
async function queryDID(did: string) {
  const client = new MsgChainClient({ rpcUrl: process.env.MSG_RPC_URL! });
  const result = await client.queryContract(
    process.env.DID_REGISTRY_ADDRESS!,
    { query_did: { did } }
  );
  console.log(JSON.stringify(result, null, 2));
  await client.close();
  return result;
}

queryDID('did:msg:msg1...').catch(console.error);

3.5 更新 DID

# Python: 更新 DID 文档
async def update_did():
    client = Client(rpc_url=os.getenv("MSG_RPC_URL"))
    wallet = Wallet.from_mnemonic(os.getenv("AGENT_MNEMONIC"))

    # 添加新的服务端点
    update_msg = {
        "update_did": {
            "did": f"did:msg:{wallet.address}",
            "verification_method": [],
            "service": [
                {
                    "id": f"did:msg:{wallet.address}#new-endpoint",
                    "type": "AIAgentEndpoint",
                    "serviceEndpoint": "https://my-agent.production.com"
                }
            ]
        }
    }

    result = await client.execute_contract(
        sender=wallet,
        contract_address=os.getenv("DID_REGISTRY_ADDRESS"),
        msg=update_msg,
        funds=[],
        gas_limit=300000
    )

    print(f"✓ DID 已更新: {result.transaction_hash}")
    await client.close()

asyncio.run(update_did())

4. 第二步:为 Agent 充值

Agent 需要在链上持有足够的 MSG 代币来支付 Gas 费用、执行合约调用以及进行 AIPAY 支付。

4.1 获取测试代币(测试网)

# 从水龙头获取测试代币
curl -X POST https://faucet.msgchain.org/claim \
  -H "Content-Type: application/json" \
  -d '{"address":"'$(msgd keys show agent_owner -a)'"}'

# 验证余额
msgd query bank balances $(msgd keys show agent_owner -a)

4.2 主网充值

# 从交易所或其他钱包转入
msgd tx bank send \
  $(msgd keys show my_wallet -a) \
  $(msgd keys show agent_owner -a) \
  1000000umsg \
  --from my_wallet \
  --gas auto \
  --gas-adjustment 1.5 \
  --fees 500umsg

4.3 创建 Agent 专用子钱包

建议为主 Agent 创建一个独立的子钱包,避免主钱包私钥暴露风险。

# 创建 Agent 专用钱包
msgd keys add agent_bot

# 转入运营资金
msgd tx bank send \
  $(msgd keys show agent_owner -a) \
  $(msgd keys show agent_bot -a) \
  500000umsg \
  --from agent_owner \
  --gas auto \
  --gas-adjustment 1.5 \
  --fees 500umsg

# 设置环境变量
export AGENT_BOT_ADDR=$(msgd keys show agent_bot -a)
echo "Agent Bot Address: $AGENT_BOT_ADDR"
# Python: 余额管理
async def check_and_fund_agent(min_balance: int = 100000):
    client = Client(rpc_url=os.getenv("MSG_RPC_URL"))
    wallet = Wallet.from_mnemonic(os.getenv("AGENT_MNEMONIC"))

    balance = await client.get_balance(wallet.address, "umsg")
    print(f"当前余额: {balance} umsg")

    if balance < min_balance:
        print(f"⚠  余额不足 (需要 {min_balance} umsg)")
        print(f"请向 {wallet.address} 转入至少 {min_balance - balance} umsg")
        return False

    print(f"✓ 余额充足")
    await client.close()
    return True

asyncio.run(check_and_fund_agent())
// TypeScript: 余额检查
async function checkBalance(minBalance: number = 100_000): Promise<boolean> {
  const client = new MsgChainClient({ rpcUrl: process.env.MSG_RPC_URL! });
  const wallet = await Wallet.fromMnemonic(process.env.AGENT_MNEMONIC!);

  const balance = await client.getBalance(wallet.address, 'umsg');
  console.log(`当前余额: ${balance} umsg`);

  if (balance < minBalance) {
    console.log(`⚠ 余额不足 (需要 ${minBalance} umsg)`);
    return false;
  }

  console.log(`✓ 余额充足`);
  await client.close();
  return true;
}

checkBalance().catch(console.error);

4.4 Gas 费用估算

# 估算交易 Gas
msgd tx wasm execute $DID_REGISTRY_ADDRESS \
  '{"create_did":{"did_document":{"id":"did:msg:test","verification_method":[],"authentication":[]}}}' \
  --from agent_owner \
  --dry-run \
  --output json | jq '.gas_info'

# 示例输出:
# {
#   "gas_wanted": "300000",
#   "gas_used": "245678"
# }
# Python: Gas 估算器
class GasEstimator:
    def __init__(self, client: Client):
        self.client = client
        self.base_gas = 100000
        self.gas_per_byte = 5000

    def estimate_register_agent(self, metadata_size_bytes: int) -> int:
        return self.base_gas + metadata_size_bytes * self.gas_per_byte

    def estimate_execute_constitution(self, steps: int) -> int:
        return self.base_gas * 2 + steps * 50000

    def estimate_aipay_transfer(self) -> int:
        return 150000

    async def simulate_tx(self, msg: dict) -> int:
        result = await self.client.simulate(msg)
        return int(result.gas_info.gas_used * 1.3)  # 加 30% 缓冲

5. 第三步:注册 AI Agent

在创建 DID 之后,下一步是在 Agent Registry 合约中注册 Agent 的元数据、能力描述和发现信息。

5.1 Agent 注册数据结构

{
  "agent_id": "msg1agentaddress...",
  "owner": "msg1owneraddress...",
  "did": "did:msg:msg1agentaddress...",
  "metadata": {
    "name": "my_first_agent",
    "description": "My first AI Agent on MSG Chain",
    "version": "1.0.0",
    "tags": ["defi", "analytics", "automation"],
    "avatar_uri": "ipfs://QmX...",
    "website": "https://my-agent.example.com"
  },
  "capabilities": [
    {
      "name": "token_swap",
      "description": "Execute token swaps via AMM DEX",
      "input_schema": {
        "type": "object",
        "properties": {
          "token_in": {"type": "string"},
          "token_out": {"type": "string"},
          "amount": {"type": "string"}
        },
        "required": ["token_in", "token_out", "amount"]
      },
      "output_schema": {
        "type": "object",
        "properties": {
          "tx_hash": {"type": "string"},
          "amount_out": {"type": "string"}
        }
      },
      "pricing": {
        "model": "per_execution",
        "amount": "10000",
        "denom": "umsg"
      }
    },
    {
      "name": "market_analysis",
      "description": "Analyze market data and provide insights",
      "input_schema": {
        "type": "object",
        "properties": {
          "pair": {"type": "string"},
          "timeframe": {"type": "string", "enum": ["1h", "24h", "7d"]}
        },
        "required": ["pair"]
      },
      "pricing": {
        "model": "per_request",
        "amount": "5000",
        "denom": "umsg"
      }
    }
  ],
  "status": "active",
  "endpoints": {
    "api": "http://localhost:8080",
    "a2a": "http://localhost:8080/a2a",
    "ws": "ws://localhost:8080/ws"
  },
  "stake": {
    "amount": "1000000",
    "denom": "umsg"
  },
  "created_at": "2026-07-06T00:00:00Z",
  "updated_at": "2026-07-06T00:00:00Z"
}

5.2 使用 CLI 注册 Agent

# 注册 Agent
msgd tx wasm execute $AGENT_REGISTRY_ADDRESS \
  '{"register_agent":{
    "metadata":{
      "name":"my_first_agent",
      "description":"My first AI Agent on MSG Chain",
      "version":"1.0.0",
      "tags":["defi","analytics","automation"]
    },
    "capabilities":[
      {
        "name":"token_swap",
        "description":"Execute token swaps via AMM DEX",
        "input_schema":{
          "type":"object",
          "properties":{
            "token_in":{"type":"string"},
            "token_out":{"type":"string"},
            "amount":{"type":"string"}
          },
          "required":["token_in","token_out","amount"]
        },
        "pricing":{"model":"per_execution","amount":"10000","denom":"umsg"}
      }
    ],
    "endpoints":{
      "api":"http://localhost:8080",
      "a2a":"http://localhost:8080/a2a",
      "ws":"ws://localhost:8080/ws"
    },
    "stake":{"amount":"1000000","denom":"umsg"}
  }}' \
  --from agent_owner \
  --gas auto \
  --gas-adjustment 1.5 \
  --fees 500umsg \
  --output json

5.3 使用 Python SDK 注册 Agent

# Python: 注册 Agent
import asyncio
import json
from msg_chain_sdk import Client, Wallet
from msg_chain_sdk.schemas import AgentMetadata, AgentCapability, AgentRegistration

async def register_agent():
    client = Client(rpc_url=os.getenv("MSG_RPC_URL"))
    wallet = Wallet.from_mnemonic(os.getenv("AGENT_MNEMONIC"))

    # 构建 Agent 元数据
    metadata = AgentMetadata(
        name=os.getenv("AGENT_NAME", "my_first_agent"),
        description=os.getenv("AGENT_DESCRIPTION", "My first AI Agent on MSG Chain"),
        version="1.0.0",
        tags=["defi", "analytics", "automation"],
        website="https://my-agent.example.com"
    )

    # 定义能力
    capabilities = [
        AgentCapability(
            name="token_swap",
            description="Execute token swaps via AMM DEX",
            input_schema={
                "type": "object",
                "properties": {
                    "token_in": {"type": "string"},
                    "token_out": {"type": "string"},
                    "amount": {"type": "string"}
                },
                "required": ["token_in", "token_out", "amount"]
            },
            output_schema={
                "type": "object",
                "properties": {
                    "tx_hash": {"type": "string"},
                    "amount_out": {"type": "string"}
                }
            },
            pricing={"model": "per_execution", "amount": "10000", "denom": "umsg"}
        ),
        AgentCapability(
            name="market_analysis",
            description="Analyze market data and provide insights",
            input_schema={
                "type": "object",
                "properties": {
                    "pair": {"type": "string"},
                    "timeframe": {"type": "string", "enum": ["1h", "24h", "7d"]}
                },
                "required": ["pair"]
            },
            pricing={"model": "per_request", "amount": "5000", "denom": "umsg"}
        )
    ]

    # 构建注册消息
    register_msg = {
        "register_agent": {
            "metadata": metadata.model_dump(),
            "capabilities": [c.model_dump() for c in capabilities],
            "endpoints": {
                "api": os.getenv("AGENT_API_ENDPOINT", "http://localhost:8080"),
                "a2a": os.getenv("AGENT_API_ENDPOINT", "http://localhost:8080") + "/a2a",
                "ws": os.getenv("AGENT_API_ENDPOINT", "http://localhost:8080").replace("http", "ws") + "/ws"
            },
            "stake": {"amount": "1000000", "denom": "umsg"}
        }
    }

    # 发送交易
    result = await client.execute_contract(
        sender=wallet,
        contract_address=os.getenv("AGENT_REGISTRY_ADDRESS"),
        msg=register_msg,
        funds=[{"denom": "umsg", "amount": "1000000"}],  # 质押金
        gas_limit=500000
    )

    agent_id = wallet.address
    print(f"✓ Agent 注册成功")
    print(f"  Agent ID: {agent_id}")
    print(f"  交易哈希: {result.transaction_hash}")
    print(f"  名称: {metadata.name}")
    print(f"  能力数: {len(capabilities)}")

    # 保存 Agent 注册信息
    registration = {
        "agent_id": agent_id,
        "metadata": metadata.model_dump(),
        "capabilities": [c.model_dump() for c in capabilities],
        "created_at": "2026-07-06T00:00:00Z"
    }
    with open("agent_registration.json", "w") as f:
        json.dump(registration, f, indent=2)

    await client.close()
    return registration

asyncio.run(register_agent())

5.4 使用 TypeScript SDK 注册 Agent

// TypeScript: 注册 Agent
import { MsgChainClient, Wallet } from '@msgchain/sdk';

interface AgentCapability {
  name: string;
  description: string;
  input_schema: Record<string, any>;
  output_schema?: Record<string, any>;
  pricing: { model: string; amount: string; denom: string };
}

interface AgentRegistration {
  agent_id: string;
  metadata: {
    name: string;
    description: string;
    version: string;
    tags: string[];
    website?: string;
  };
  capabilities: AgentCapability[];
  created_at: string;
}

async function registerAgent(): Promise<AgentRegistration> {
  const client = new MsgChainClient({ rpcUrl: process.env.MSG_RPC_URL! });
  const wallet = await Wallet.fromMnemonic(process.env.AGENT_MNEMONIC!);

  const capabilities: AgentCapability[] = [
    {
      name: 'token_swap',
      description: 'Execute token swaps via AMM DEX',
      input_schema: {
        type: 'object',
        properties: {
          token_in: { type: 'string' },
          token_out: { type: 'string' },
          amount: { type: 'string' },
        },
        required: ['token_in', 'token_out', 'amount'],
      },
      output_schema: {
        type: 'object',
        properties: {
          tx_hash: { type: 'string' },
          amount_out: { type: 'string' },
        },
      },
      pricing: { model: 'per_execution', amount: '10000', denom: 'umsg' },
    },
    {
      name: 'market_analysis',
      description: 'Analyze market data and provide insights',
      input_schema: {
        type: 'object',
        properties: {
          pair: { type: 'string' },
          timeframe: { type: 'string', enum: ['1h', '24h', '7d'] },
        },
        required: ['pair'],
      },
      pricing: { model: 'per_request', amount: '5000', denom: 'umsg' },
    },
  ];

  const registerMsg = {
    register_agent: {
      metadata: {
        name: process.env.AGENT_NAME!,
        description: process.env.AGENT_DESCRIPTION!,
        version: '1.0.0',
        tags: ['defi', 'analytics', 'automation'],
      },
      capabilities,
      endpoints: {
        api: process.env.AGENT_API_ENDPOINT!,
        a2a: `${process.env.AGENT_API_ENDPOINT!}/a2a`,
        ws: `${process.env.AGENT_API_ENDPOINT!.replace('http', 'ws')}/ws`,
      },
      stake: { amount: '1000000', denom: 'umsg' },
    },
  };

  const result = await client.executeContract(
    wallet,
    process.env.AGENT_REGISTRY_ADDRESS!,
    registerMsg,
    [{ denom: 'umsg', amount: '1000000' }],
    500_000
  );

  const agentId = wallet.address;
  console.log(`✓ Agent 注册成功`);
  console.log(`  Agent ID: ${agentId}`);
  console.log(`  交易哈希: ${result.transactionHash}`);

  await client.close();
  return {
    agent_id: agentId,
    metadata: registerMsg.register_agent.metadata,
    capabilities,
    created_at: new Date().toISOString(),
  };
}

registerAgent().catch(console.error);

5.5 查询 Agent 注册信息

# 按 Agent ID 查询
msgd query wasm contract-state smart $AGENT_REGISTRY_ADDRESS \
  '{"query_agent":{"agent_id":"msg1agentaddress..."}}' \
  --output json

# 按能力搜索
msgd query wasm contract-state smart $AGENT_REGISTRY_ADDRESS \
  '{"query_agents_by_capability":{"capability":"token_swap"}}' \
  --output json

# 按标签搜索
msgd query wasm contract-state smart $AGENT_REGISTRY_ADDRESS \
  '{"query_agents_by_tag":{"tag":"defi"}}' \
  --output json

# 列出所有活跃 Agent
msgd query wasm contract-state smart $AGENT_REGISTRY_ADDRESS \
  '{"list_active_agents":{"limit":20,"start_after":null}}' \
  --output json
# Python: 查询 Agent
async def discover_agents(capability: str = None, tag: str = None):
    client = Client(rpc_url=os.getenv("MSG_RPC_URL"))

    if capability:
        query = {"query_agents_by_capability": {"capability": capability}}
        print(f"搜索能力: {capability}")
    elif tag:
        query = {"query_agents_by_tag": {"tag": tag}}
        print(f"搜索标签: {tag}")
    else:
        query = {"list_active_agents": {"limit": 20, "start_after": None}}
        print("列出所有活跃 Agent")

    result = await client.query_contract(
        os.getenv("AGENT_REGISTRY_ADDRESS"),
        query
    )

    agents = result.get("agents", [])
    print(f"找到 {len(agents)} 个 Agent:")

    for agent in agents:
        print(f"  - {agent['metadata']['name']} ({agent['agent_id'][:20]}...)")
        print(f"    能力: {[c['name'] for c in agent.get('capabilities', [])]}")

    await client.close()
    return agents

asyncio.run(discover_agents(capability="token_swap"))
// TypeScript: 发现 Agent
async function discoverAgents(capability?: string, tag?: string) {
  const client = new MsgChainClient({ rpcUrl: process.env.MSG_RPC_URL! });

  let query: Record<string, any>;
  if (capability) {
    query = { query_agents_by_capability: { capability } };
  } else if (tag) {
    query = { query_agents_by_tag: { tag } };
  } else {
    query = { list_active_agents: { limit: 20, start_after: null } };
  }

  const result = await client.queryContract(
    process.env.AGENT_REGISTRY_ADDRESS!,
    query
  );

  const agents = result.agents || [];
  console.log(`找到 ${agents.length} 个 Agent:`);
  agents.forEach((agent: any) => {
    console.log(`  - ${agent.metadata.name} (${agent.agent_id.slice(0, 20)}...)`);
  });

  await client.close();
  return agents;
}

discoverAgents({ capability: 'token_swap' }).catch(console.error);

5.6 更新 Agent 信息

# Python: 更新 Agent
async def update_agent():
    client = Client(rpc_url=os.getenv("MSG_RPC_URL"))
    wallet = Wallet.from_mnemonic(os.getenv("AGENT_MNEMONIC"))

    update_msg = {
        "update_agent": {
            "metadata": {
                "name": "my_first_agent_v2",
                "description": "Updated: Now with A2A support",
                "version": "2.0.0",
                "tags": ["defi", "analytics", "automation", "a2a"]
            },
            "capabilities": [
                {
                    "name": "a2a_query",
                    "description": "Query other agents via A2A protocol",
                    "input_schema": {
                        "type": "object",
                        "properties": {
                            "target_agent": {"type": "string"},
                            "query": {"type": "string"}
                        },
                        "required": ["target_agent", "query"]
                    },
                    "pricing": {"model": "per_request", "amount": "2000", "denom": "umsg"}
                }
            ],
            "status": "active"
        }
    }

    result = await client.execute_contract(
        sender=wallet,
        contract_address=os.getenv("AGENT_REGISTRY_ADDRESS"),
        msg=update_msg,
        funds=[],
        gas_limit=300000
    )

    print(f"✓ Agent 已更新: {result.transaction_hash}")
    await client.close()

asyncio.run(update_agent())

5.7 注销 Agent

# 注销 Agent(取回质押金)
msgd tx wasm execute $AGENT_REGISTRY_ADDRESS \
  '{"deregister_agent":{}}' \
  --from agent_owner \
  --gas auto \
  --gas-adjustment 1.5 \
  --fees 500umsg
# Python: 注销 Agent
async def deregister_agent():
    client = Client(rpc_url=os.getenv("MSG_RPC_URL"))
    wallet = Wallet.from_mnemonic(os.getenv("AGENT_MNEMONIC"))

    result = await client.execute_contract(
        sender=wallet,
        contract_address=os.getenv("AGENT_REGISTRY_ADDRESS"),
        msg={"deregister_agent": {}},
        funds=[],
        gas_limit=200000
    )

    print(f"✓ Agent 已注销,质押金已退还")
    print(f"  交易哈希: {result.transaction_hash}")
    await client.close()

asyncio.run(deregister_agent())

6. 第四步:设定宪章与策略引擎

AI Agent 宪章(Constitution)是约束 Agent 行为的最高规则集合。宪章定义了 Agent 的目标、权限边界、决策原则和行动策略。

6.1 宪章结构

{
  "agent_id": "msg1agentaddress...",
  "constitution": {
    "name": "My Agent Constitution v1",
    "version": "1.0.0",
    "purpose": "Provide decentralized analytics and trading assistance",
    "principles": [
      {
        "id": "principle-001",
        "name": "Safety First",
        "statement": "Never execute actions that could cause irreversible loss of user funds without explicit multi-signature approval",
        "priority": 1
      },
      {
        "id": "principle-002",
        "name": "Transparency",
        "statement": "All decisions and actions must be logged on-chain when possible",
        "priority": 2
      },
      {
        "id": "principle-003",
        "name": "User Autonomy",
        "statement": "Users must always have the final say in any financial decision",
        "priority": 3
      }
    ],
    "policies": [
      {
        "id": "policy-001",
        "name": "Maximum Trade Size",
        "statement": "Single trade limit: 1000 MSG. Requires MPC approval for larger trades",
        "conditions": {
          "type": "amount_limit",
          "max_amount": "1000000000",
          "denom": "umsg",
          "requires_mpc_above": "1000000000"
        }
      },
      {
        "id": "policy-002",
        "name": "Allowed Tokens",
        "statement": "Only trade whitelisted tokens",
        "conditions": {
          "type": "token_whitelist",
          "tokens": ["umsg", "uusdc", "uatom"]
        }
      }
    ],
    "strategy_engine": {
      "model": "rule_based",
      "rules": [
        {
          "id": "rule-001",
          "trigger": "on_market_alert",
          "condition": "price_change > 5%",
          "action": "notify_user",
          "params": {"channel": "telegram"}
        },
        {
          "id": "rule-002",
          "trigger": "on_arbitrage_opportunity",
          "condition": "profit_potential > 0.5% AND gas_cost < profit_potential * 0.3",
          "action": "execute_arbitrage",
          "params": {"max_slippage": "0.5%"}
        }
      ]
    },
    "constraints": [
      "Must maintain minimum balance of 50000 umsg for gas",
      "Must log all trades to chain",
      "Must verify user signature for any withdrawal",
      "Must not interact with blacklisted contracts"
    ],
    "amendments": [
      {
        "id": "amendment-001",
        "description": "Added A2A communication policy",
        "approved_by": ["msg1owneraddress"],
        "timestamp": "2026-07-06T00:00:00Z"
      }
    ]
  }
}

6.2 部署宪章合约

# 设置宪章
msgd tx wasm execute $CONSTITUTION_ADDRESS \
  '{"set_constitution":{
    "constitution":{
      "name":"My Agent Constitution v1",
      "version":"1.0.0",
      "purpose":"Provide decentralized analytics and trading assistance",
      "principles":[
        {"id":"principle-001","name":"Safety First","statement":"Never execute actions that could cause irreversible loss of user funds","priority":1},
        {"id":"principle-002","name":"Transparency","statement":"All decisions must be logged on-chain","priority":2},
        {"id":"principle-003","name":"User Autonomy","statement":"User has final say","priority":3}
      ],
      "policies":[
        {"id":"policy-001","name":"Maximum Trade Size","statement":"Single trade limit","conditions":{"type":"amount_limit","max_amount":"1000000000","denom":"umsg"}}
      ],
      "strategy_engine":{
        "model":"rule_based",
        "rules":[{"id":"rule-001","trigger":"on_market_alert","condition":"price_change > 5%","action":"notify_user","params":{}}]
      },
      "constraints":["Must maintain minimum balance of 50000 umsg"],
      "amendments":[]
    }
  }}' \
  --from agent_owner \
  --gas auto \
  --gas-adjustment 1.5 \
  --fees 500umsg

6.3 使用 Python SDK 设置宪章

# Python: 设置宪章
import asyncio
from msg_chain_sdk import Client, Wallet

async def set_constitution():
    client = Client(rpc_url=os.getenv("MSG_RPC_URL"))
    wallet = Wallet.from_mnemonic(os.getenv("AGENT_MNEMONIC"))

    constitution = {
        "name": "My Agent Constitution v1",
        "version": "1.0.0",
        "purpose": "Provide decentralized analytics and trading assistance",
        "principles": [
            {
                "id": "principle-001",
                "name": "Safety First",
                "statement": "Never execute actions that could cause irreversible loss of user funds without explicit multi-signature approval",
                "priority": 1
            },
            {
                "id": "principle-002",
                "name": "Transparency",
                "statement": "All decisions and actions must be logged on-chain when possible",
                "priority": 2
            },
            {
                "id": "principle-003",
                "name": "User Autonomy",
                "statement": "Users must always have the final say in any financial decision",
                "priority": 3
            }
        ],
        "policies": [
            {
                "id": "policy-001",
                "name": "Maximum Trade Size",
                "statement": "Single trade limit: 1000 MSG. Requires MPC approval for larger trades",
                "conditions": {
                    "type": "amount_limit",
                    "max_amount": "1000000000",
                    "denom": "umsg",
                    "requires_mpc_above": "1000000000"
                }
            },
            {
                "id": "policy-002",
                "name": "Allowed Tokens",
                "statement": "Only trade whitelisted tokens",
                "conditions": {
                    "type": "token_whitelist",
                    "tokens": ["umsg", "uusdc", "uatom"]
                }
            }
        ],
        "strategy_engine": {
            "model": "rule_based",
            "rules": [
                {
                    "id": "rule-001",
                    "trigger": "on_market_alert",
                    "condition": "price_change > 5%",
                    "action": "notify_user",
                    "params": {"channel": "telegram"}
                },
                {
                    "id": "rule-002",
                    "trigger": "on_arbitrage_opportunity",
                    "condition": "profit_potential > 0.5% AND gas_cost < profit_potential * 0.3",
                    "action": "execute_arbitrage",
                    "params": {"max_slippage": "0.5%"}
                }
            ]
        },
        "constraints": [
            "Must maintain minimum balance of 50000 umsg for gas",
            "Must log all trades to chain",
            "Must verify user signature for any withdrawal"
        ],
        "amendments": []
    }

    set_msg = {
        "set_constitution": {
            "constitution": constitution
        }
    }

    result = await client.execute_contract(
        sender=wallet,
        contract_address=os.getenv("CONSTITUTION_ADDRESS"),
        msg=set_msg,
        funds=[],
        gas_limit=600000
    )

    print(f"✓ 宪章设置成功")
    print(f"  宪章名称: {constitution['name']}")
    print(f"  原则数: {len(constitution['principles'])}")
    print(f"  策略规则数: {len(constitution['strategy_engine']['rules'])}")
    print(f"  交易哈希: {result.transaction_hash}")

    await client.close()
    return constitution

asyncio.run(set_constitution())

6.4 使用 TypeScript SDK 设置宪章

// TypeScript: 设置宪章
import { MsgChainClient, Wallet } from '@msgchain/sdk';

interface Constitution {
  name: string;
  version: string;
  purpose: string;
  principles: Array<{ id: string; name: string; statement: string; priority: number }>;
  policies: Array<{ id: string; name: string; statement: string; conditions: Record<string, any> }>;
  strategy_engine: {
    model: string;
    rules: Array<{ id: string; trigger: string; condition: string; action: string; params: Record<string, any> }>;
  };
  constraints: string[];
  amendments: Array<{ id: string; description: string; approved_by: string[]; timestamp: string }>;
}

async function setConstitution(): Promise<Constitution> {
  const client = new MsgChainClient({ rpcUrl: process.env.MSG_RPC_URL! });
  const wallet = await Wallet.fromMnemonic(process.env.AGENT_MNEMONIC!);

  const constitution: Constitution = {
    name: 'My Agent Constitution v1',
    version: '1.0.0',
    purpose: 'Provide decentralized analytics and trading assistance',
    principles: [
      {
        id: 'principle-001',
        name: 'Safety First',
        statement: 'Never execute actions that could cause irreversible loss of user funds',
        priority: 1,
      },
      {
        id: 'principle-002',
        name: 'Transparency',
        statement: 'All decisions and actions must be logged on-chain when possible',
        priority: 2,
      },
    ],
    policies: [
      {
        id: 'policy-001',
        name: 'Maximum Trade Size',
        statement: 'Single trade limit: 1000 MSG',
        conditions: {
          type: 'amount_limit',
          max_amount: '1000000000',
          denom: 'umsg',
        },
      },
    ],
    strategy_engine: {
      model: 'rule_based',
      rules: [
        {
          id: 'rule-001',
          trigger: 'on_market_alert',
          condition: 'price_change > 5%',
          action: 'notify_user',
          params: { channel: 'telegram' },
        },
      ],
    },
    constraints: ['Must maintain minimum balance of 50000 umsg for gas'],
    amendments: [],
  };

  const result = await client.executeContract(
    wallet,
    process.env.CONSTITUTION_ADDRESS!,
    { set_constitution: { constitution } },
    [],
    600_000
  );

  console.log(`✓ 宪章设置成功`);
  console.log(`  名称: ${constitution.name}`);
  console.log(`  交易哈希: ${result.transactionHash}`);

  await client.close();
  return constitution;
}

setConstitution().catch(console.error);

6.5 策略引擎运行时

# Python: 策略引擎执行
import asyncio
from enum import Enum
from typing import Any, Callable, Dict, List, Optional

class RuleTrigger(Enum):
    ON_MARKET_ALERT = "on_market_alert"
    ON_ARBITRAGE = "on_arbitrage_opportunity"
    ON_NEW_TRANSACTION = "on_new_transaction"
    ON_PRICE_UPDATE = "on_price_update"
    ON_SCHEDULE = "on_schedule"

class StrategyEngine:
    def __init__(self, constitution: Dict):
        self.constitution = constitution
        self.rules = constitution["strategy_engine"]["rules"]
        self.policies = constitution["policies"]
        self.principles = constitution["principles"]
        self._handlers: Dict[str, Callable] = {}

    def register_handler(self, trigger: str, handler: Callable):
        self._handlers[trigger] = handler

    async def evaluate(self, trigger: str, context: Dict) -> Optional[Dict]:
        """评估所有匹配的规则并返回最高优先级的动作"""
        matching_rules = [r for r in self.rules if r["trigger"] == trigger]

        if not matching_rules:
            return None

        # 按条件复杂度排序
        for rule in matching_rules:
            if await self._evaluate_condition(rule["condition"], context):
                # 检查策略约束
                if await self._check_policies(rule["action"], context):
                    print(f"✓ 规则触发: {rule['id']} -> {rule['action']}")
                    return {
                        "rule_id": rule["id"],
                        "action": rule["action"],
                        "params": rule.get("params", {})
                    }

        return None

    async def _evaluate_condition(self, condition: str, context: Dict) -> bool:
        """评估规则条件(简化示例,生产环境使用安全沙箱)"""
        try:
            # 替换上下文变量
            eval_context = {}
            for key, value in context.items():
                if isinstance(value, (int, float, str, bool)):
                    eval_context[key] = value

            # 安全的条件评估
            # 注意: 生产环境应使用受限的表达式求值器
            safe_globals = {"__builtins__": {}}
            safe_locals = eval_context

            result = eval(condition, safe_globals, safe_locals)
            return bool(result)
        except Exception as e:
            print(f"⚠ 条件评估失败: {e}")
            return False

    async def _check_policies(self, action: str, context: Dict) -> bool:
        """检查动作是否符合所有策略"""
        for policy in self.policies:
            condition = policy["conditions"]
            if condition["type"] == "amount_limit":
                amount = int(context.get("amount", 0))
                max_amount = int(condition["max_amount"])
                if amount > max_amount:
                    print(f"⚠ 策略限制: {policy['name']} (限额 {max_amount})")
                    return False
            elif condition["type"] == "token_whitelist":
                token = context.get("token", "")
                if token and token not in condition["tokens"]:
                    print(f"⚠ 策略限制: {policy['name']} (代币不在白名单)")
                    return False
        return True

    async def execute_action(self, action: str, params: Dict):
        """执行动作"""
        handler = self._handlers.get(action)
        if handler:
            await handler(params)
        else:
            print(f"⚠ 未注册处理器: {action}")

# 使用示例
async def strategy_example():
    # 加载宪章
    with open("constitution.json") as f:
        import json
        constitution = json.load(f)

    engine = StrategyEngine(constitution)

    # 注册动作处理器
    async def notify_user(params):
        print(f"通知用户: {params}")

    async def execute_arbitrage(params):
        print(f"执行套利: {params}")

    engine.register_handler("notify_user", notify_user)
    engine.register_handler("execute_arbitrage", execute_arbitrage)

    # 模拟市场告警事件
    context = {
        "price_change": 6.2,
        "pair": "MSG/USDC",
        "current_price": 1.25,
        "amount": 500000
    }

    result = await engine.evaluate("on_market_alert", context)
    if result:
        await engine.execute_action(result["action"], result["params"])

asyncio.run(strategy_example())

6.6 查询宪章状态

msgd query wasm contract-state smart $CONSTITUTION_ADDRESS \
  '{"query_constitution":{"agent_id":"msg1agentaddress..."}}' \
  --output json
# Python: 查询宪章
async def query_constitution(agent_id: str):
    client = Client(rpc_url=os.getenv("MSG_RPC_URL"))
    result = await client.query_contract(
        os.getenv("CONSTITUTION_ADDRESS"),
        {"query_constitution": {"agent_id": agent_id}}
    )
    print(json.dumps(result, indent=2, ensure_ascii=False))
    await client.close()
    return result

asyncio.run(query_constitution("msg1agentaddress..."))

7. 第五步:连接 Agent API 网关

Agent API 网关是外界与 Agent 交互的入口。它处理身份验证、请求路由、速率限制和与链上合约的通信。

7.1 API 网关架构

┌──────────┐     ┌──────────────┐     ┌──────────────┐     ┌──────────┐
│ 客户端    │────▶│  API 网关     │────▶│  策略引擎     │────▶│ 链上合约  │
│ (用户/    │     │  :8080       │     │  (宪章检查)   │     │          │
│  其他Agent)│     │  TLS + 认证   │     │  规则评估     │     │  Registry│
└──────────┘     └──────────────┘     └──────────────┘     └──────────┘
                        │
                        ▼
                 ┌──────────────┐
                 │  事件队列     │
                 │  (Redis/NATS) │
                 └──────────────┘

7.2 Python FastAPI 网关实现

# Python: Agent API 网关 (gateway.py)
import asyncio
import json
import os
import hashlib
from datetime import datetime
from typing import Optional, Dict, Any, List

from fastapi import FastAPI, HTTPException, Request, Depends, Header
from fastapi.middleware.cors import CORSMiddleware
from fastapi.responses import JSONResponse
from pydantic import BaseModel, Field
import httpx

from msg_chain_sdk import Client, Wallet
from msg_chain_sdk.crypto import Dilithium5, verify_message

# ==================== 数据模型 ====================

class ExecuteRequest(BaseModel):
    capability: str
    params: Dict[str, Any]
    user_address: Optional[str] = None
    agent_session_id: Optional[str] = None
    idempotency_key: Optional[str] = None

class ExecuteResponse(BaseModel):
    success: bool
    result: Optional[Dict[str, Any]] = None
    error: Optional[str] = None
    session_id: str
    gas_used: Optional[int] = None

class AgentStatus(BaseModel):
    agent_id: str
    name: str
    status: str
    uptime: float
    total_requests: int
    total_errors: int
    balance: str
    version: str

class QueryRequest(BaseModel):
    query_type: str
    params: Dict[str, Any]

# ==================== 网关应用 ====================

app = FastAPI(
    title="AI Agent API Gateway",
    description="MSG Chain AI Agent Gateway",
    version="1.0.0"
)

app.add_middleware(
    CORSMiddleware,
    allow_origins=["*"],
    allow_credentials=True,
    allow_methods=["*"],
    allow_headers=["*"],
)

# ==================== 全局状态 ====================

class AgentState:
    def __init__(self):
        self.wallet: Optional[Wallet] = None
        self.client: Optional[Client] = None
        self.dilithium_key: Optional[Dilithium5] = None
        self.chain_client: Optional[MsgChainClient] = None
        self.start_time = datetime.utcnow()
        self.total_requests = 0
        self.total_errors = 0
        self.idempotency_cache: Dict[str, Any] = {}
        self.session_store: Dict[str, Dict] = {}

state = AgentState()

# ==================== 依赖注入 ====================

async def get_chain_client() -> Client:
    if state.client is None:
        state.client = Client(rpc_url=os.getenv("MSG_RPC_URL"))
        if os.getenv("AGENT_MNEMONIC"):
            state.wallet = Wallet.from_mnemonic(os.getenv("AGENT_MNEMONIC"))
    return state.client

async def verify_agent_signature(
    request: Request,
    x_signature: Optional[str] = Header(None),
    x_timestamp: Optional[str] = Header(None),
    x_agent_id: Optional[str] = Header(None)
):
    """验证 Agent 签名(用于 Agent 间通信)"""
    if not x_signature or not x_timestamp or not x_agent_id:
        # 允许未签名请求(用户查询等只读操作)
        return None

    # 验证时间戳(防止重放攻击,5 分钟窗口)
    try:
        ts = int(x_timestamp)
        if abs(datetime.utcnow().timestamp() - ts) > 300:
            raise HTTPException(status_code=401, detail="签名已过期")
    except ValueError:
        raise HTTPException(status_code=401, detail="无效的时间戳")

    # 查询发送方 Agent 的公钥
    client = await get_chain_client()
    did_result = await client.query_contract(
        os.getenv("DID_REGISTRY_ADDRESS"),
        {"query_did": {"did": f"did:msg:{x_agent_id}"}}
    )

    # 验证签名
    body = await request.body()
    message = body.decode() + x_timestamp
    # 简化验证 - 生产环境使用完整的 Dilithium-5 验证
    return x_agent_id

# ==================== 路由 ====================

@app.get("/health")
async def health_check():
    return {
        "status": "ok",
        "agent_id": state.wallet.address if state.wallet else None,
        "uptime": (datetime.utcnow() - state.start_time).total_seconds(),
        "timestamp": datetime.utcnow().isoformat()
    }

@app.get("/agent/status", response_model=AgentStatus)
async def get_agent_status(client: Client = Depends(get_chain_client)):
    if not state.wallet:
        raise HTTPException(status_code=500, detail="Agent 未初始化")

    balance = await client.get_balance(state.wallet.address, "umsg")

    return AgentStatus(
        agent_id=state.wallet.address,
        name=os.getenv("AGENT_NAME", "unknown"),
        status="active",
        uptime=(datetime.utcnow() - state.start_time).total_seconds(),
        total_requests=state.total_requests,
        total_errors=state.total_errors,
        balance=str(balance),
        version="1.0.0"
    )

@app.post("/execute", response_model=ExecuteResponse)
async def execute_capability(
    req: ExecuteRequest,
    client: Client = Depends(get_chain_client)
):
    state.total_requests += 1

    # 幂等性检查
    if req.idempotency_key:
        cached = state.idempotency_cache.get(req.idempotency_key)
        if cached:
            return ExecuteResponse(**cached)

    # 生成会话 ID
    session_id = hashlib.sha256(
        f"{state.wallet.address}:{datetime.utcnow().timestamp()}:{req.capability}".encode()
    ).hexdigest()[:16]

    try:
        # 步骤 1: 验证宪章约束
        constitution_check = await client.query_contract(
            os.getenv("CONSTITUTION_ADDRESS"),
            {
                "check_action": {
                    "agent_id": state.wallet.address,
                    "capability": req.capability,
                    "params": req.params
                }
            }
        )

        if not constitution_check.get("allowed", True):
            raise HTTPException(
                status_code=403,
                detail=f"宪章约束拒绝: {constitution_check.get('reason', 'unknown')}"
            )

        # 步骤 2: 执行能力
        result = await execute_capability_handler(req.capability, req.params, client)

        # 步骤 3: 记录会话
        session = {
            "session_id": session_id,
            "capability": req.capability,
            "params": req.params,
            "result": result,
            "timestamp": datetime.utcnow().isoformat(),
            "user_address": req.user_address
        }
        state.session_store[session_id] = session

        response = ExecuteResponse(
            success=True,
            result=result,
            session_id=session_id
        )

        # 缓存幂等键
        if req.idempotency_key:
            state.idempotency_cache[req.idempotency_key] = response.model_dump()

        return response

    except HTTPException:
        raise
    except Exception as e:
        state.total_errors += 1
        return ExecuteResponse(
            success=False,
            error=str(e),
            session_id=session_id
        )

@app.post("/query")
async def query_chain(
    req: QueryRequest,
    client: Client = Depends(get_chain_client)
):
    """链上查询代理(只读操作,不需要 Gas)"""
    try:
        result = await client.query_contract(
            os.getenv("AGENT_REGISTRY_ADDRESS"),
            {req.query_type: req.params}
        )
        return {"success": True, "result": result}
    except Exception as e:
        raise HTTPException(status_code=400, detail=str(e))

@app.get("/sessions/{session_id}")
async def get_session(session_id: str):
    """查询会话详情"""
    session = state.session_store.get(session_id)
    if not session:
        raise HTTPException(status_code=404, detail="会话不存在")
    return session

# ==================== 能力执行器 ====================

async def execute_capability_handler(
    capability: str,
    params: Dict[str, Any],
    client: Client
) -> Dict[str, Any]:
    """实际的能力执行逻辑"""
    handlers = {
        "token_swap": handle_token_swap,
        "market_analysis": handle_market_analysis,
        "a2a_query": handle_a2a_query,
    }

    handler = handlers.get(capability)
    if not handler:
        raise HTTPException(status_code=400, detail=f"不支持的能力: {capability}")

    return await handler(params, client)

async def handle_token_swap(params: Dict[str, Any], client: Client) -> Dict[str, Any]:
    """代币交换处理"""
    token_in = params.get("token_in")
    token_out = params.get("token_out")
    amount = params.get("amount")

    if not all([token_in, token_out, amount]):
        raise HTTPException(status_code=400, detail="缺少必要参数")

    # 这里调用 DEX 合约
    print(f"执行代币交换: {amount} {token_in} -> {token_out}")

    # 模拟执行
    return {
        "tx_hash": "simulated_tx_hash_123",
        "amount_out": str(int(amount) * 98 // 100),  # 2% slippage
        "token_out": token_out
    }

async def handle_market_analysis(params: Dict[str, Any], client: Client) -> Dict[str, Any]:
    """市场分析处理"""
    pair = params.get("pair")
    timeframe = params.get("timeframe", "24h")

    # 从链上或预言机获取数据
    return {
        "pair": pair,
        "timeframe": timeframe,
        "current_price": "1.25",
        "price_change_24h": "+3.2%",
        "volume_24h": "1250000",
        "analysis": "Bullish trend with increasing volume"
    }

async def handle_a2a_query(params: Dict[str, Any], client: Client) -> Dict[str, Any]:
    """A2A 通信处理"""
    target_agent = params.get("target_agent")
    query = params.get("query")

    # 查询目标 Agent 的 API 端点
    agent_info = await client.query_contract(
        os.getenv("AGENT_REGISTRY_ADDRESS"),
        {"query_agent": {"agent_id": target_agent}}
    )

    endpoint = agent_info.get("endpoints", {}).get("a2a")
    if not endpoint:
        raise HTTPException(status_code=404, detail="目标 Agent 不支持 A2A")

    # 调用目标 Agent
    async with httpx.AsyncClient() as http_client:
        response = await http_client.post(
            f"{endpoint}/a2a/query",
            json={"query": query, "source_agent": state.wallet.address},
            timeout=30
        )
        return response.json()

# ==================== 启动 ====================

@app.on_event("startup")
async def startup():
    # 初始化链客户端
    state.client = Client(rpc_url=os.getenv("MSG_RPC_URL"))
    if os.getenv("AGENT_MNEMONIC"):
        state.wallet = Wallet.from_mnemonic(os.getenv("AGENT_MNEMONIC"))
        print(f"Agent 钱包: {state.wallet.address}")

    # 加载 Dilithium-5 密钥
    if os.path.exists("keys/dilithium_priv.bin"):
        state.dilithium_key = Dilithium5.load_private_key("keys/dilithium_priv.bin")
        print("✓ Dilithium-5 密钥已加载")

@app.on_event("shutdown")
async def shutdown():
    if state.client:
        await state.client.close()

# ==================== 主入口 ====================

if __name__ == "__main__":
    import uvicorn
    port = int(os.getenv("AGENT_PORT", "8080"))
    uvicorn.run(app, host="0.0.0.0", port=port, log_level="info")

7.3 启动 API 网关

# 使用 uvicorn 启动
uvicorn gateway:app --host 0.0.0.0 --port 8080 --reload

# 或使用 Python 直接运行
python gateway.py

7.4 使用 TypeScript Express 实现

// TypeScript: Agent API 网关 (src/gateway.ts)
import express, { Request, Response, NextFunction } from 'express';
import { MsgChainClient, Wallet } from '@msgchain/sdk';
import { Dilithium5 } from '@msgchain/sdk/crypto';
import crypto from 'crypto';
import axios from 'axios';

const app = express();
app.use(express.json());

// ==================== 类型定义 ====================

interface ExecuteRequest {
  capability: string;
  params: Record<string, any>;
  user_address?: string;
  agent_session_id?: string;
  idempotency_key?: string;
}

interface ExecuteResponse {
  success: boolean;
  result?: Record<string, any>;
  error?: string;
  session_id: string;
  gas_used?: number;
}

interface AgentState {
  wallet?: Wallet;
  chainClient?: MsgChainClient;
  dilithiumKey?: Dilithium5;
  startTime: Date;
  totalRequests: number;
  totalErrors: number;
  idempotencyCache: Map<string, ExecuteResponse>;
  sessionStore: Map<string, any>;
}

// ==================== 状态 ====================

const state: AgentState = {
  startTime: new Date(),
  totalRequests: 0,
  totalErrors: 0,
  idempotencyCache: new Map(),
  sessionStore: new Map(),
};

// ==================== 中间件 ====================

async function authMiddleware(req: Request, res: Response, next: NextFunction) {
  const signature = req.headers['x-signature'] as string;
  const timestamp = req.headers['x-timestamp'] as string;
  const agentId = req.headers['x-agent-id'] as string;

  if (!signature || !timestamp || !agentId) {
    // 只读请求不需要签名
    return next();
  }

  // 验证时间戳(5 分钟窗口)
  const ts = parseInt(timestamp);
  if (isNaN(ts) || Math.abs(Date.now() / 1000 - ts) > 300) {
    return res.status(401).json({ error: '签名已过期' });
  }

  // 生产环境:验证 Dilithium-5 签名
  next();
}

app.use(authMiddleware);

// ==================== 路由 ====================

app.get('/health', (req: Request, res: Response) => {
  res.json({
    status: 'ok',
    agent_id: state.wallet?.address || null,
    uptime: (Date.now() - state.startTime.getTime()) / 1000,
  });
});

app.get('/agent/status', async (req: Request, res: Response) => {
  if (!state.wallet || !state.chainClient) {
    return res.status(500).json({ error: 'Agent 未初始化' });
  }

  const balance = await state.chainClient.getBalance(state.wallet.address, 'umsg');

  res.json({
    agent_id: state.wallet.address,
    name: process.env.AGENT_NAME || 'unknown',
    status: 'active',
    uptime: (Date.now() - state.startTime.getTime()) / 1000,
    total_requests: state.totalRequests,
    total_errors: state.totalErrors,
    balance: balance.toString(),
    version: '1.0.0',
  });
});

app.post('/execute', async (req: Request, res: Response) => {
  state.totalRequests++;
  const execReq = req.body as ExecuteRequest;

  // 幂等性检查
  if (execReq.idempotency_key) {
    const cached = state.idempotencyCache.get(execReq.idempotency_key);
    if (cached) {
      return res.json(cached);
    }
  }

  const sessionId = crypto
    .createHash('sha256')
    .update(`${state.wallet?.address}:${Date.now()}:${execReq.capability}`)
    .digest('hex')
    .slice(0, 16);

  try {
    // 宪章约束检查
    if (state.chainClient) {
      const constitutionCheck = await state.chainClient.queryContract(
        process.env.CONSTITUTION_ADDRESS!,
        {
          check_action: {
            agent_id: state.wallet?.address,
            capability: execReq.capability,
            params: execReq.params,
          },
        }
      );

      if (!constitutionCheck.allowed) {
        return res.status(403).json({
          success: false,
          error: `宪章约束拒绝: ${constitutionCheck.reason || 'unknown'}`,
          session_id: sessionId,
        });
      }
    }

    // 执行能力
    const result = await executeCapability(execReq.capability, execReq.params);

    const response: ExecuteResponse = {
      success: true,
      result,
      session_id: sessionId,
    };

    // 缓存幂等键
    if (execReq.idempotency_key) {
      state.idempotencyCache.set(execReq.idempotency_key, response);
    }

    return res.json(response);
  } catch (error: any) {
    state.totalErrors++;
    return res.status(500).json({
      success: false,
      error: error.message,
      session_id: sessionId,
    });
  }
});

// ==================== 能力执行器 ====================

async function executeCapability(
  capability: string,
  params: Record<string, any>
): Promise<Record<string, any>> {
  switch (capability) {
    case 'token_swap':
      return handleTokenSwap(params);
    case 'market_analysis':
      return handleMarketAnalysis(params);
    case 'a2a_query':
      return handleA2AQuery(params);
    default:
      throw new Error(`不支持的能力: ${capability}`);
  }
}

async function handleTokenSwap(params: Record<string, any>): Promise<Record<string, any>> {
  const { token_in, token_out, amount } = params;
  if (!token_in || !token_out || !amount) {
    throw new Error('缺少必要参数');
  }
  return {
    tx_hash: `sim_${crypto.randomBytes(16).toString('hex')}`,
    amount_out: String(Math.floor(Number(amount) * 0.98)),
    token_out,
  };
}

async function handleMarketAnalysis(params: Record<string, any>): Promise<Record<string, any>> {
  return {
    pair: params.pair,
    timeframe: params.timeframe || '24h',
    current_price: '1.25',
    price_change_24h: '+3.2%',
    analysis: 'Bullish trend with increasing volume',
  };
}

async function handleA2AQuery(params: Record<string, any>): Promise<Record<string, any>> {
  const { target_agent, query } = params;

  const agentInfo = await state.chainClient?.queryContract(
    process.env.AGENT_REGISTRY_ADDRESS!,
    { query_agent: { agent_id: target_agent } }
  );

  const endpoint = agentInfo?.endpoints?.a2a;
  if (!endpoint) {
    throw new Error('目标 Agent 不支持 A2A');
  }

  const response = await axios.post(`${endpoint}/a2a/query`, {
    query,
    source_agent: state.wallet?.address,
  });

  return response.data;
}

// ==================== 启动 ====================

async function main() {
  const port = parseInt(process.env.AGENT_PORT || '8080');

  // 初始化链客户端
  state.chainClient = new MsgChainClient({
    rpcUrl: process.env.MSG_RPC_URL!,
  });

  // 加载钱包
  if (process.env.AGENT_MNEMONIC) {
    state.wallet = await Wallet.fromMnemonic(process.env.AGENT_MNEMONIC);
    console.log(`Agent 钱包: ${state.wallet.address}`);
  }

  // 加载 Dilithium-5 密钥
  const fs = await import('fs');
  if (fs.existsSync('keys/dilithium_priv.bin')) {
    state.dilithiumKey = Dilithium5.loadPrivateKey('keys/dilithium_priv.bin');
    console.log('✓ Dilithium-5 密钥已加载');
  }

  app.listen(port, () => {
    console.log(`Agent API 网关运行在 http://0.0.0.0:${port}`);
  });
}

main().catch(console.error);

export default app;

7.5 测试 API

# 健康检查
curl http://localhost:8080/health

# 获取 Agent 状态
curl http://localhost:8080/agent/status

# 执行能力
curl -X POST http://localhost:8080/execute \
  -H "Content-Type: application/json" \
  -d '{
    "capability":"market_analysis",
    "params":{"pair":"MSG/USDC","timeframe":"24h"}
  }'

# 链上查询
curl -X POST http://localhost:8080/query \
  -H "Content-Type: application/json" \
  -d '{
    "query_type":"query_agent",
    "params":{"agent_id":"msg1agentaddress..."}
  }'

# 查询会话
curl http://localhost:8080/sessions/{session_id}

8. 第六步:订阅链上事件

Agent 需要监听链上事件来自动响应市场变化、合约调用和治理提案。

8.1 WebSocket 连接

# 使用 wscat 测试 WebSocket 连接
wscat -c wss://rpc.msgchain.org:443/websocket

# 订阅新区块事件
> {"jsonrpc":"2.0","method":"subscribe","id":"1","params":{"query":"tm.event='NewBlock'"}}

# 订阅交易事件
> {"jsonrpc":"2.0","method":"subscribe","id":"2","params":{"query":"tm.event='Tx'"}}

# 订阅合约事件(按合约地址过滤)
> {"jsonrpc":"2.0","method":"subscribe","id":"3","params":{"query":"wasm._contract_address='msg1qypqxpq9kcrn2c9afea5lq35ef37c5x7jqylz3' AND wasm.action='agent_registered'"}}

8.2 Python 事件订阅器

# Python: 事件订阅器 (event_subscriber.py)
import asyncio
import json
import os
from datetime import datetime
from typing import Callable, Dict, List, Optional, Any
from enum import Enum

import websockets
import httpx

class EventType(Enum):
    NEW_BLOCK = "tm.event='NewBlock'"
    NEW_TX = "tm.event='Tx'"
    AGENT_REGISTERED = "wasm._contract_address='{}' AND wasm.action='agent_registered'"
    CONSTITUTION_UPDATED = "wasm._contract_address='{}' AND wasm.action='constitution_updated'"
    AIPAY_RECEIVED = "wasm._contract_address='{}' AND wasm.action='payment_received'"

class EventSubscriber:
    def __init__(
        self,
        ws_url: str,
        agent_address: str,
        rpc_url: str,
        handlers: Optional[Dict[str, Callable]] = None
    ):
        self.ws_url = ws_url
        self.agent_address = agent_address
        self.rpc_url = rpc_url
        self.handlers = handlers or {}
        self._running = False
        self._subscriptions: Dict[str, int] = {}
        self._request_id = 0

    def register_handler(self, event_type: str, handler: Callable):
        """注册事件处理器"""
        self.handlers[event_type] = handler
        print(f"✓ 已注册处理器: {event_type}")

    async def connect(self):
        """建立 WebSocket 连接并订阅事件"""
        self._running = True

        async for websocket in websockets.connect(
            self.ws_url,
            ping_interval=30,
            ping_timeout=10
        ):
            try:
                print(f"✓ WebSocket 已连接: {self.ws_url}")

                # 订阅多个事件
                subscriptions = [
                    ("new_block", f"tm.event='NewBlock'"),
                    ("new_tx", f"tm.event='Tx'"),
                    (
                        "agent_registered",
                        f"wasm._contract_address='{os.getenv('AGENT_REGISTRY_ADDRESS')}'"
                    ),
                    (
                        "constitution_updated",
                        f"wasm._contract_address='{os.getenv('CONSTITUTION_ADDRESS')}'"
                    ),
                    (
                        "payment_received",
                        f"wasm._contract_address='{os.getenv('AIPAY_MODULE_ADDRESS')}'"
                    ),
                ]

                for name, query in subscriptions:
                    await self._subscribe(websocket, name, query)

                # 持续监听事件
                async for message in websocket:
                    if not self._running:
                        break
                    await self._handle_message(json.loads(message))

            except websockets.ConnectionClosed:
                print("⚠ WebSocket 断开,3 秒后重连...")
                await asyncio.sleep(3)
            except Exception as e:
                print(f"⚠ WebSocket 错误: {e}")
                await asyncio.sleep(5)

    async def _subscribe(self, websocket, name: str, query: str):
        """发送订阅请求"""
        self._request_id += 1
        subscribe_msg = {
            "jsonrpc": "2.0",
            "method": "subscribe",
            "id": str(self._request_id),
            "params": {"query": query}
        }
        await websocket.send(json.dumps(subscribe_msg))
        self._subscriptions[name] = self._request_id
        print(f"  ✓ 已订阅: {name}")

    async def _handle_message(self, message: Dict):
        """处理收到的消息"""
        try:
            result = message.get("result", {})
            data = result.get("data", {})
            value = data.get("value", {})

            event_type = self._detect_event_type(value)
            if event_type and event_type in self.handlers:
                await self.handlers[event_type](value)
            else:
                # 默认处理:记录到日志
                self._log_event(value)

        except Exception as e:
            print(f"⚠ 消息处理错误: {e}")

    def _detect_event_type(self, value: Dict) -> Optional[str]:
        """检测事件类型"""
        events = value.get("events", [])
        for event in events:
            event_type = event.get("type", "")
            attributes = event.get("attributes", [])
            for attr in attributes:
                key = attr.get("key", "")
                if key == "action":
                    return attr.get("value", "")
        return None

    def _log_event(self, value: Dict):
        """记录事件到日志文件"""
        timestamp = datetime.utcnow().isoformat()
        log_entry = {
            "timestamp": timestamp,
            "height": value.get("header", {}).get("height"),
            "data": value
        }

        with open("events.log", "a") as f:
            f.write(json.dumps(log_entry) + "\n")

    def stop(self):
        """停止订阅"""
        self._running = False

# ==================== 使用示例 ====================

async def handle_new_block(data: Dict):
    """新区块处理器"""
    height = data.get("header", {}).get("height", "unknown")
    print(f"📦 新区块: #{height}")

async def handle_agent_registered(data: Dict):
    """Agent 注册事件处理器"""
    print(f"🤖 新 Agent 注册!")
    for event in data.get("events", []):
        print(f"  Event: {event.get('type', 'unknown')}")

async def handle_payment_received(data: Dict):
    """支付事件处理器"""
    print(f"💰 收到支付!")
    events = data.get("events", [])
    for event in events:
        for attr in event.get("attributes", []):
            if attr.get("key") == "amount":
                print(f"  金额: {attr.get('value')}")

async def run_subscriber():
    subscriber = EventSubscriber(
        ws_url=os.getenv("MSG_WEBSOCKET_URL", "wss://rpc.msgchain.org:443/websocket"),
        agent_address=os.getenv("AGENT_BOT_ADDR", ""),
        rpc_url=os.getenv("MSG_RPC_URL", "https://rpc.msgchain.org:443")
    )

    # 注册处理器
    subscriber.register_handler("new_block", handle_new_block)
    subscriber.register_handler("agent_registered", handle_agent_registered)
    subscriber.register_handler("payment_received", handle_payment_received)

    print("🚀 事件订阅器启动中...")
    await subscriber.connect()

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

8.3 TypeScript 事件订阅器

// TypeScript: 事件订阅器 (src/eventSubscriber.ts)
import WebSocket from 'ws';
import * as fs from 'fs';
import * as path from 'path';

interface EventSubscription {
  name: string;
  query: string;
}

interface EventHandler {
  (data: any): Promise<void>;
}

class EventSubscriber {
  private ws: WebSocket | null = null;
  private running = false;
  private requestId = 0;
  private handlers: Map<string, EventHandler> = new Map();
  private subscriptions: Map<string, number> = new Map();

  constructor(
    private wsUrl: string,
    private agentAddress: string
  ) {}

  registerHandler(eventType: string, handler: EventHandler): void {
    this.handlers.set(eventType, handler);
    console.log(`✓ 已注册处理器: ${eventType}`);
  }

  async connect(): Promise<void> {
    this.running = true;

    const connectLoop = async () => {
      while (this.running) {
        try {
          await this._connect();
        } catch (error) {
          console.log(`⚠ WebSocket 断开,3 秒后重连...`);
          await new Promise(resolve => setTimeout(resolve, 3000));
        }
      }
    };

    await connectLoop();
  }

  private _connect(): Promise<void> {
    return new Promise((resolve, reject) => {
      this.ws = new WebSocket(this.wsUrl);

      this.ws.on('open', async () => {
        console.log(`✓ WebSocket 已连接: ${this.wsUrl}`);

        const subs: EventSubscription[] = [
          { name: 'new_block', query: "tm.event='NewBlock'" },
          { name: 'new_tx', query: "tm.event='Tx'" },
          {
            name: 'agent_registered',
            query: `wasm._contract_address='${process.env.AGENT_REGISTRY_ADDRESS}'`,
          },
          {
            name: 'payment_received',
            query: `wasm._contract_address='${process.env.AIPAY_MODULE_ADDRESS}'`,
          },
        ];

        for (const sub of subs) {
          this._subscribe(sub);
        }

        resolve();
      });

      this.ws.on('message', (data: WebSocket.Data) => {
        try {
          const message = JSON.parse(data.toString());
          this._handleMessage(message);
        } catch (error) {
          console.error('消息解析错误:', error);
        }
      });

      this.ws.on('close', () => {
        console.log('⚠ WebSocket 连接关闭');
        if (this.running) {
          setTimeout(() => {
            this._connect().catch(reject);
          }, 3000);
        }
      });

      this.ws.on('error', (error) => {
        console.error('WebSocket 错误:', error);
        reject(error);
      });
    });
  }

  private _subscribe(subscription: EventSubscription): void {
    this.requestId++;
    const subscribeMsg = {
      jsonrpc: '2.0',
      method: 'subscribe',
      id: String(this.requestId),
      params: { query: subscription.query },
    };
    this.ws?.send(JSON.stringify(subscribeMsg));
    this.subscriptions.set(subscription.name, this.requestId);
    console.log(`  ✓ 已订阅: ${subscription.name}`);
  }

  private _handleMessage(message: any): void {
    try {
      const result = message?.result;
      const data = result?.data?.value;
      if (!data) return;

      const eventType = this._detectEventType(data);
      const handler = this.handlers.get(eventType);
      if (handler) {
        handler(data).catch(console.error);
      } else {
        this._logEvent(data);
      }
    } catch (error) {
      console.error('消息处理错误:', error);
    }
  }

  private _detectEventType(value: any): string | null {
    const events = value?.events || [];
    for (const event of events) {
      const attributes = event?.attributes || [];
      for (const attr of attributes) {
        if (attr.key === 'action') {
          return attr.value;
        }
      }
    }
    return null;
  }

  private _logEvent(value: any): void {
    const logEntry = {
      timestamp: new Date().toISOString(),
      data: value,
    };
    fs.appendFileSync('events.log', JSON.stringify(logEntry) + '\n');
  }

  stop(): void {
    this.running = false;
    this.ws?.close();
  }
}

// ==================== 使用示例 ====================

async function main() {
  const subscriber = new EventSubscriber(
    process.env.MSG_WEBSOCKET_URL || 'wss://rpc.msgchain.org:443/websocket',
    process.env.AGENT_BOT_ADDR || ''
  );

  subscriber.registerHandler('new_block', async (data) => {
    const height = data?.header?.height || 'unknown';
    console.log(`📦 新区块: #${height}`);
  });

  subscriber.registerHandler('payment_received', async (data) => {
    console.log(`💰 收到支付!`);
  });

  console.log('🚀 事件订阅器启动中...');
  await subscriber.connect();
}

main().catch(console.error);

8.4 事件驱动的自动响应

# Python: 事件驱动的自动响应
class AutoResponder:
    def __init__(self, subscriber: EventSubscriber, strategy_engine):
        self.subscriber = subscriber
        self.strategy_engine = strategy_engine

    async def start(self):
        # 注册事件处理器到策略引擎
        @self.subscriber.register_handler("new_block")
        async def on_new_block(data):
            """每个新区块触发策略评估"""
            height = data.get("header", {}).get("height")
            context = {
                "block_height": height,
                "timestamp": datetime.utcnow().isoformat()
            }

            # 检查是否有计划任务
            result = await self.strategy_engine.evaluate(
                "on_schedule", context
            )
            if result:
                await self.strategy_engine.execute_action(
                    result["action"], result["params"]
                )

        @self.subscriber.register_handler("price_update")
        async def on_price_update(data):
            """价格更新时触发市场策略"""
            context = {
                "price_change": float(data.get("change_pct", 0)),
                "pair": data.get("pair", ""),
                "current_price": float(data.get("price", 0))
            }

            result = await self.strategy_engine.evaluate(
                "on_market_alert", context
            )
            if result:
                await self.strategy_engine.execute_action(
                    result["action"], result["params"]
                )

        # 启动订阅
        await self.subscriber.connect()

9. 第七步:集成 MPC 钱包与 AIPAY 支付

MPC(Multi-Party Computation)钱包为 Agent 提供多重签名安全保障。AIPAY 是 MSG Chain 上的 AI Agent 支付协议。

9.1 MPC 钱包架构

┌─────────────────────────────────────────────────────────────┐
│                     MPC 钱包架构                              │
├─────────────────────────────────────────────────────────────┤
│                                                             │
│  ┌──────────┐    ┌──────────┐    ┌──────────┐              │
│  │ 参与方 A  │    │ 参与方 B  │    │ 参与方 C  │              │
│  │ (用户)    │    │ (Agent)  │    │ (仲裁节点)│              │
│  └─────┬────┘    └─────┬────┘    └─────┬────┘              │
│        │               │               │                    │
│  ┌─────▼───────────────▼───────────────▼─────┐              │
│  │          MPC 协议层 (GG20/CMP)             │              │
│  │   密钥分片 | 分布式签名 | 多方计算          │              │
│  └───────────────────┬───────────────────────┘              │
│                      │                                      │
│  ┌───────────────────▼───────────────────────┐              │
│  │          MSG Chain 合约层                  │              │
│  │   MPC Wallet Factory | Wallet Proxy       │              │
│  └───────────────────────────────────────────┘              │
│                                                             │
└─────────────────────────────────────────────────────────────┘

9.2 创建 MPC 钱包

# Python: 创建 MPC 钱包
import asyncio
from msg_chain_sdk import Client, Wallet
from msg_chain_sdk.mpc import MPCClient, MPCParticipant

async def create_mpc_wallet():
    client = Client(rpc_url=os.getenv("MSG_RPC_URL"))
    wallet = Wallet.from_mnemonic(os.getenv("AGENT_MNEMONIC"))

    # 定义 MPC 参与方
    participants = [
        MPCParticipant(
            address=wallet.address,
            weight=2,  # 用户权重 2
            role="owner"
        ),
        MPCParticipant(
            address="msg1arbitrator...",
            weight=1,  # 仲裁节点权重 1
            role="arbitrator"
        )
    ]

    # 创建 MPC 钱包
    create_msg = {
        "create_mpc_wallet": {
            "participants": [p.model_dump() for p in participants],
            "threshold": 2,  # 需要 2/3 签名
            "agent_id": wallet.address
        }
    }

    result = await client.execute_contract(
        sender=wallet,
        contract_address=os.getenv("MPC_FACTORY_ADDRESS"),
        msg=create_msg,
        funds=[],
        gas_limit=500000
    )

    mpc_wallet_address = result.wallet_address
    print(f"✓ MPC 钱包创建成功")
    print(f"  MPC 钱包地址: {mpc_wallet_address}")
    print(f"  参与方: {len(participants)}")
    print(f"  签名阈值: 2/{len(participants)}")

    await client.close()
    return mpc_wallet_address

asyncio.run(create_mpc_wallet())

9.3 发起 MPC 签名交易

# Python: MPC 签名交易
async def propose_mpc_transaction():
    client = Client(rpc_url=os.getenv("MSG_RPC_URL"))
    wallet = Wallet.from_mnemonic(os.getenv("AGENT_MNEMONIC"))
    mpc_client = MPCClient(endpoint="https://mpc.msgchain.org:9443")

    # 构造交易
    tx = {
        "to": os.getenv("AIPAY_MODULE_ADDRESS"),
        "value": "50000",
        "denom": "umsg",
        "data": {
            "aipay": {
                "payment": {
                    "receiver": "msg1receiver...",
                    "amount": "50000",
                    "denom": "umsg",
                    "session_id": "session_123"
                }
            }
        }
    }

    # 发起 MPC 签名请求
    proposal = await mpc_client.create_proposal(
        wallet_address=os.getenv("MPC_WALLET_ADDRESS"),
        transaction=tx,
        proposer=wallet.address,
        description="AIPAY payment for data analysis"
    )

    print(f"✓ MPC 签名请求已创建")
    print(f"  提案 ID: {proposal.proposal_id}")
    print(f"  需要签名数: {proposal.required_signatures}")

    # 等待其他参与方签名(模拟)
    await asyncio.sleep(2)

    # 收集签名并广播
    signatures = await mpc_client.collect_signatures(proposal.proposal_id)
    if len(signatures) >= proposal.required_signatures:
        tx_hash = await client.broadcast_tx(signatures.signed_tx)
        print(f"✓ 交易已广播: {tx_hash}")
    else:
        print(f"⚠ 签名不足: {len(signatures)}/{proposal.required_signatures}")

    await client.close()
    return proposal.proposal_id

asyncio.run(propose_mpc_transaction())

9.4 AIPAY 支付集成

# Python: AIPAY 支付系统
from msg_chain_sdk.aipay import AIPAYClient, PaymentSession, PaymentStatus

class AIPAYService:
    def __init__(self, client: Client, wallet: Wallet):
        self.client = client
        self.wallet = wallet
        self.aipay_client = AIPAYClient(
            client=client,
            module_address=os.getenv("AIPAY_MODULE_ADDRESS")
        )
        self.active_sessions: Dict[str, PaymentSession] = {}

    async def create_payment_session(
        self,
        receiver: str,
        amount: int,
        denom: str = "umsg",
        metadata: Optional[Dict] = None
    ) -> PaymentSession:
        """创建支付会话"""
        session = await self.aipay_client.create_session(
            sender=self.wallet.address,
            receiver=receiver,
            amount=str(amount),
            denom=denom,
            metadata=metadata or {}
        )

        self.active_sessions[session.session_id] = session
        print(f"💰 支付会话已创建: {session.session_id}")
        print(f"  发送方: {session.sender}")
        print(f"  接收方: {session.receiver}")
        print(f"  金额: {session.amount} {session.denom}")

        return session

    async def execute_payment(self, session_id: str) -> str:
        """执行支付"""
        session = self.active_sessions.get(session_id)
        if not session:
            raise ValueError(f"会话不存在: {session_id}")

        # 执行支付(如果是小额直接支付,否则需要 MPC)
        if int(session.amount) < 1000000:  # < 1 MSG
            result = await self.aipay_client.direct_pay(
                wallet=self.wallet,
                session=session
            )
        else:
            # 大额支付需要通过 MPC
            result = await self.aipay_client.mpc_pay(
                wallet=self.wallet,
                session=session,
                mpc_endpoint=os.getenv("MPC_ENDPOINT")
            )

        print(f"✓ 支付成功: {result.tx_hash}")
        return result.tx_hash

    async def verify_payment(self, session_id: str) -> PaymentStatus:
        """验证支付状态"""
        status = await self.aipay_client.get_session_status(session_id)
        print(f"支付状态: {status.status.value}")
        print(f"  确认数: {status.confirmations}")
        return status

    async def refund_payment(self, session_id: str) -> str:
        """退款"""
        result = await self.aipay_client.refund(
            wallet=self.wallet,
            session_id=session_id
        )
        print(f"✓ 退款成功: {result.tx_hash}")
        return result.tx_hash

# 使用示例
async def aipay_example():
    client = Client(rpc_url=os.getenv("MSG_RPC_URL"))
    wallet = Wallet.from_mnemonic(os.getenv("AGENT_MNEMONIC"))
    aipay = AIPAYService(client, wallet)

    # Agent A 向 Agent B 支付分析费用
    session = await aipay.create_payment_session(
        receiver="msg1agent_b_address",
        amount=50000,
        denom="umsg",
        metadata={"service": "market_analysis", "request_id": "req_001"}
    )

    tx_hash = await aipay.execute_payment(session.session_id)
    status = await aipay.verify_payment(session.session_id)

    await client.close()

asyncio.run(aipay_example())
// TypeScript: AIPAY 支付集成
import { AIPAYClient, PaymentSession, PaymentStatus } from '@msgchain/sdk/aipay';

class AIPAYService {
  private aipayClient: AIPAYClient;
  private activeSessions: Map<string, PaymentSession> = new Map();

  constructor(
    private client: MsgChainClient,
    private wallet: Wallet
  ) {
    this.aipayClient = new AIPAYClient({
      client,
      moduleAddress: process.env.AIPAY_MODULE_ADDRESS!,
    });
  }

  async createPaymentSession(
    receiver: string,
    amount: number,
    denom: string = 'umsg',
    metadata?: Record<string, any>
  ): Promise<PaymentSession> {
    const session = await this.aipayClient.createSession({
      sender: this.wallet.address,
      receiver,
      amount: String(amount),
      denom,
      metadata: metadata || {},
    });

    this.activeSessions.set(session.sessionId, session);
    console.log(`💰 支付会话已创建: ${session.sessionId}`);
    return session;
  }

  async executePayment(sessionId: string): Promise<string> {
    const session = this.activeSessions.get(sessionId);
    if (!session) throw new Error(`会话不存在: ${sessionId}`);

    const result = Number(session.amount) < 1_000_000
      ? await this.aipayClient.directPay(this.wallet, session)
      : await this.aipayClient.mpcPay(this.wallet, session, process.env.MPC_ENDPOINT!);

    console.log(`✓ 支付成功: ${result.txHash}`);
    return result.txHash;
  }

  async verifyPayment(sessionId: string): Promise<PaymentStatus> {
    return this.aipayClient.getSessionStatus(sessionId);
  }
}

9.5 AIPAY 按使用量计费

# Python: AIPAY 按能力计费
class AIPAYMeteredBilling:
    def __init__(self, aipay_service: AIPAYService):
        self.aipay = aipay_service
        self.usage_records: Dict[str, List[Dict]] = {}

    async def charge_for_execution(
        self,
        capability: str,
        user_address: str,
        pricing: Dict
    ) -> Optional[str]:
        """按执行次数计费"""
        amount = int(pricing.get("amount", 0))
        if amount == 0:
            return None

        # 记录使用量
        if user_address not in self.usage_records:
            self.usage_records[user_address] = []
        self.usage_records[user_address].append({
            "capability": capability,
            "amount": amount,
            "timestamp": datetime.utcnow().isoformat()
        })

        # 创建并执行支付
        session = await self.aipay.create_payment_session(
            receiver=self.aipay.wallet.address,
            amount=amount,
            denom=pricing.get("denom", "umsg"),
            metadata={"capability": capability, "user": user_address}
        )

        return await self.aipay.execute_payment(session.session_id)

    async def get_user_billing(self, user_address: str) -> Dict:
        """查询用户账单"""
        records = self.usage_records.get(user_address, [])
        total = sum(r["amount"] for r in records)
        return {
            "user": user_address,
            "total_charged": total,
            "total_requests": len(records),
            "records": records[-10:]  # 最近 10 条
        }

10. 第八步:Agent 间通信(A2A)

A2A(Agent-to-Agent)协议允许 MSG Chain 上的 AI Agent 直接通信和协作。

10.1 A2A 协议消息格式

{
  "a2a_version": "1.0.0",
  "message_id": "msg_a2a_001",
  "source_agent": "msg1agent_a_address",
  "target_agent": "msg1agent_b_address",
  "message_type": "request",
  "timestamp": "2026-07-06T00:00:00Z",
  "ttl_seconds": 300,
  "conversation_id": "conv_001",
  "payload": {
    "method": "query",
    "params": {
      "capability": "market_analysis",
      "query": "What is the current MSG/USDC price?"
    }
  },
  "signature": {
    "algorithm": "Dilithium5",
    "value": "base64_encoded_signature..."
  }
}

10.2 Python A2A 服务端

# Python: A2A 通信处理器 (a2a_handler.py)
import asyncio
import json
import os
import uuid
from datetime import datetime
from typing import Optional, Dict, Any, Callable
from fastapi import APIRouter, HTTPException, Request
from pydantic import BaseModel, Field

from msg_chain_sdk import Client, Wallet
from msg_chain_sdk.crypto import Dilithium5, verify_message

# ==================== A2A 数据模型 ====================

class A2ARequest(BaseModel):
    a2a_version: str = "1.0.0"
    message_id: str = Field(default_factory=lambda: f"msg_{uuid.uuid4().hex[:12]}")
    source_agent: str
    target_agent: str
    message_type: str = "request"
    timestamp: str = Field(default_factory=lambda: datetime.utcnow().isoformat())
    ttl_seconds: int = 300
    conversation_id: Optional[str] = None
    payload: Dict[str, Any]
    signature: Optional[Dict[str, str]] = None

class A2AResponse(BaseModel):
    a2a_version: str = "1.0.0"
    message_id: str = Field(default_factory=lambda: f"resp_{uuid.uuid4().hex[:12]}")
    source_agent: str
    target_agent: str
    message_type: str = "response"
    timestamp: str = Field(default_factory=lambda: datetime.utcnow().isoformat())
    conversation_id: str
    payload: Dict[str, Any]
    signature: Optional[Dict[str, str]] = None

# ==================== A2A 路由器 ====================

class A2ARouter:
    def __init__(self, agent_address: str, dilithium_key: Dilithium5):
        self.agent_address = agent_address
        self.dilithium_key = dilithium_key
        self.method_handlers: Dict[str, Callable] = {}
        self.conversations: Dict[str, list] = {}
        self.router = APIRouter(prefix="/a2a", tags=["a2a"])
        self._setup_routes()

    def register_method(self, method: str, handler: Callable):
        """注册 A2A 方法处理器"""
        self.method_handlers[method] = handler
        print(f"✓ A2A 方法已注册: {method}")

    def _setup_routes(self):
        @self.router.post("/query")
        async def handle_a2a_query(req: A2ARequest):
            return await self._process_request(req)

        @self.router.get("/health")
        async def a2a_health():
            return {
                "agent_id": self.agent_address,
                "a2a_version": "1.0.0",
                "methods": list(self.method_handlers.keys())
            }

    async def _process_request(self, req: A2ARequest) -> A2AResponse:
        """处理 A2A 请求"""
        # 验证 TTL
        req_time = datetime.fromisoformat(req.timestamp)
        if (datetime.utcnow() - req_time).total_seconds() > req.ttl_seconds:
            raise HTTPException(status_code=408, detail="请求已过期")

        # 验证源 Agent(可选,生产环境必需)
        if req.signature:
            valid = await self._verify_signature(req)
            if not valid:
                raise HTTPException(status_code=401, detail="签名验证失败")

        # 记录会话
        conv_id = req.conversation_id or f"conv_{uuid.uuid4().hex[:12]}"
        if conv_id not in self.conversations:
            self.conversations[conv_id] = []
        self.conversations[conv_id].append(req.model_dump())

        # 执行请求
        method = req.payload.get("method")
        params = req.payload.get("params", {})

        handler = self.method_handlers.get(method)
        if not handler:
            raise HTTPException(status_code=400, detail=f"未知方法: {method}")

        try:
            result = await handler(params)

            response = A2AResponse(
                source_agent=self.agent_address,
                target_agent=req.source_agent,
                conversation_id=conv_id,
                payload={"result": result}
            )

            # 签名响应
            response.signature = await self._sign_message(response.model_dump())

            # 记录响应
            self.conversations[conv_id].append(response.model_dump())

            return response

        except Exception as e:
            return A2AResponse(
                source_agent=self.agent_address,
                target_agent=req.source_agent,
                conversation_id=conv_id,
                payload={"error": str(e)},
                message_type="error"
            )

    async def _verify_signature(self, req: A2ARequest) -> bool:
        """验证 A2A 消息签名"""
        # 查询源 Agent 的公钥
        client = Client(rpc_url=os.getenv("MSG_RPC_URL"))
        did_doc = await client.query_contract(
            os.getenv("DID_REGISTRY_ADDRESS"),
            {"query_did": {"did": f"did:msg:{req.source_agent}"}}
        )

        # 提取 Dilithium-5 公钥
        for vm in did_doc.get("verificationMethod", []):
            if vm["type"] == "Dilithium5VerificationKey2026":
                pub_key = vm["publicKeyMultibase"]
                break
        else:
            return False

        # 验证签名
        message = json.dumps(req.model_dump(exclude={"signature"}), sort_keys=True)
        return verify_message(message, req.signature["value"], pub_key)

    async def _sign_message(self, message: Dict) -> Dict:
        """签名消息"""
        message_str = json.dumps(message, sort_keys=True)
        signature = self.dilithium_key.sign(message_str.encode())
        return {
            "algorithm": "Dilithium5",
            "value": signature.hex()
        }

    def get_router(self):
        return self.router

# ==================== A2A 客户端 ====================

class A2AClient:
    def __init__(self, agent_address: str, dilithium_key: Dilithium5):
        self.agent_address = agent_address
        self.dilithium_key = dilithium_key
        self.http_client = httpx.AsyncClient(timeout=30)

    async def send_request(
        self,
        target_agent_id: str,
        method: str,
        params: Dict,
        target_endpoint: str = None
    ) -> Dict:
        """向目标 Agent 发送 A2A 请求"""
        # 查询目标 Agent 的 A2A 端点
        if not target_endpoint:
            client = Client(rpc_url=os.getenv("MSG_RPC_URL"))
            agent_info = await client.query_contract(
                os.getenv("AGENT_REGISTRY_ADDRESS"),
                {"query_agent": {"agent_id": target_agent_id}}
            )
            target_endpoint = agent_info.get("endpoints", {}).get("a2a")

        if not target_endpoint:
            raise ValueError(f"目标 Agent 无 A2A 端点: {target_agent_id}")

        # 构造 A2A 请求
        request = A2ARequest(
            source_agent=self.agent_address,
            target_agent=target_agent_id,
            payload={"method": method, "params": params}
        )

        # 签名请求
        request.signature = {
            "algorithm": "Dilithium5",
            "value": self.dilithium_key.sign(
                json.dumps(request.model_dump(exclude={"signature"}), sort_keys=True).encode()
            ).hex()
        }

        # 发送请求
        response = await self.http_client.post(
            f"{target_endpoint}/query",
            json=request.model_dump()
        )
        response.raise_for_status()
        return response.json()

    async def close(self):
        await self.http_client.aclose()

11. 第九步:监控与日志

生产环境下的 Agent 需要完善的监控、日志和告警系统来确保稳定运行。

11.1 日志系统

# Python: 结构化日志 (logger.py)
import json
import logging
import sys
from datetime import datetime
from typing import Optional

class AgentLogger:
    def __init__(self, name: str, level: str = "INFO", log_file: Optional[str] = None):
        self.logger = logging.getLogger(name)
        self.logger.setLevel(getattr(logging, level.upper()))

        formatter = logging.Formatter(
            '%(message)s'  # JSON 格式,自行序列化
        )

        # 控制台输出
        console_handler = logging.StreamHandler(sys.stdout)
        console_handler.setFormatter(formatter)
        self.logger.addHandler(console_handler)

        # 文件输出
        if log_file:
            file_handler = logging.FileHandler(log_file)
            file_handler.setFormatter(formatter)
            self.logger.addHandler(file_handler)

    def _make_record(self, level: str, message: str, extra: Optional[dict] = None):
        record = {
            "timestamp": datetime.utcnow().isoformat(),
            "level": level,
            "logger": self.logger.name,
            "message": message,
            **(extra or {})
        }
        return json.dumps(record, ensure_ascii=False)

    def info(self, message: str, **kwargs):
        self.logger.info(self._make_record("INFO", message, kwargs))

    def warning(self, message: str, **kwargs):
        self.logger.warning(self._make_record("WARNING", message, kwargs))

    def error(self, message: str, **kwargs):
        self.logger.error(self._make_record("ERROR", message, kwargs))

    def critical(self, message: str, **kwargs):
        self.logger.critical(self._make_record("CRITICAL", message, kwargs))

# ==================== 使用示例 ====================
log = AgentLogger("my_agent", level="INFO", log_file="agent.log")

log.info("Agent 启动",
    agent_id="msg1agent...",
    version="1.0.0",
    port=8080
)

log.warning("Gas 费用偏高",
    tx_hash="tx_hash...",
    gas_used=450000,
    gas_limit=500000
)

log.error("合约调用失败",
    contract="did_registry",
    error="out of gas",
    block_height=123456
)
// TypeScript: 结构化日志 (src/logger.ts)
import * as fs from 'fs';
import * as path from 'path';

type LogLevel = 'INFO' | 'WARNING' | 'ERROR' | 'CRITICAL';

interface LogRecord {
  timestamp: string;
  level: LogLevel;
  logger: string;
  message: string;
  [key: string]: any;
}

class AgentLogger {
  private stream?: fs.WriteStream;

  constructor(
    private name: string,
    private level: LogLevel = 'INFO',
    logFile?: string
  ) {
    if (logFile) {
      this.stream = fs.createWriteStream(logFile, { flags: 'a' });
    }
  }

  private log(level: LogLevel, message: string, extra?: Record<string, any>): void {
    if (this._shouldLog(level)) {
      const record: LogRecord = {
        timestamp: new Date().toISOString(),
        level,
        logger: this.name,
        message,
        ...extra,
      };
      const output = JSON.stringify(record) + '\n';
      process.stdout.write(output);
      this.stream?.write(output);
    }
  }

  private _shouldLog(level: LogLevel): boolean {
    const levels: LogLevel[] = ['INFO', 'WARNING', 'ERROR', 'CRITICAL'];
    return levels.indexOf(level) >= levels.indexOf(this.level);
  }

  info(message: string, extra?: Record<string, any>): void {
    this.log('INFO', message, extra);
  }

  warning(message: string, extra?: Record<string, any>): void {
    this.log('WARNING', message, extra);
  }

  error(message: string, extra?: Record<string, any>): void {
    this.log('ERROR', message, extra);
  }

  critical(message: string, extra?: Record<string, any>): void {
    this.log('CRITICAL', message, extra);
  }

  close(): void {
    this.stream?.end();
  }
}

const log = new AgentLogger('my_agent', 'INFO', 'agent.log');
log.info('Agent 启动', { agent_id: 'msg1...', version: '1.0.0' });

11.2 Prometheus 指标

# Python: Prometheus 指标 (metrics.py)
from prometheus_client import Counter, Histogram, Gauge, start_http_server
import time
from functools import wraps

# 指标定义
REQUEST_COUNT = Counter(
    'agent_requests_total',
    'Total number of agent requests',
    ['capability', 'status']
)

REQUEST_DURATION = Histogram(
    'agent_request_duration_seconds',
    'Request duration in seconds',
    ['capability'],
    buckets=[0.1, 0.5, 1.0, 2.0, 5.0, 10.0]
)

ACTIVE_SESSIONS = Gauge(
    'agent_active_sessions',
    'Number of active sessions',
    ['type']
)

GAS_USED = Counter(
    'agent_gas_used_total',
    'Total gas used',
    ['contract']
)

BALANCE_GAUGE = Gauge(
    'agent_balance',
    'Current agent balance',
    ['denom']
)

def monitor_capability(name: str):
    """监控能力执行的装饰器"""
    def decorator(func):
        @wraps(func)
        async def wrapper(*args, **kwargs):
            start = time.time()
            try:
                result = await func(*args, **kwargs)
                REQUEST_COUNT.labels(capability=name, status="success").inc()
                return result
            except Exception as e:
                REQUEST_COUNT.labels(capability=name, status="error").inc()
                raise
            finally:
                duration = time.time() - start
                REQUEST_DURATION.labels(capability=name).observe(duration)
        return wrapper
    return decorator

# 启动指标服务器
def start_metrics_server(port: int = 9090):
    start_http_server(port)
    print(f"📊 Prometheus 指标端点: http://0.0.0.0:{port}/metrics")

11.3 健康检查与自愈

# Python: 健康检查与自动恢复 (health.py)
import asyncio
import os
import subprocess
import time
from datetime import datetime

class HealthChecker:
    def __init__(self, agent_id: str, endpoints: Dict[str, str]):
        self.agent_id = agent_id
        self.endpoints = endpoints
        self.last_healthy = datetime.utcnow()
        self.consecutive_failures = 0
        self.max_failures = 3

    async def check_all(self) -> Dict[str, bool]:
        """执行所有健康检查"""
        results = {}

        # 1. API 网关健康检查
        results["api_gateway"] = await self._check_http(
            self.endpoints.get("api", "http://localhost:8080/health")
        )

        # 2. 链节点连接
        results["chain_connection"] = await self._check_chain()

        # 3. 钱包余额
        results["balance"] = await self._check_balance()

        # 4. 事件订阅器
        results["event_subscriber"] = await self._check_websocket()

        # 5. 合约可访问性
        results["contracts"] = await self._check_contracts()

        # 更新状态
        all_healthy = all(results.values())
        if all_healthy:
            self.last_healthy = datetime.utcnow()
            self.consecutive_failures = 0
        else:
            self.consecutive_failures += 1

        return results

    async def _check_http(self, url: str) -> bool:
        """HTTP 端点检查"""
        try:
            async with httpx.AsyncClient() as client:
                response = await client.get(url, timeout=10)
                return response.status_code == 200
        except Exception:
            return False

    async def _check_chain(self) -> bool:
        """链节点连接检查"""
        try:
            client = Client(rpc_url=os.getenv("MSG_RPC_URL"))
            status = await client.get_status()
            await client.close()
            return status.sync_info is not None
        except Exception:
            return False

    async def _check_balance(self) -> bool:
        """余额检查"""
        try:
            client = Client(rpc_url=os.getenv("MSG_RPC_URL"))
            wallet = Wallet.from_mnemonic(os.getenv("AGENT_MNEMONIC"))
            balance = await client.get_balance(wallet.address, "umsg")
            await client.close()
            return int(balance) > 50000  # 至少 50000 umsg
        except Exception:
            return False

    async def _check_websocket(self) -> bool:
        """WebSocket 连接检查"""
        try:
            async with websockets.connect(
                os.getenv("MSG_WEBSOCKET_URL"),
                ping_interval=10,
                ping_timeout=5,
                close_timeout=5
            ) as ws:
                return ws.open
        except Exception:
            return False

    async def _check_contracts(self) -> bool:
        """合约可访问性检查"""
        try:
            client = Client(rpc_url=os.getenv("MSG_RPC_URL"))
            for name, addr in [
                ("DID Registry", os.getenv("DID_REGISTRY_ADDRESS")),
                ("Agent Registry", os.getenv("AGENT_REGISTRY_ADDRESS")),
                ("Constitution", os.getenv("CONSTITUTION_ADDRESS"))
            ]:
                code_id = await client.get_contract_code_id(addr)
                if not code_id:
                    return False
            await client.close()
            return True
        except Exception:
            return False

    async def auto_heal(self):
        """自动恢复"""
        log = AgentLogger("health_checker")

        if self.consecutive_failures >= self.max_failures:
            log.critical("连续失败超过阈值,执行自动恢复",
                failures=self.consecutive_failures,
                threshold=self.max_failures
            )

            # 重启 API 网关
            log.info("重启 API 网关...")
            subprocess.run(["pkill", "-f", "uvicorn"], capture_output=True)
            subprocess.Popen(
                ["uvicorn", "gateway:app", "--host", "0.0.0.0", "--port", "8080"],
                stdout=subprocess.DEVNULL,
                stderr=subprocess.DEVNULL
            )
            await asyncio.sleep(5)

            # 重启事件订阅器
            log.info("重启事件订阅器...")
            # (在实际部署中通过进程管理器如 systemd/supervisor 管理)

            self.consecutive_failures = 0

    async def monitoring_loop(self, interval: int = 60):
        """定期监控循环"""
        log = AgentLogger("monitor")

        while True:
            log.info("执行健康检查...")
            results = await self.check_all()

            for name, healthy in results.items():
                status = "✓" if healthy else "✗"
                level = log.info if healthy else log.warning
                level(f"  {status} {name}")

            await self.auto_heal()

            # 更新 Prometheus 指标
            BALANCE_GAUGE.labels(denom="umsg").set(
                await self._get_balance_value()
            )

            await asyncio.sleep(interval)

    async def _get_balance_value(self) -> float:
        try:
            client = Client(rpc_url=os.getenv("MSG_RPC_URL"))
            wallet = Wallet.from_mnemonic(os.getenv("AGENT_MNEMONIC"))
            balance = await client.get_balance(wallet.address, "umsg")
            await client.close()
            return float(balance)
        except Exception:
            return 0.0

11.4 监控面板配置

# docker-compose.monitoring.yml
version: '3.8'

services:
  prometheus:
    image: prom/prometheus:latest
    volumes:
      - ./prometheus.yml:/etc/prometheus/prometheus.yml
    ports:
      - "9090:9090"

  grafana:
    image: grafana/grafana:latest
    ports:
      - "3000:3000"
    environment:
      - GF_SECURITY_ADMIN_PASSWORD=admin
    volumes:
      - grafana_data:/var/lib/grafana

volumes:
  grafana_data:
# prometheus.yml
global:
  scrape_interval: 15s

scrape_configs:
  - job_name: 'ai_agent'
    static_configs:
      - targets: ['host.docker.internal:9090']
    metrics_path: '/metrics'

12. 第十步:一键启动脚本

完整的 Agent 一键启动脚本,整合上述所有步骤。

12.1 启动脚本

#!/usr/bin/env bash
# =============================================================================
# MSG Chain AI Agent 一键启动脚本
# 版本: 1.0.0
# 描述: 自动化完成 Agent 的完整启动流程
# 使用方法:
#   ./start_agent.sh                          # 交互式启动
#   ./start_agent.sh --config config.json     # 从配置文件启动
#   ./start_agent.sh --quick                  # 快速启动(使用默认值)
# =============================================================================

set -euo pipefail

# ==================== 颜色定义 ====================
RED='\033[0;31m'
GREEN='\033[0;32m'
YELLOW='\033[1;33m'
BLUE='\033[0;34m'
NC='\033[0m'

info()  { echo -e "${BLUE}[INFO]${NC} $1"; }
ok()    { echo -e "${GREEN}[✓]${NC} $1"; }
warn()  { echo -e "${YELLOW}[⚠]${NC} $1"; }
error() { echo -e "${RED}[✗]${NC} $1"; }

# ==================== 配置 ====================

CONFIG_FILE="${CONFIG_FILE:-agent_config.json}"
LOG_DIR="./logs"
KEYS_DIR="./keys"
PID_FILE="./agent.pid"

# 默认配置
MSG_CHAIN_ID="msg-chain-1"
MSG_RPC_URL="https://rpc.msgchain.org:443"
MSG_REST_URL="https://rest.msgchain.org:443"
MSG_WS_URL="wss://rpc.msgchain.org:443/websocket"
AGENT_PORT=8080
AGENT_NAME="my_ai_agent"
LOG_LEVEL="INFO"

# ==================== 函数 ====================

print_banner() {
    echo ""
    echo "╔═══════════════════════════════════════════════╗"
    echo "║      MSG Chain AI Agent 启动器 v1.0.0        ║"
    echo "║      https://msgchain.org/agents             ║"
    echo "╚═══════════════════════════════════════════════╝"
    echo ""
}

check_dependencies() {
    local missing=0
    local deps=("msgd" "python3" "uvicorn" "wscat")

    for dep in "${deps[@]}"; do
        if ! command -v "$dep" &> /dev/null; then
            warn "缺少依赖: $dep"
            missing=1
        fi
    done

    if [ $missing -eq 1 ]; then
        info "请安装缺少的依赖后重试"
        info "参考: docs/AI Agent 引导启动流程指南.md 第零步"
        exit 1
    fi
    ok "所有依赖已安装"
}

load_config() {
    if [ -f "$CONFIG_FILE" ]; then
        info "从 $CONFIG_FILE 加载配置..."
        MSG_RPC_URL=$(python3 -c "import json; print(json.load(open('$CONFIG_FILE'))['rpc_url'])")
        AGENT_NAME=$(python3 -c "import json; print(json.load(open('$CONFIG_FILE'))['agent_name'])")
        AGENT_PORT=$(python3 -c "import json; print(json.load(open('$CONFIG_FILE'))['port'])")
        ok "配置已加载"
    fi
}

setup_directories() {
    mkdir -p "$LOG_DIR" "$KEYS_DIR"
    ok "目录已创建: $LOG_DIR $KEYS_DIR"
}

check_chain_connection() {
    info "检查链节点连接..."
    if msgd status --node "$MSG_RPC_URL" &>/dev/null; then
        local block_height=$(msgd status --node "$MSG_RPC_URL" 2>/dev/null | \
            python3 -c "import sys,json; print(json.load(sys.stdin)['sync_info']['latest_block_height'])")
        ok "节点已连接 (最新区块: $block_height)"
    else
        error "无法连接到节点: $MSG_RPC_URL"
        exit 1
    fi
}

check_or_create_wallet() {
    if msgd keys show agent_owner &>/dev/null; then
        ok "钱包已存在: $(msgd keys show agent_owner -a)"
    else
        info "创建新钱包..."
        msgd keys add agent_owner
        warn "请安全保存助记词!"
        echo "按回车键继续..."
        read -r
    fi

    # 检查余额
    local balance=$(msgd query bank balances "$(msgd keys show agent_owner -a)" \
        --node "$MSG_RPC_URL" 2>/dev/null | \
        python3 -c "import sys,json; coins=json.load(sys.stdin).get('balances',[]); print([c['amount'] for c in coins if c['denom']=='umsg'][0] if any(c['denom']=='umsg' for c in coins) else '0')" 2>/dev/null || echo "0")

    if [ "$balance" = "0" ]; then
        warn "钱包余额为 0,请先充值"
        info "水龙头: https://faucet.msgchain.org/claim"
        info "或转账至: $(msgd keys show agent_owner -a)"
        echo "按回车键继续(如果已充值)..."
        read -r
    fi
    ok "钱包余额: $balance umsg"
}

generate_keys() {
    if [ ! -f "$KEYS_DIR/dilithium_priv.bin" ]; then
        info "生成 Dilithium-5 后量子密钥..."
        python3 -c "
from msg_chain_sdk.crypto import Dilithium5
import os

dk = Dilithium5.generate()
os.makedirs('$KEYS_DIR', exist_ok=True)
dk.save_private_key('$KEYS_DIR/dilithium_priv.bin')
dk.save_public_key('$KEYS_DIR/dilithium_pub.bin')
print(f'公钥: {dk.public_key.hex()[:64]}...')
"
        ok "Dilithium-5 密钥已生成"
    else
        ok "Dilithium-5 密钥已存在"
    fi
}

create_did() {
    info "检查 DID 是否已创建..."
    local agent_addr=$(msgd keys show agent_owner -a)
    local did="did:msg:${agent_addr}"

    if msgd query wasm contract-state smart "$DID_REGISTRY_ADDRESS" \
        "{\"query_did\":{\"did\":\"${did}\"}}" \
        --node "$MSG_RPC_URL" &>/dev/null; then
        ok "DID 已存在: $did"
    else
        info "创建 DID..."
        python3 -c "
import asyncio, os, json
from msg_chain_sdk import Client, Wallet
from msg_chain_sdk.crypto import Dilithium5

async def main():
    client = Client(rpc_url='$MSG_RPC_URL')
    wallet = Wallet.from_mnemonic(os.getenv('AGENT_MNEMONIC'))

    dk = Dilithium5.load_private_key('$KEYS_DIR/dilithium_priv.bin')

    did_doc = {
        '@context': ['https://www.w3.org/ns/did/v1', 'https://msgchain.org/ns/did-config/v1'],
        'id': f'did:msg:{wallet.address}',
        'alsoKnownAs': ['$AGENT_NAME'],
        'controller': [f'did:msg:{wallet.address}'],
        'verificationMethod': [
            {'id': f'did:msg:{wallet.address}#keys-1', 'type': 'EcdsaSecp256k1RecoveryMethod2020',
             'controller': f'did:msg:{wallet.address}', 'blockchainAccountId': wallet.address},
            {'id': f'did:msg:{wallet.address}#dilithium-5-keys-1', 'type': 'Dilithium5VerificationKey2026',
             'controller': f'did:msg:{wallet.address}', 'publicKeyMultibase': dk.public_key_multibase()}
        ],
        'authentication': [f'did:msg:{wallet.address}#keys-1', f'did:msg:{wallet.address}#dilithium-5-keys-1'],
        'service': [{'id': f'did:msg:{wallet.address}#agent-endpoint', 'type': 'AIAgentEndpoint',
                     'serviceEndpoint': 'http://localhost:$AGENT_PORT'}]
    }

    result = await client.execute_contract(
        sender=wallet,
        contract_address=os.getenv('DID_REGISTRY_ADDRESS'),
        msg={'create_did': {'did_document': did_doc}},
        funds=[],
        gas_limit=500000
    )
    print(f'DID 创建成功: {result.transaction_hash}')

    with open('keys/did_document.json', 'w') as f:
        json.dump(did_doc, f, indent=2)
    print('DID 文档已保存')

    await client.close()

asyncio.run(main())
"
        ok "DID 创建成功"
    fi
}

register_agent() {
    info "检查 Agent 是否已注册..."
    local agent_addr=$(msgd keys show agent_owner -a)

    if msgd query wasm contract-state smart "$AGENT_REGISTRY_ADDRESS" \
        "{\"query_agent\":{\"agent_id\":\"${agent_addr}\"}}" \
        --node "$MSG_RPC_URL" &>/dev/null; then
        ok "Agent 已注册: $agent_addr"
    else
        info "注册 Agent..."
        python3 -c "
import asyncio, os, json
from msg_chain_sdk import Client, Wallet

async def main():
    client = Client(rpc_url='$MSG_RPC_URL')
    wallet = Wallet.from_mnemonic(os.getenv('AGENT_MNEMONIC'))

    register_msg = {
        'register_agent': {
            'metadata': {
                'name': '$AGENT_NAME',
                'description': 'AI Agent on MSG Chain',
                'version': '1.0.0',
                'tags': ['defi', 'analytics']
            },
            'capabilities': [
                {
                    'name': 'market_analysis',
                    'description': 'Market data analysis',
                    'input_schema': {'type': 'object', 'properties': {'pair': {'type': 'string'}},
                                     'required': ['pair']},
                    'pricing': {'model': 'per_request', 'amount': '5000', 'denom': 'umsg'}
                }
            ],
            'endpoints': {
                'api': 'http://localhost:$AGENT_PORT',
                'a2a': 'http://localhost:$AGENT_PORT/a2a',
                'ws': 'ws://localhost:$AGENT_PORT/ws'
            },
            'stake': {'amount': '1000000', 'denom': 'umsg'}
        }
    }

    result = await client.execute_contract(
        sender=wallet,
        contract_address=os.getenv('AGENT_REGISTRY_ADDRESS'),
        msg=register_msg,
        funds=[{'denom': 'umsg', 'amount': '1000000'}],
        gas_limit=500000
    )
    print(f'Agent 注册成功: {result.transaction_hash}')

    with open('agent_registration.json', 'w') as f:
        json.dump(register_msg, f, indent=2)

    await client.close()

asyncio.run(main())
"
        ok "Agent 注册成功"
    fi
}

set_constitution() {
    info "检查宪章是否已设置..."
    local agent_addr=$(msgd keys show agent_owner -a)

    if msgd query wasm contract-state smart "$CONSTITUTION_ADDRESS" \
        "{\"query_constitution\":{\"agent_id\":\"${agent_addr}\"}}" \
        --node "$MSG_RPC_URL" &>/dev/null; then
        ok "宪章已设置"
    else
        info "设置宪章..."
        python3 -c "
import asyncio, os
from msg_chain_sdk import Client, Wallet

async def main():
    client = Client(rpc_url='$MSG_RPC_URL')
    wallet = Wallet.from_mnemonic(os.getenv('AGENT_MNEMONIC'))

    constitution = {
        'name': '$AGENT_NAME Constitution',
        'version': '1.0.0',
        'purpose': 'Provide decentralized analytics and trading assistance',
        'principles': [
            {'id': 'p1', 'name': 'Safety', 'statement': 'Never cause irreversible loss', 'priority': 1},
            {'id': 'p2', 'name': 'Transparency', 'statement': 'Log all actions on-chain', 'priority': 2}
        ],
        'policies': [
            {'id': 'pol1', 'name': 'Trade Limit', 'statement': 'Max 1000 MSG per trade',
             'conditions': {'type': 'amount_limit', 'max_amount': '1000000000', 'denom': 'umsg'}}
        ],
        'strategy_engine': {
            'model': 'rule_based',
            'rules': [
                {'id': 'r1', 'trigger': 'on_market_alert', 'condition': 'price_change > 5',
                 'action': 'notify_user', 'params': {}}
            ]
        },
        'constraints': ['Maintain min 50000 umsg balance'],
        'amendments': []
    }

    result = await client.execute_contract(
        sender=wallet,
        contract_address=os.getenv('CONSTITUTION_ADDRESS'),
        msg={'set_constitution': {'constitution': constitution}},
        funds=[],
        gas_limit=600000
    )
    print(f'宪章设置成功: {result.transaction_hash}')

    await client.close()

asyncio.run(main())
"
        ok "宪章设置成功"
    fi
}

start_api_gateway() {
    info "启动 API 网关 (端口: $AGENT_PORT)..."

    # 检查是否已在运行
    if [ -f "$PID_FILE" ]; then
        local old_pid=$(cat "$PID_FILE")
        if kill -0 "$old_pid" 2>/dev/null; then
            ok "API 网关已在运行 (PID: $old_pid)"
            return
        fi
    fi

    nohup python3 gateway.py > "$LOG_DIR/gateway.log" 2>&1 &
    local pid=$!
    echo "$pid" > "$PID_FILE"
    sleep 3

    if kill -0 "$pid" 2>/dev/null; then
        ok "API 网关已启动 (PID: $pid)"
    else
        error "API 网关启动失败,检查日志: $LOG_DIR/gateway.log"
        exit 1
    fi
}

start_event_subscriber() {
    info "启动事件订阅器..."
    nohup python3 -c "
import asyncio
from event_subscriber import EventSubscriber

async def main():
    sub = EventSubscriber(
        ws_url='$MSG_WS_URL',
        agent_address='$(msgd keys show agent_owner -a)',
        rpc_url='$MSG_RPC_URL'
    )
    await sub.connect()

asyncio.run(main())
" > "$LOG_DIR/subscriber.log" 2>&1 &

    local pid=$!
    echo "$pid" > "$PID_FILE.subscriber"
    ok "事件订阅器已启动 (PID: $pid)"
}

start_monitoring() {
    info "启动监控服务..."

    # 启动 Prometheus 指标端点
    nohup python3 -c "
from prometheus_client import start_http_server
start_http_server(9090)
import time
while True:
    time.sleep(1)
" > "$LOG_DIR/metrics.log" 2>&1 &

    local pid=$!
    echo "$pid" > "$PID_FILE.metrics"
    ok "监控服务已启动 (PID: $pid)"
}

print_summary() {
    local agent_addr=$(msgd keys show agent_owner -a)

    echo ""
    echo "╔═══════════════════════════════════════════════╗"
    echo "║        AI Agent 启动完成!                    ║"
    echo "╠═══════════════════════════════════════════════╣"
    echo "║  Agent ID: ${agent_addr:0:20}...            "
    echo "║  DID:      did:msg:${agent_addr:0:20}...    "
    echo "║  API:      http://localhost:$AGENT_PORT      "
    echo "║  A2A:      http://localhost:$AGENT_PORT/a2a  "
    echo "║  WS:       ws://localhost:$AGENT_PORT/ws     "
    echo "║  Metrics:  http://localhost:9090/metrics     "
    echo "║                                              "
    echo "║  日志:     $LOG_DIR/                          "
    echo "║  密钥:     $KEYS_DIR/                         "
    echo "╚═══════════════════════════════════════════════╝"
    echo ""
}

# ==================== 主流程 ====================

main() {
    print_banner

    # 加载配置
    load_config

    # 导出环境变量
    export MSG_CHAIN_ID
    export MSG_RPC_URL
    export MSG_REST_URL
    export MSG_WEBSOCKET_URL="$MSG_WS_URL"
    export AGENT_PORT
    export AGENT_NAME
    export LOG_LEVEL

    # 步骤 0: 环境检查
    info "步骤 0/10: 环境准备"
    check_dependencies
    setup_directories

    # 步骤 0.5: 链连接
    info "步骤 0.5/10: 链连接"
    check_chain_connection
    check_or_create_wallet

    # 步骤 1: DID
    info "步骤 1/10: 创建 DID"
    generate_keys
    create_did

    # 步骤 2: 注册
    info "步骤 2/10: 注册 Agent"
    register_agent

    # 步骤 3: 宪章
    info "步骤 3/10: 设置宪章"
    set_constitution

    # 步骤 4: 启动服务
    info "步骤 4/10: 启动服务"
    start_api_gateway
    start_event_subscriber
    start_monitoring

    # 完成
    print_summary
}

# ==================== 入口 ====================

case "${1:-}" in
    --config)
        CONFIG_FILE="$2"
        main
        ;;
    --quick)
        CONFIG_FILE=""
        main
        ;;
    --stop)
        info "停止 Agent..."
        for pid_file in "$PID_FILE" "$PID_FILE.subscriber" "$PID_FILE.metrics"; do
            if [ -f "$pid_file" ]; then
                kill "$(cat "$pid_file")" 2>/dev/null || true
                rm -f "$pid_file"
            fi
        done
        ok "Agent 已停止"
        ;;
    --status)
        info "Agent 状态:"
        for pid_file in "$PID_FILE" "$PID_FILE.subscriber" "$PID_FILE.metrics"; do
            if [ -f "$pid_file" ] && kill -0 "$(cat "$pid_file")" 2>/dev/null; then
                ok "  $(basename "$pid_file"): 运行中 (PID: $(cat "$pid_file"))"
            else
                warn "  $(basename "$pid_file"): 未运行"
            fi
        done
        ;;
    --help|-h)
        echo "使用方法: $0 [选项]"
        echo ""
        echo "选项:"
        echo "  --config <文件>  从配置文件启动"
        echo "  --quick          快速启动(默认配置)"
        echo "  --stop           停止 Agent"
        echo "  --status         查看运行状态"
        echo "  --help, -h       显示帮助"
        ;;
    *)
        main
        ;;
esac

保存脚本并赋予执行权限:

chmod +x start_agent.sh

12.2 配置文件示例

{
  "chain_id": "msg-chain-1",
  "rpc_url": "https://rpc.msgchain.org:443",
  "rest_url": "https://rest.msgchain.org:443",
  "ws_url": "wss://rpc.msgchain.org:443/websocket",
  "agent_name": "my_ai_agent",
  "port": 8080,
  "log_level": "INFO",
  "min_balance": "50000",
  "stake_amount": "1000000",
  "auto_heal": true,
  "health_check_interval": 60,
  "metrics_port": 9090
}

12.3 使用 Makefile 管理

# Makefile
.PHONY: help setup start stop restart status clean

help:
	@echo "MSG Chain AI Agent 管理"
	@echo ""
	@echo "用法:"
	@echo "  make setup     初始化环境(安装依赖)"
	@echo "  make start     启动 Agent"
	@echo "  make stop      停止 Agent"
	@echo "  make restart   重启 Agent"
	@echo "  make status    查看状态"
	@echo "  make clean     清理日志和数据"

setup:
	python3 -m venv agent-env
	. agent-env/bin/activate && pip install -r requirements.txt
	npm install
	chmod +x start_agent.sh

start:
	./start_agent.sh --quick

stop:
	./start_agent.sh --stop

restart: stop start

status:
	./start_agent.sh --status

clean:
	rm -rf logs/ keys/ *.pid __pycache__/ .env

12.4 Docker 部署

# Dockerfile
FROM python:3.11-slim

WORKDIR /app

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

# 安装 msgd
COPY --from=msgchain/msgd:latest /usr/local/bin/msgd /usr/local/bin/msgd

# 复制应用
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt

COPY . .

# 创建目录
RUN mkdir -p logs keys

# 暴露端口
EXPOSE 8080 9090

# 启动命令
CMD ["./start_agent.sh", "--quick"]
# docker-compose.yml
version: '3.8'

services:
  agent:
    build: .
    ports:
      - "8080:8080"
      - "9090:9090"
    volumes:
      - ./keys:/app/keys
      - ./logs:/app/logs
      - ./.env:/app/.env
    environment:
      - AGENT_PORT=8080
      - LOG_LEVEL=INFO
    restart: unless-stopped

  prometheus:
    image: prom/prometheus:latest
    ports:
      - "9091:9090"
    volumes:
      - ./prometheus.yml:/etc/prometheus/prometheus.yml
    restart: unless-stopped

  grafana:
    image: grafana/grafana:latest
    ports:
      - "3000:3000"
    environment:
      - GF_SECURITY_ADMIN_PASSWORD=admin
    restart: unless-stopped

13. 附录

附录 A:Agent 启动状态机

                  ┌──────────┐
                  │  初始状态  │
                  └────┬─────┘
                       │
                  ┌────▼─────┐
                  │ 环境检查  │ ◄──── 失败重试
                  └────┬─────┘
                       │
                  ┌────▼─────┐
                  │  DID 创建  │
                  └────┬─────┘
                       │
                  ┌────▼─────┐
                  │ 钱包充值   │ ◄──── 余额不足时等待
                  └────┬─────┘
                       │
                  ┌────▼─────┐
                  │ Agent注册  │
                  └────┬─────┘
                       │
                  ┌────▼─────┐
                  │ 宪章部署   │
                  └────┬─────┘
                       │
                  ┌────▼─────┐
                  │ API网关启动│ ◄──── 健康检查
                  └────┬─────┘
                       │
                  ┌────▼─────┐
                  │ 事件订阅   │
                  └────┬─────┘
                       │
                  ┌────▼─────┐
                  │  运行中    │
                  └──────────┘

附录 B:常见问题

Q: 创建 DID 时返回 "out of gas"
A: 增加 gas_limit。DID 文档较大时需 500000+ gas。使用 --gas auto --gas-adjustment 2.0 参数。

Q: 注册 Agent 时质押金不足
A: Agent 注册需要至少 1 MSG (1000000 umsg) 的质押金。请先充值。

Q: WebSocket 连接频繁断开
A: 设置心跳间隔(ping_interval: 30s),实现自动重连逻辑。参考事件订阅器中的重连代码。

Q: MPC 签名超时
A: 确保所有参与方在线。设置合理的超时时间(默认 300 秒),超时后可重新发起提案。

Q: 宪章规则不生效
A: 检查策略引擎是否正确加载了宪章。API 网关启动时会自动查询并缓存宪章。

Q: A2A 通信失败
A: 确认目标 Agent 已注册 A2A 端点。检查防火墙规则(端口 8080)。验证签名。

附录 C:安全清单

附录 D:参考资源

资源 链接
MSG Chain 官网 https://msgchain.org
区块浏览器 https://explorer.msgchain.org
水龙头 https://faucet.msgchain.org
开发文档 https://docs.msgchain.org
GitHub https://github.com/msgchain
Discord https://discord.gg/msgchain
RPC 端点 https://rpc.msgchain.org:443
REST 端点 https://rest.msgchain.org:443
WebSocket wss://rpc.msgchain.org:443/websocket

附录 E:合约接口速查

# DID Registry 接口
DID_INTERFACE = {
    "execute": ["create_did", "update_did", "deactivate_did"],
    "query": ["query_did", "query_did_by_also_known_as", "list_active_dids"]
}

# Agent Registry 接口
REGISTRY_INTERFACE = {
    "execute": ["register_agent", "update_agent", "deregister_agent"],
    "query": ["query_agent", "query_agents_by_capability", "query_agents_by_tag",
              "list_active_agents", "list_agents_by_owner"]
}

# Constitution 接口
CONSTITUTION_INTERFACE = {
    "execute": ["set_constitution", "update_constitution", "propose_amendment",
                "approve_amendment", "check_action"],
    "query": ["query_constitution", "query_strategy_rules", "query_amendment_history"]
}

# AIPAY 接口
AIPAY_INTERFACE = {
    "execute": ["create_session", "direct_pay", "mpc_pay", "refund"],
    "query": ["get_session", "get_agent_balance", "get_transaction_history"]
}

# MPC Wallet 接口
MPC_INTERFACE = {
    "execute": ["create_mpc_wallet", "propose_transaction", "sign_transaction",
                "execute_signed_transaction"],
    "query": ["get_wallet_info", "get_pending_proposals", "get_signature_status"]
}

附录 F:Agent 自检命令

# ==================== 一键自检命令 ====================

# 1. 检查链连接
echo "=== 链连接 ==="
msgd status 2>/dev/null | python3 -c "import sys,json; s=json.load(sys.stdin); print(f'区块高度: {s[\"sync_info\"][\"latest_block_height\"]}')" 2>/dev/null || echo "✗ 连接失败"

# 2. 检查钱包
echo "=== 钱包 ==="
msgd keys list 2>/dev/null | python3 -c "import sys,json; [print(f'  {k[\"name\"]}: {k[\"address\"]}') for k in json.load(sys.stdin)]" 2>/dev/null

# 3. 检查余额
echo "=== 余额 ==="
msgd query bank balances $(msgd keys show agent_owner -a 2>/dev/null) 2>/dev/null | python3 -c "import sys,json;[print(f'  {c[\"denom\"]}: {c[\"amount\"]}') for c in json.load(sys.stdin).get('balances',[])]" 2>/dev/null

# 4. 检查 DID
echo "=== DID ==="
msgd query wasm contract-state smart $DID_REGISTRY_ADDRESS \
  "{\"query_did\":{\"did\":\"did:msg:$(msgd keys show agent_owner -a 2>/dev/null)\"}}" \
  --node $MSG_RPC_URL 2>/dev/null | python3 -c "import sys,json; d=json.load(sys.stdin); print(f'  DID: {d.get(\"did\",{}).get(\"id\",\"未找到\")}')" 2>/dev/null || echo "  DID 未创建"

# 5. 检查 Agent 注册
echo "=== Agent 注册 ==="
msgd query wasm contract-state smart $AGENT_REGISTRY_ADDRESS \
  "{\"query_agent\":{\"agent_id\":\"$(msgd keys show agent_owner -a 2>/dev/null)\"}}" \
  --node $MSG_RPC_URL 2>/dev/null | python3 -c "import sys,json; a=json.load(sys.stdin); print(f'  名称: {a.get(\"metadata\",{}).get(\"name\",\"未注册\")}')" 2>/dev/null || echo "  Agent 未注册"

# 6. 检查宪章
echo "=== 宪章 ==="
msgd query wasm contract-state smart $CONSTITUTION_ADDRESS \
  "{\"query_constitution\":{\"agent_id\":\"$(msgd keys show agent_owner -a 2>/dev/null)\"}}" \
  --node $MSG_RPC_URL 2>/dev/null | python3 -c "import sys,json; c=json.load(sys.stdin); print(f'  宪章: {c.get(\"constitution\",{}).get(\"name\",\"未设置\")}')" 2>/dev/null || echo "  宪章未设置"

# 7. 检查 API 网关
echo "=== API 网关 ==="
curl -s http://localhost:8080/health 2>/dev/null | python3 -c "import sys,json; h=json.load(sys.stdin); print(f'  状态: {h.get(\"status\",\"未运行\")}')" 2>/dev/null || echo "  API 网关未运行"

# 8. 检查 Metrics
echo "=== 监控 ==="
curl -s http://localhost:9090/metrics 2>/dev/null | head -5 || echo "  Metrics 未运行"

文档结束

本文档提供了在 MSG Chain 上启动 AI Agent 的完整指南。
如有问题,请参考 开发文档 或加入 Discord 社区。

祝你的 Agent 运行顺利!🚀


本文档内容基于 MSGChain 代码库真实状态编写,非 AI 自动生成。
主网状态: No-Go | 白皮书: https://msgchain.org/whitepaper/