dApp Docs/AI Agent 沙箱与模拟测试环境指南
Development reference. Not independently verified for production.

AI Agent 沙箱与模拟测试环境指南

适用链: msg-chain-1 | 地址前缀: msg | 版本: v1.0

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


目录

  1. 概述
  2. 本地链沙箱
  3. Mock Agent Registry
  4. Mock A2A 通信
  5. Mock AIPAY 支付
  6. Agent 行为模拟框架
  7. 宪法与策略验证
  8. CI/CD 集成

1. 概述

1.1 为什么需要沙箱测试

AI Agent 在区块链上运行涉及不可逆的交易和真实资产。沙箱测试环境允许开发者在无风险条件下验证 Agent 行为、发现漏洞并优化策略。在 MSG Chain 上构建 Agent 时,需要测试以下关键维度:

1.2 沙箱架构总览

沙箱环境采用三层架构设计:

+------------------------------------------------------+
|                  Agent 行为模拟层                        |
|  +----------+  +----------+  +----------+             |
|  | Agent A  |  | Agent B  |  | Agent C  |             |
|  +----+-----+  +----+-----+  +----+-----+             |
|       |              |              |                   |
|  +----v--------------v--------------v-----+            |
|  |         Mock 服务层                     |            |
|  |  +---------+ +---------+ +----------+  |            |
|  |  | Mock    | | Mock A2A| | Mock     |  |            |
|  |  |Registry | | Bus     | |AIPAY     |  |            |
|  |  +---------+ +---------+ +----------+  |            |
|  +----------------+-----------------------+            |
|                   |                                     |
|  +----------------v-----------------------+            |
|  |          msgd 本地链节点                |            |
|  |   +----------+  +------------------+   |            |
|  |   | CometBFT |  | CosmWasm VM      |   |            |
|  |   +----------+  +------------------+   |            |
|  +----------------------------------------+            |
+------------------------------------------------------+

1.3 确定性模拟

沙箱的核心价值在于确定性:相同的输入 + 相同的种子 = 完全相同的结果。这使得:

1.4 沙箱 vs 测试网

维度 沙箱 测试网
启动速度 秒级 N/A
确定性 完全确定 受共识影响
费用 零成本 需领取代币
隔离性 完全隔离 与其他测试者共享
快照/回放 支持 不支持
调试能力 完整日志/断点 有限

1.5 前置要求

# 最低需求
- msgd >= v1.2.0
- Python >= 3.10
- Rust >= 1.75
- Docker >= 24.0
- 磁盘: 2GB 可用空间
- 内存: 4GB 推荐

2. 本地链沙箱

2.1 使用 msgd init 启动单节点

最简单的沙箱方式是启动一个本地的 msgd 单节点实例。

# 初始化验证节点
msgd init sandbox-node \
  --chain-id msg-chain-1 \
  --default-denom umsg

# 创建核心配置文件
msgd config chain-id msg-chain-1
msgd config keyring-backend test
msgd config output json
# 创建测试账户
msgd keys add alice --keyring-backend test
msgd keys add bob --keyring-backend test
msgd keys add charlie --keyring-backend test

# 为测试账户预充值
msgd add-genesis-account $(msgd keys show alice -a --keyring-backend test) \
  1000000000000umsg

msgd add-genesis-account $(msgd keys show bob -a --keyring-backend test) \
  1000000000000umsg

msgd add-genesis-account $(msgd keys show charlie -a --keyring-backend test) \
  1000000000000umsg

# 生成创世交易
msgd gentx alice 1000000000umsg \
  --chain-id msg-chain-1 \
  --keyring-backend test

# 收集创世交易
msgd collect-gentxs
# 启动节点
msgd start \
  --minimum-gas-prices 0.001umsg \
  --rpc.laddr tcp://0.0.0.0:26657 \
  --grpc.enable true \
  --grpc.address 0.0.0.0:9090

2.2 Docker Compose 一键部署

对于需要快速重置的 CI 环境,推荐使用 Docker Compose:

# docker-compose-msg-sandbox.yml
version: "3.9"

services:
  msg-sandbox:
    image: msgchain/msgd:latest
    container_name: msg-sandbox-node
    ports:
      - "26657:26657"   # RPC
      - "9090:9090"     # gRPC
      - "1317:1317"     # REST API
    environment:
      - MONIKER=sandbox-node
      - CHAIN_ID=msg-chain-1
      - KEYRING_BACKEND=test
      - MIN_GAS_PRICES=0.001umsg
      - DENOM=umsg
    volumes:
      - msg-sandbox-data:/root/.msgd
    command: >
      sh -c "
        msgd init $${MONIKER} --chain-id $${CHAIN_ID} --default-denom $${DENOM} &&
        msgd keys add faucet --keyring-backend $${KEYRING_BACKEND} 2>&1 &&
        msgd add-genesis-account $$(msgd keys show faucet -a --keyring-backend $${KEYRING_BACKEND}) 1000000000000000umsg &&
        msgd gentx faucet 1000000000umsg --chain-id $${CHAIN_ID} --keyring-backend $${KEYRING_BACKEND} &&
        msgd collect-gentxs &&
        msgd start
      "
    healthcheck:
      test: ["CMD", "curl", "-f", "http://localhost:26657/status"]
      interval: 5s
      timeout: 3s
      retries: 15
      start_period: 20s

  init-accounts:
    image: msgchain/msgd:latest
    container_name: msg-init-accounts
    depends_on:
      msg-sandbox:
        condition: service_healthy
    volumes:
      - msg-sandbox-data:/root/.msgd
    command: >
      sh -c "
        sleep 5 &&
        for name in agent_alice agent_bob agent_charlie trader_dave oracle_eve; do
          msgd keys add $$name --keyring-backend test --home /root/.msgd 2>/dev/null &&
          addr=$$(msgd keys show $$name -a --keyring-backend test --home /root/.msgd) &&
          msgd tx bank send faucet $$addr 1000000000000umsg \
            --chain-id msg-chain-1 \
            --keyring-backend test \
            --home /root/.msgd \
            --node tcp://msg-sandbox:26657 \
            -y > /dev/null 2>&1
        done
      "

volumes:
  msg-sandbox-data:
# 启动沙箱
docker compose -f docker-compose-msg-sandbox.yml up -d

# 查看状态
curl http://localhost:26657/status | jq .

# 验证账户余额
msgd query bank balances \
  $(msgd keys show agent_alice -a --keyring-backend test) \
  --node tcp://localhost:26657

2.3 Gas-Free 模式

开发测试阶段,gas 消耗是阻碍快速迭代的主要因素。以下配置可实现免费交易:

# 方案 A: 修改 app.toml
sed -i 's/minimum-gas-prices = "0.001umsg"/minimum-gas-prices = "0"/g' \
  ~/.msgd/config/app.toml

# 方案 B: 使用 --gas-prices 标志
msgd tx bank send alice bob 1000umsg \
  --gas-prices 0umsg \
  --gas auto \
  --gas-adjustment 1.5 \
  --chain-id msg-chain-1 \
  --keyring-backend test \
  --node tcp://localhost:26657 \
  -y

# 方案 C: 修改创世 gas 上限
sed -i 's/"max_gas": "-1"/"max_gas": "50000000"/g' \
  ~/.msgd/config/genesis.json
# gas_free_mode.py - Gas-Free 模式管理工具
import json
import subprocess
from pathlib import Path


class GasFreeManager:

    CONFIG_PATHS = {
        "app": Path.home() / ".msgd" / "config" / "app.toml",
        "genesis": Path.home() / ".msgd" / "config" / "genesis.json",
    }

    def __init__(self, node_rpc: str = "tcp://localhost:26657"):
        self.node_rpc = node_rpc

    def enable_gas_free(self) -> None:
        app_toml = self.CONFIG_PATHS["app"].read_text()
        app_toml = app_toml.replace(
            'minimum-gas-prices = "0.001umsg"',
            'minimum-gas-prices = "0"',
        )
        self.CONFIG_PATHS["app"].write_text(app_toml)
        genesis = json.loads(self.CONFIG_PATHS["genesis"].read_text())
        genesis["consensus_params"]["block"]["max_gas"] = "50000000"
        self.CONFIG_PATHS["genesis"].write_text(json.dumps(genesis, indent=2))
        subprocess.run(["msgd", "unsafe-reset-all"], check=False)
        print("Gas-Free mode enabled, restart msgd start")

    def send_free_tx(self, from_key, to_addr, amount="1000umsg", memo="gas-free test"):
        result = subprocess.run(
            [
                "msgd", "tx", "bank", "send", from_key, to_addr, amount,
                "--gas-prices", "0umsg", "--gas", "auto",
                "--gas-adjustment", "1.5", "--chain-id", "msg-chain-1",
                "--keyring-backend", "test", "--node", self.node_rpc,
                "--memo", memo, "-y", "--output", "json",
            ],
            capture_output=True, text=True,
        )
        return json.loads(result.stdout)

    def get_balance(self, address: str) -> int:
        result = subprocess.run(
            [
                "msgd", "query", "bank", "balances", address,
                "--node", self.node_rpc, "--output", "json",
            ],
            capture_output=True, text=True,
        )
        data = json.loads(result.stdout)
        for coin in data.get("balances", []):
            if coin["denom"] == "umsg":
                return int(coin["amount"])
        return 0

2.4 预置测试账户

# sandbox_accounts.py
from dataclasses import dataclass


@dataclass
class SandboxAccount:
    name: str
    mnemonic: str
    address: str
    initial_balance: str
    role: str


SANDBOX_ACCOUNTS = [
    SandboxAccount(
        name="agent_alice",
        mnemonic=("abandon abandon abandon abandon abandon abandon "
                  "abandon abandon abandon abandon abandon able"),
        address="msg1qy0ag94n9q3q7q7q7q7q7q7q7q7q7q7q7q7q7q",
        initial_balance="1000000000000000umsg",
        role="agent_operator",
    ),
    SandboxAccount(
        name="agent_bob",
        mnemonic=("abandon abandon abandon abandon abandon abandon "
                  "abandon abandon abandon abandon abandon acid"),
        address="msg1qy0ag94n9q3q7q7q7q7q7q7q7q7q7q7q7q7q7r",
        initial_balance="1000000000000000umsg",
        role="agent_operator",
    ),
    SandboxAccount(
        name="agent_charlie",
        mnemonic=("abandon abandon abandon abandon abandon abandon "
                  "abandon abandon abandon abandon abandon actor"),
        address="msg1qy0ag94n9q3q7q7q7q7q7q7q7q7q7q7q7q7q8q",
        initial_balance="1000000000000000umsg",
        role="agent_operator",
    ),
    SandboxAccount(
        name="trader_dave",
        mnemonic=("abandon abandon abandon abandon abandon abandon "
                  "abandon abandon abandon abandon abandon adapt"),
        address="msg1qy0ag94n9q3q7q7q7q7q7q7q7q7q7q7q7q7q9q",
        initial_balance="500000000000000umsg",
        role="trader",
    ),
    SandboxAccount(
        name="oracle_eve",
        mnemonic=("abandon abandon abandon abandon abandon abandon "
                  "abandon abandon abandon abandon abandon address"),
        address="msg1qy0ag94n9q3q7q7q7q7q7q7q7q7q7q7q7q8q0q",
        initial_balance="200000000000000umsg",
        role="oracle",
    ),
    SandboxAccount(
        name="governance_frank",
        mnemonic=("abandon abandon abandon abandon abandon abandon "
                  "abandon abandon abandon abandon abandon adjust"),
        address="msg1qy0ag94n9q3q7q7q7q7q7q7q7q7q7q7q7q8q1q",
        initial_balance="100000000000000umsg",
        role="governance",
    ),
    SandboxAccount(
        name="relay_grace",
        mnemonic=("abandon abandon abandon abandon abandon abandon "
                  "abandon abandon abandon abandon abandon advice"),
        address="msg1qy0ag94n9q3q7q7q7q7q7q7q7q7q7q7q7q8q2q",
        initial_balance="100000000000000umsg",
        role="relayer",
    ),
]


def init_sandbox_accounts(node_rpc: str = "tcp://localhost:26657"):
    import subprocess
    for acct in SANDBOX_ACCOUNTS:
        result = subprocess.run(
            ["msgd", "keys", "add", acct.name, "--recover", "--keyring-backend", "test"],
            input=acct.mnemonic + "\n",
            capture_output=True, text=True,
        )
        if result.returncode == 0:
            print(f"Imported: {acct.name} ({acct.role})")
        else:
            print(f"Skipped: {acct.name}")
        subprocess.run(
            ["msgd", "tx", "bank", "send", "faucet", acct.address, acct.initial_balance,
             "--chain-id", "msg-chain-1", "--keyring-backend", "test",
             "--node", node_rpc, "-y"],
            capture_output=True,
        )


def list_sandbox_accounts():
    import subprocess, json
    for acct in SANDBOX_ACCOUNTS:
        result = subprocess.run(
            ["msgd", "query", "bank", "balances", acct.address,
             "--node", "tcp://localhost:26657", "--output", "json"],
            capture_output=True, text=True,
        )
        data = json.loads(result.stdout)
        balance = ", ".join(
            f"{c['amount']}{c['denom']}" for c in data.get("balances", [])
        ) or "0umsg"
        print(f"{acct.name:30s} | {acct.role:15s} | {balance}")

2.5 创世文件定制

# 修改创世参数以获得更快的测试

# 1. 缩短出块时间(默认为 5000ms)
jq '.consensus_params.block.time_iota_ms = "500"' \
  ~/.msgd/config/genesis.json > tmp.json && mv tmp.json ~/.msgd/config/genesis.json

# 2. 减少投票周期
jq '.consensus_params.evidence.max_age_duration = "300000000000"' \
  ~/.msgd/config/genesis.json > tmp.json && mv tmp.json ~/.msgd/config/genesis.json

# 3. 启用 Agent 模块
jq '.app_state += {"agent": {"params": {"max_agents": 1000, "registration_fee": "0umsg"}}}' \
  ~/.msgd/config/genesis.json > tmp.json && mv tmp.json ~/.msgd/config/genesis.json

2.6 快照与恢复

# 创建快照
function sandbox_snapshot() {
    local NAME="${1:-sandbox-snapshot}"
    local SNAPSHOT_DIR="/tmp/msg-sandbox-snapshots/${NAME}"
    echo "Creating snapshot: ${NAME}"
    mkdir -p "${SNAPSHOT_DIR}"
    pkill msgd 2>/dev/null; sleep 2
    cp -r ~/.msgd/data "${SNAPSHOT_DIR}/data"
    cp ~/.msgd/config/genesis.json "${SNAPSHOT_DIR}/"
    cp ~/.msgd/config/priv_validator_key.json "${SNAPSHOT_DIR}/"
    echo "Snapshot saved: ${SNAPSHOT_DIR}"
    msgd start --minimum-gas-prices 0umsg &
}

# 恢复快照
function sandbox_restore() {
    local NAME="${1:-sandbox-snapshot}"
    local SNAPSHOT_DIR="/tmp/msg-sandbox-snapshots/${NAME}"
    if [ ! -d "${SNAPSHOT_DIR}" ]; then
        echo "Snapshot not found: ${SNAPSHOT_DIR}"
        return 1
    fi
    echo "Restoring snapshot: ${NAME}"
    pkill msgd 2>/dev/null; sleep 2
    msgd unsafe-reset-all 2>/dev/null
    cp -r "${SNAPSHOT_DIR}/data" ~/.msgd/data
    cp "${SNAPSHOT_DIR}/genesis.json" ~/.msgd/config/
    cp "${SNAPSHOT_DIR}/priv_validator_key.json" ~/.msgd/config/
    echo "Snapshot restored"
    msgd start --minimum-gas-prices 0umsg &
}

2.7 常用调试命令

# 沙箱调试命令速查表

# 查询当前区块高度
msgd status --node tcp://localhost:26657 | jq .SyncInfo.latest_block_height

# 查询验证者集合
msgd query staking validators --node tcp://localhost:26657 --output json | \
  jq '.validators[] | {operator_address, tokens, status}'

# 查看事件日志
msgd query tx --node tcp://localhost:26657 --output json --type=hash <TX_HASH> | jq '.logs'

# 检查链上 Agent 模块状态
msgd query agent list --node tcp://localhost:26657 --output json

# 导出区块链状态
msgd export > /tmp/sandbox-export.json

# 重置沙箱(快速清空)
msgd unsafe-reset-all && rm -rf ~/.msgd/data/*

3. Mock Agent Registry

3.1 Rust 实现

// mock_registry/src/lib.rs
use std::collections::HashMap;
use serde::{Deserialize, Serialize};

#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub enum AgentStatus {
    Active,
    Paused,
    Disabled,
    Slashed,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AgentCapability {
    pub name: String,
    pub version: String,
    pub parameters: Vec<String>,
    pub required_balance: u128,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MockAgent {
    pub id: u64,
    pub owner: String,
    pub address: String,
    pub name: String,
    pub description: String,
    pub capabilities: Vec<AgentCapability>,
    pub status: AgentStatus,
    pub registered_at: u64,
    pub metadata: HashMap<String, String>,
    pub delegation_level: u8,
}

#[derive(Debug, Default)]
pub struct AgentFilter {
    pub status: Option<AgentStatus>,
    pub capability: Option<String>,
    pub owner: Option<String>,
    pub min_delegation: Option<u8>,
}

pub struct MockRegistry {
    pub agents: HashMap<u64, MockAgent>,
    pub next_id: u64,
    pub owner_index: HashMap<String, Vec<u64>>,
    pub capability_index: HashMap<String, Vec<u64>>,
}

impl MockRegistry {
    pub fn new() -> Self {
        Self {
            agents: HashMap::new(),
            next_id: 1,
            owner_index: HashMap::new(),
            capability_index: HashMap::new(),
        }
    }

    pub fn register_agent(
        &mut self,
        owner: String,
        name: String,
        description: String,
        capabilities: Vec<AgentCapability>,
        delegation_level: u8,
    ) -> Result<MockAgent, RegistryError> {
        if name.is_empty() {
            return Err(RegistryError::EmptyName);
        }
        if capabilities.is_empty() {
            return Err(RegistryError::NoCapabilities);
        }
        if delegation_level > 5 {
            return Err(RegistryError::DelegationOutOfRange);
        }
        let id = self.next_id;
        self.next_id += 1;
        let address = format!("msg1agent{:016x}", id);
        let agent = MockAgent {
            id,
            owner: owner.clone(),
            address,
            name,
            description,
            capabilities: capabilities.clone(),
            status: AgentStatus::Active,
            registered_at: self.current_time(),
            metadata: HashMap::new(),
            delegation_level,
        };
        self.agents.insert(id, agent.clone());
        self.owner_index.entry(owner).or_default().push(id);
        for cap in &capabilities {
            self.capability_index.entry(cap.name.clone()).or_default().push(id);
        }
        Ok(agent)
    }

    pub fn get_agent(&self, id: u64) -> Option<&MockAgent> {
        self.agents.get(&id)
    }

    pub fn get_agent_by_address(&self, address: &str) -> Option<&MockAgent> {
        self.agents.values().find(|a| a.address == address)
    }

    pub fn list_agents(&self, filter: Option<AgentFilter>) -> Vec<&MockAgent> {
        let agents: Vec<&MockAgent> = self.agents.values().collect();
        match filter {
            Some(f) => agents.into_iter().filter(|a| {
                let mut matches = true;
                if let Some(ref status) = f.status { matches &= a.status == *status; }
                if let Some(ref cap) = f.capability {
                    matches &= a.capabilities.iter().any(|c| c.name == *cap);
                }
                if let Some(ref owner) = f.owner { matches &= a.owner == *owner; }
                if let Some(min_del) = f.min_delegation {
                    matches &= a.delegation_level >= min_del;
                }
                matches
            }).collect(),
            None => agents,
        }
    }

    pub fn update_status(&mut self, id: u64, status: AgentStatus) -> Result<(), RegistryError> {
        match self.agents.get_mut(&id) {
            Some(agent) => { agent.status = status; Ok(()) }
            None => Err(RegistryError::AgentNotFound(id)),
        }
    }

    pub fn find_by_capability(&self, capability: &str) -> Vec<&MockAgent> {
        self.capability_index.get(capability)
            .map(|ids| ids.iter().filter_map(|id| self.agents.get(id)).collect())
            .unwrap_or_default()
    }

    pub fn find_by_owner(&self, owner: &str) -> Vec<&MockAgent> {
        self.owner_index.get(owner)
            .map(|ids| ids.iter().filter_map(|id| self.agents.get(id)).collect())
            .unwrap_or_default()
    }

    pub fn deregister(&mut self, id: u64) -> Result<MockAgent, RegistryError> {
        let agent = self.agents.remove(&id).ok_or(RegistryError::AgentNotFound(id))?;
        if let Some(ids) = self.owner_index.get_mut(&agent.owner) {
            ids.retain(|&i| i != id);
        }
        for cap in &agent.capabilities {
            if let Some(ids) = self.capability_index.get_mut(&cap.name) {
                ids.retain(|&i| i != id);
            }
        }
        Ok(agent)
    }

    pub fn stats(&self) -> RegistryStats {
        let total = self.agents.len() as u64;
        let active = self.agents.values()
            .filter(|a| a.status == AgentStatus::Active).count() as u64;
        let by_capability: HashMap<String, u64> = self.capability_index
            .iter().map(|(k, v)| (k.clone(), v.len() as u64)).collect();
        RegistryStats { total, active, by_capability }
    }

    fn current_time(&self) -> u64 { 1700000000 + self.agents.len() as u64 }
}

#[derive(Debug)]
pub enum RegistryError {
    AgentNotFound(u64),
    EmptyName,
    NoCapabilities,
    DelegationOutOfRange,
    DuplicateEntry,
}

impl std::fmt::Display for RegistryError {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        match self {
            RegistryError::AgentNotFound(id) => write!(f, "Agent {} not found", id),
            RegistryError::EmptyName => write!(f, "Agent name cannot be empty"),
            RegistryError::NoCapabilities => write!(f, "At least one capability required"),
            RegistryError::DelegationOutOfRange => write!(f, "Delegation level must be 0-5"),
            RegistryError::DuplicateEntry => write!(f, "Duplicate agent entry"),
        }
    }
}

impl std::error::Error for RegistryError {}

#[derive(Debug, Serialize, Deserialize)]
pub struct RegistryStats {
    pub total: u64,
    pub active: u64,
    pub by_capability: HashMap<String, u64>,
}

3.2 Python 绑定

# mock_registry_py/bridge.py
from __future__ import annotations
import json
from dataclasses import dataclass, field
from enum import Enum
from typing import Optional


class AgentStatus(str, Enum):
    ACTIVE = "active"
    PAUSED = "paused"
    DISABLED = "disabled"
    SLASHED = "slashed"


@dataclass
class AgentCapability:
    name: str
    version: str
    parameters: list[str] = field(default_factory=list)
    required_balance: int = 0


@dataclass
class MockAgent:
    id: int
    owner: str
    address: str
    name: str
    description: str
    capabilities: list[AgentCapability]
    status: AgentStatus = AgentStatus.ACTIVE
    registered_at: int = 0
    metadata: dict[str, str] = field(default_factory=dict)
    delegation_level: int = 0


@dataclass
class RegistryStats:
    total: int = 0
    active: int = 0
    by_capability: dict[str, int] = field(default_factory=dict)


class MockRegistry:

    def __init__(self):
        self._agents: dict[int, MockAgent] = {}
        self._next_id: int = 1
        self._owner_index: dict[str, list[int]] = {}
        self._capability_index: dict[str, list[int]] = {}

    def register_agent(self, owner, name, description="",
                       capabilities=None, delegation_level=0):
        if not name:
            raise ValueError("Agent name cannot be empty")
        if not capabilities:
            raise ValueError("At least one capability required")
        if not (0 <= delegation_level <= 5):
            raise ValueError("Delegation level must be 0-5")
        agent_id = self._next_id
        self._next_id += 1
        address = f"msg1agent{agent_id:016x}"
        agent = MockAgent(
            id=agent_id, owner=owner, address=address, name=name,
            description=description, capabilities=capabilities,
            status=AgentStatus.ACTIVE, registered_at=1700000000,
            delegation_level=delegation_level,
        )
        self._agents[agent_id] = agent
        self._owner_index.setdefault(owner, []).append(agent_id)
        for cap in capabilities:
            self._capability_index.setdefault(cap.name, []).append(agent_id)
        return agent

    def get_agent(self, agent_id):
        return self._agents.get(agent_id)

    def list_agents(self, status=None, capability=None, owner=None, min_delegation=None):
        result = list(self._agents.values())
        if status is not None:
            result = [a for a in result if a.status == status]
        if capability is not None:
            result = [a for a in result if any(c.name == capability for c in a.capabilities)]
        if owner is not None:
            result = [a for a in result if a.owner == owner]
        if min_delegation is not None:
            result = [a for a in result if a.delegation_level >= min_delegation]
        return result

    def update_status(self, agent_id, status):
        if agent_id not in self._agents:
            raise KeyError(f"Agent {agent_id} not found")
        self._agents[agent_id].status = status

    def find_by_capability(self, capability):
        ids = self._capability_index.get(capability, [])
        return [self._agents[i] for i in ids if i in self._agents]

    def find_by_owner(self, owner):
        ids = self._owner_index.get(owner, [])
        return [self._agents[i] for i in ids if i in self._agents]

    def deregister(self, agent_id):
        if agent_id not in self._agents:
            raise KeyError(f"Agent {agent_id} not found")
        agent = self._agents.pop(agent_id)
        if agent.owner in self._owner_index:
            self._owner_index[agent.owner] = [
                i for i in self._owner_index[agent.owner] if i != agent_id
            ]
        for cap in agent.capabilities:
            if cap.name in self._capability_index:
                self._capability_index[cap.name] = [
                    i for i in self._capability_index[cap.name] if i != agent_id
                ]
        return agent

    def stats(self):
        total = len(self._agents)
        active = sum(1 for a in self._agents.values() if a.status == AgentStatus.ACTIVE)
        by_capability = {k: len(v) for k, v in self._capability_index.items()}
        return RegistryStats(total=total, active=active, by_capability=by_capability)

    def reset(self):
        self._agents.clear()
        self._next_id = 1
        self._owner_index.clear()
        self._capability_index.clear()

    def to_json(self):
        data = {
            "agents": {
                str(k): {
                    "id": v.id, "owner": v.owner, "address": v.address,
                    "name": v.name, "description": v.description,
                    "capabilities": [
                        {"name": c.name, "version": c.version,
                         "parameters": c.parameters,
                         "required_balance": c.required_balance}
                        for c in v.capabilities
                    ],
                    "status": v.status.value, "registered_at": v.registered_at,
                    "metadata": v.metadata, "delegation_level": v.delegation_level,
                }
                for k, v in self._agents.items()
            },
            "next_id": self._next_id,
        }
        return json.dumps(data, indent=2, ensure_ascii=False)

    @classmethod
    def from_json(cls, json_str):
        data = json.loads(json_str)
        registry = cls()
        registry._next_id = data["next_id"]
        for _, agent_data in data["agents"].items():
            agent = MockAgent(
                id=agent_data["id"], owner=agent_data["owner"],
                address=agent_data["address"], name=agent_data["name"],
                description=agent_data["description"],
                capabilities=[AgentCapability(**c) for c in agent_data["capabilities"]],
                status=AgentStatus(agent_data["status"]),
                registered_at=agent_data["registered_at"],
                metadata=agent_data["metadata"],
                delegation_level=agent_data["delegation_level"],
            )
            registry._agents[agent.id] = agent
            registry._owner_index.setdefault(agent.owner, []).append(agent.id)
            for cap in agent.capabilities:
                registry._capability_index.setdefault(cap.name, []).append(agent.id)
        return registry


def create_default_test_registry():
    registry = MockRegistry()
    owner = "msg1owner000000000000000000000000000000000"

    registry.register_agent(owner, "PriceOracleAgent",
        "提供实时代币价格数据",
        [AgentCapability("price_feed", "1.0.0", ["pair", "interval"], 100_000_000),
         AgentCapability("market_analysis", "1.1.0", ["token", "timeframe"], 50_000_000)],
        delegation_level=3)

    registry.register_agent(owner, "SwapExecutorAgent",
        "执行 DEX 交易和路由",
        [AgentCapability("swap", "2.0.0", ["from", "to", "amount"], 500_000_000),
         AgentCapability("route_optimization", "1.0.0",
                        ["token_in", "token_out", "max_hops"], 200_000_000)],
        delegation_level=4)

    registry.register_agent(owner, "RiskManagerAgent",
        "监控链上风险并触发保护措施",
        [AgentCapability("risk_assessment", "1.2.0", ["protocol", "threshold"], 100_000_000),
         AgentCapability("circuit_breaker", "1.0.0", ["condition", "action"], 300_000_000)],
        delegation_level=5)

    registry.register_agent(owner, "LiquidityManagerAgent",
        "自动管理流动性池",
        [AgentCapability("liquidity_provision", "1.0.0",
                        ["pool", "amount", "range"], 1_000_000_000),
         AgentCapability("yield_optimization", "2.0.0",
                        ["strategy", "min_apy"], 500_000_000)],
        delegation_level=2)

    registry.register_agent(owner, "RelayAgent",
        "处理跨链消息和资产转移",
        [AgentCapability("cross_chain_transfer", "1.0.0",
                        ["destination", "asset", "amount"], 200_000_000)],
        delegation_level=3)

    return registry

3.3 注册表测试

#[cfg(test)]
mod tests {
    use mock_registry::*;

    #[test]
    fn test_register_and_query() {
        let mut registry = MockRegistry::new();
        let owner = "msg1test".to_string();
        let agent = registry.register_agent(
            owner.clone(), "TestAgent".into(), "desc".into(),
            vec![AgentCapability {
                name: "test_cap".into(), version: "1.0".into(),
                parameters: vec![], required_balance: 0,
            }],
            1,
        ).unwrap();
        assert_eq!(agent.id, 1);
        assert_eq!(agent.status, AgentStatus::Active);
        let fetched = registry.get_agent(1).unwrap();
        assert_eq!(fetched.name, "TestAgent");
    }

    #[test]
    fn test_filter_agents() {
        let mut registry = MockRegistry::new();
        registry.register_agent("owner".into(), "A".into(), "".into(),
            vec![AgentCapability {
                name: "cap_a".into(), version: "1".into(),
                parameters: vec![], required_balance: 0,
            }], 1,
        ).unwrap();
        registry.register_agent("owner".into(), "B".into(), "".into(),
            vec![AgentCapability {
                name: "cap_b".into(), version: "1".into(),
                parameters: vec![], required_balance: 0,
            }], 1,
        ).unwrap();
        let filtered = registry.list_agents(Some(AgentFilter {
            capability: Some("cap_a".into()), ..Default::default()
        }));
        assert_eq!(filtered.len(), 1);
    }

    #[test]
    fn test_status_transitions() {
        let mut registry = MockRegistry::new();
        registry.register_agent("owner".into(), "A".into(), "".into(),
            vec![AgentCapability {
                name: "c".into(), version: "1".into(),
                parameters: vec![], required_balance: 0,
            }], 0,
        ).unwrap();
        registry.update_status(1, AgentStatus::Paused).unwrap();
        assert_eq!(registry.get_agent(1).unwrap().status, AgentStatus::Paused);
        assert!(registry.update_status(999, AgentStatus::Active).is_err());
    }

    #[test]
    fn test_validation_errors() {
        let mut registry = MockRegistry::new();
        assert!(registry.register_agent(
            "o".into(), "".into(), "".into(), vec![], 0
        ).is_err());
        assert!(registry.register_agent(
            "o".into(), "A".into(), "".into(), vec![], 0
        ).is_err());
    }
}

4. Mock A2A 通信

4.1 内存消息总线

// mock_a2a/src/message_bus.rs
use std::collections::{HashMap, VecDeque};
use std::time::{Duration, Instant};
use serde::{Deserialize, Serialize};

#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum MessagePriority {
    Low = 0,
    Normal = 1,
    High = 2,
    Critical = 3,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct A2AMessage {
    pub id: u64,
    pub sender: String,
    pub recipient: String,
    pub msg_type: String,
    pub payload: Vec<u8>,
    pub priority: MessagePriority,
    pub timestamp: u64,
    pub ttl: u64,
    pub correlation_id: Option<String>,
}

#[derive(Debug, Clone)]
pub struct DelayConfig {
    pub base_delay_ms: u64,
    pub jitter_ms: u64,
    pub drop_rate: f64,
}

impl Default for DelayConfig {
    fn default() -> Self {
        Self { base_delay_ms: 100, jitter_ms: 50, drop_rate: 0.0 }
    }
}

struct DelayedMessage {
    message: A2AMessage,
    deliver_at: Instant,
}

pub struct MockA2ABus {
    inbox: HashMap<String, VecDeque<A2AMessage>>,
    delay_queue: VecDeque<DelayedMessage>,
    sent_count: u64,
    delivered_count: u64,
    dropped_count: u64,
    failed_count: u64,
    message_id_counter: u64,
    delay_config: DelayConfig,
    rng_state: u64,
}

impl MockA2ABus {
    pub fn new() -> Self {
        Self::with_config(DelayConfig::default())
    }

    pub fn with_config(config: DelayConfig) -> Self {
        Self {
            inbox: HashMap::new(),
            delay_queue: VecDeque::new(),
            sent_count: 0,
            delivered_count: 0,
            dropped_count: 0,
            failed_count: 0,
            message_id_counter: 0,
            delay_config: config,
            rng_state: 42,
        }
    }

    pub fn send(&mut self, sender: &str, recipient: &str,
                msg_type: &str, payload: Vec<u8>) -> u64 {
        self.send_with_priority(sender, recipient, msg_type, payload, MessagePriority::Normal)
    }

    pub fn send_with_priority(&mut self, sender: &str, recipient: &str,
                              msg_type: &str, payload: Vec<u8>,
                              priority: MessagePriority) -> u64 {
        self.message_id_counter += 1;
        let message_id = self.message_id_counter;

        if self.should_drop() {
            self.dropped_count += 1;
            return message_id;
        }

        let msg = A2AMessage {
            id: message_id,
            sender: sender.to_string(),
            recipient: recipient.to_string(),
            msg_type: msg_type.to_string(),
            payload,
            priority,
            timestamp: 1700000000,
            ttl: 60000,
            correlation_id: None,
        };

        let delay_ms = self.compute_delay();
        let deliver_at = Instant::now() + Duration::from_millis(delay_ms);
        self.delay_queue.push_back(DelayedMessage { message: msg, deliver_at });
        self.sent_count += 1;
        message_id
    }

    pub fn poll(&mut self) -> usize {
        let now = Instant::now();
        let mut delivered = 0;
        while let Some(front) = self.delay_queue.front() {
            if front.deliver_at > now { break; }
            let entry = self.delay_queue.pop_front().unwrap();
            self.inbox.entry(entry.message.recipient.clone())
                .or_default().push_back(entry.message);
            delivered += 1;
        }
        self.delivered_count += delivered;
        delivered
    }

    pub fn receive(&mut self, recipient: &str) -> Option<A2AMessage> {
        self.poll();
        self.inbox.get_mut(recipient).and_then(|q| q.pop_front())
    }

    pub fn receive_all(&mut self, recipient: &str) -> Vec<A2AMessage> {
        self.poll();
        self.inbox.get_mut(recipient)
            .map(|q| q.drain(..).collect()).unwrap_or_default()
    }

    pub fn inbox_size(&self, recipient: &str) -> usize {
        self.inbox.get(recipient).map(|q| q.len()).unwrap_or(0)
    }

    pub fn reset(&mut self) {
        self.inbox.clear();
        self.delay_queue.clear();
        self.sent_count = 0;
        self.delivered_count = 0;
        self.dropped_count = 0;
        self.failed_count = 0;
    }

    pub fn stats(&self) -> A2AStats {
        A2AStats {
            sent: self.sent_count,
            delivered: self.delivered_count,
            dropped: self.dropped_count,
            failed: self.failed_count,
            in_flight: self.delay_queue.len() as u64,
            in_inboxes: self.inbox.values().map(|q| q.len() as u64).sum(),
        }
    }

    fn should_drop(&mut self) -> bool {
        if self.delay_config.drop_rate <= 0.0 { return false; }
        let val = self.next_random() as f64 / u64::MAX as f64;
        val < self.delay_config.drop_rate
    }

    fn compute_delay(&mut self) -> u64 {
        let jitter = (self.next_random() % self.delay_config.jitter_ms.max(1)) as u64;
        self.delay_config.base_delay_ms + jitter
    }

    fn next_random(&mut self) -> u64 {
        self.rng_state = self.rng_state
            .wrapping_mul(6364136223846793005)
            .wrapping_add(1442695040888963407);
        self.rng_state
    }
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct A2AStats {
    pub sent: u64,
    pub delivered: u64,
    pub dropped: u64,
    pub failed: u64,
    pub in_flight: u64,
    pub in_inboxes: u64,
}

4.2 消息类型定义

// mock_a2a/src/protocol.rs
pub mod msg_types {
    pub const DISCOVERY_QUERY: &str = "a2a/discovery/query";
    pub const DISCOVERY_RESPONSE: &str = "a2a/discovery/response";
    pub const CAPABILITY_ANNOUNCE: &str = "a2a/capability/announce";
    pub const TASK_ASSIGN: &str = "a2a/task/assign";
    pub const TASK_ACCEPT: &str = "a2a/task/accept";
    pub const TASK_REJECT: &str = "a2a/task/reject";
    pub const TASK_PROGRESS: &str = "a2a/task/progress";
    pub const TASK_COMPLETE: &str = "a2a/task/complete";
    pub const TASK_FAILED: &str = "a2a/task/failed";
    pub const NEGOTIATE_START: &str = "a2a/negotiate/start";
    pub const NEGOTIATE_OFFER: &str = "a2a/negotiate/offer";
    pub const NEGOTIATE_COUNTER: &str = "a2a/negotiate/counter";
    pub const NEGOTIATE_ACCEPT: &str = "a2a/negotiate/accept";
    pub const NEGOTIATE_REJECT: &str = "a2a/negotiate/reject";
    pub const PAYMENT_REQUEST: &str = "a2a/payment/request";
    pub const PAYMENT_AUTHORIZE: &str = "a2a/payment/authorize";
    pub const PAYMENT_CONFIRM: &str = "a2a/payment/confirm";
    pub const PAYMENT_FAILED: &str = "a2a/payment/failed";
}

4.3 Python 端 A2A 总线

# mock_a2a_py/bus.py
from __future__ import annotations
import json
import time
import random
from dataclasses import dataclass, field
from enum import IntEnum
from heapq import heappush, heappop
from typing import Optional


class MessagePriority(IntEnum):
    LOW = 0
    NORMAL = 1
    HIGH = 2
    CRITICAL = 3


@dataclass(order=True)
class ScheduledMessage:
    deliver_at: float
    priority: int
    sequence: int
    message: dict = field(compare=False)


@dataclass
class A2AStats:
    sent: int = 0
    delivered: int = 0
    dropped: int = 0
    failed: int = 0
    in_flight: int = 0
    in_inboxes: int = 0


class MockA2ABus:

    def __init__(self, base_delay_ms=100.0, jitter_ms=50.0,
                 drop_rate=0.0, seed=42):
        self._inbox: dict[str, list[dict]] = {}
        self._delay_queue: list[ScheduledMessage] = []
        self._sequence = 0
        self._message_id = 0
        self.base_delay_ms = base_delay_ms
        self.jitter_ms = jitter_ms
        self.drop_rate = drop_rate
        self._rng = random.Random(seed)
        self.sent_count = 0
        self.delivered_count = 0
        self.dropped_count = 0
        self.failed_count = 0

    def send(self, sender, recipient, msg_type, payload,
             priority=MessagePriority.NORMAL):
        self._message_id += 1
        if self._should_drop():
            self.dropped_count += 1
            return self._message_id

        if isinstance(payload, dict):
            payload_bytes = json.dumps(payload, ensure_ascii=False).encode()
        elif isinstance(payload, str):
            payload_bytes = payload.encode()
        else:
            payload_bytes = payload

        message = {
            "id": self._message_id,
            "sender": sender,
            "recipient": recipient,
            "msg_type": msg_type,
            "payload": payload_bytes,
            "priority": priority.value,
            "timestamp": int(time.time() * 1000),
            "ttl": 60000,
            "correlation_id": None,
        }

        delay_ms = self._compute_delay()
        deliver_at = time.time() + delay_ms / 1000.0
        self._sequence += 1
        heappush(self._delay_queue,
                 ScheduledMessage(deliver_at, priority.value, self._sequence, message))
        self.sent_count += 1
        return self._message_id

    def poll(self):
        now = time.time()
        delivered = 0
        while self._delay_queue and self._delay_queue[0].deliver_at <= now:
            entry = heappop(self._delay_queue)
            self._inbox.setdefault(entry.message["recipient"], []).append(entry.message)
            delivered += 1
        self.delivered_count += delivered
        return delivered

    def receive(self, recipient):
        self.poll()
        inbox = self._inbox.get(recipient)
        return inbox.pop(0) if inbox else None

    def receive_all(self, recipient):
        self.poll()
        return self._inbox.pop(recipient, [])

    def receive_by_type(self, recipient, msg_type):
        all_msgs = self.receive_all(recipient)
        return [m for m in all_msgs if m["msg_type"] == msg_type]

    def inbox_size(self, recipient):
        return len(self._inbox.get(recipient, []))

    def pending_count(self):
        return len(self._delay_queue)

    def reset(self):
        self._inbox.clear()
        self._delay_queue.clear()
        self._sequence = 0
        self._message_id = 0
        self.sent_count = 0
        self.delivered_count = 0
        self.dropped_count = 0
        self.failed_count = 0

    def stats(self):
        in_inboxes = sum(len(q) for q in self._inbox.values())
        return A2AStats(
            sent=self.sent_count, delivered=self.delivered_count,
            dropped=self.dropped_count, failed=self.failed_count,
            in_flight=self.pending_count(), in_inboxes=in_inboxes,
        )

    def set_drop_rate(self, rate):
        self.drop_rate = max(0.0, min(1.0, rate))

    def set_delay(self, base_ms, jitter_ms):
        self.base_delay_ms = base_ms
        self.jitter_ms = jitter_ms

    def _should_drop(self):
        if self.drop_rate <= 0.0:
            return False
        return self._rng.random() < self.drop_rate

    def _compute_delay(self):
        jitter = self._rng.uniform(0, self.jitter_ms)
        return self.base_delay_ms + jitter

4.4 通信测试

# test_a2a.py
from bus import MockA2ABus, MessagePriority


def test_basic_send_receive():
    bus = MockA2ABus(base_delay_ms=0, jitter_ms=0)
    bus.send("alice", "bob", "a2a/task/assign", {"task": "hello"})
    bus.poll()
    msg = bus.receive("bob")
    assert msg is not None
    assert msg["sender"] == "alice"
    print("test_basic_send_receive passed")


def test_priority_ordering():
    bus = MockA2ABus(base_delay_ms=0, jitter_ms=0)
    bus.send("alice", "bob", "normal", {}, MessagePriority.NORMAL)
    bus.send("alice", "bob", "critical", {}, MessagePriority.CRITICAL)
    bus.send("alice", "bob", "high", {}, MessagePriority.HIGH)
    bus.poll()
    msgs = bus.receive_all("bob")
    types = [m["msg_type"] for m in msgs]
    assert types == ["critical", "high", "normal"]
    print("test_priority_ordering passed")


def test_message_loss():
    bus = MockA2ABus(base_delay_ms=0, jitter_ms=0, drop_rate=0.5, seed=42)
    for i in range(100):
        bus.send("alice", "bob", f"msg_{i}", {"seq": i})
    bus.poll()
    received = len(bus.receive_all("bob"))
    assert 30 <= received <= 70
    print(f"test_message_loss passed: {received}/100")


def test_delay_simulation():
    import time
    bus = MockA2ABus(base_delay_ms=200, jitter_ms=0)
    bus.send("alice", "bob", "delayed", {})
    bus.poll()
    assert bus.receive("bob") is None
    time.sleep(0.25)
    bus.poll()
    assert bus.receive("bob") is not None
    print("test_delay_simulation passed")


def test_bus_stats():
    bus = MockA2ABus(base_delay_ms=0, jitter_ms=0, drop_rate=0.2, seed=42)
    for i in range(50):
        bus.send("alice", "bob", f"msg_{i}", {})
    bus.poll()
    stats = bus.stats()
    assert stats.sent == 50
    assert stats.dropped > 0
    assert stats.delivered + stats.dropped == 50
    print(f"test_bus_stats passed")


if __name__ == "__main__":
    test_basic_send_receive()
    test_priority_ordering()
    test_message_loss()
    test_delay_simulation()
    test_bus_stats()
    print("All A2A tests passed")

5. Mock AIPAY 支付

5.1 Rust 实现

// mock_aipay/src/lib.rs
use std::collections::HashMap;
use serde::{Deserialize, Serialize};

#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub enum PaymentStatus {
    Pending,
    Authorized,
    Completed,
    Failed,
    Refunded,
    Expired,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PaymentSession {
    pub id: String,
    pub payer: String,
    pub payee: String,
    pub amount: u128,
    pub denom: String,
    pub purpose: PaymentPurpose,
    pub status: PaymentStatus,
    pub created_at: u64,
    pub expires_at: u64,
    pub metadata: HashMap<String, String>,
    pub failure_reason: Option<String>,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum PaymentPurpose {
    AgentService,
    TaskReward,
    Subscription,
    Slash,
    Transfer,
    Custom(String),
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AccountBalance {
    pub address: String,
    pub balances: HashMap<String, u128>,
}

pub struct MockAIPAY {
    sessions: HashMap<String, PaymentSession>,
    accounts: HashMap<String, AccountBalance>,
    session_counter: u64,
}

impl MockAIPAY {
    pub fn new() -> Self {
        Self {
            sessions: HashMap::new(),
            accounts: HashMap::new(),
            session_counter: 0,
        }
    }

    pub fn register_account(&mut self, address: &str, balances: HashMap<String, u128>) {
        self.accounts.insert(address.to_string(), AccountBalance {
            address: address.to_string(), balances,
        });
    }

    pub fn create_session(&mut self, payer: &str, payee: &str,
                          amount: u128, denom: &str,
                          purpose: PaymentPurpose) -> Result<String, PaymentError> {
        let balance = self.get_balance(payer, denom)
            .ok_or(PaymentError::AccountNotFound(payer.to_string()))?;
        if balance < amount {
            return Err(PaymentError::InsufficientBalance {
                account: payer.to_string(), required: amount, available: balance,
            });
        }
        self.session_counter += 1;
        let session_id = format!("pay-{:016x}", self.session_counter);
        let session = PaymentSession {
            id: session_id.clone(),
            payer: payer.to_string(),
            payee: payee.to_string(),
            amount,
            denom: denom.to_string(),
            purpose,
            status: PaymentStatus::Pending,
            created_at: 1700000000,
            expires_at: 1700000300,
            metadata: HashMap::new(),
            failure_reason: None,
        };
        self.sessions.insert(session_id.clone(), session);
        Ok(session_id)
    }

    pub fn authorize(&mut self, session_id: &str) -> Result<(), PaymentError> {
        let session = self.sessions.get_mut(session_id)
            .ok_or(PaymentError::SessionNotFound(session_id.to_string()))?;
        if session.status != PaymentStatus::Pending {
            return Err(PaymentError::InvalidState {
                session: session_id.to_string(),
                expected: PaymentStatus::Pending,
                actual: session.status.clone(),
            });
        }
        let balance = self.get_balance(&session.payer, &session.denom)
            .ok_or(PaymentError::AccountNotFound(session.payer.clone()))?;
        if balance < session.amount {
            session.status = PaymentStatus::Failed;
            session.failure_reason = Some("Insufficient balance".into());
            return Err(PaymentError::InsufficientBalance {
                account: session.payer.clone(),
                required: session.amount,
                available: balance,
            });
        }
        session.status = PaymentStatus::Authorized;
        Ok(())
    }

    pub fn complete(&mut self, session_id: &str) -> Result<(), PaymentError> {
        let session = self.sessions.get(session_id)
            .ok_or(PaymentError::SessionNotFound(session_id.to_string()))?;
        if session.status != PaymentStatus::Authorized {
            return Err(PaymentError::InvalidState {
                session: session_id.to_string(),
                expected: PaymentStatus::Authorized,
                actual: session.status.clone(),
            });
        }
        let payer = session.payer.clone();
        let payee = session.payee.clone();
        let amount = session.amount;
        let denom = session.denom.clone();
        self.deduct(&payer, &denom, amount)?;
        self.credit(&payee, &denom, amount)?;
        self.sessions.get_mut(session_id).unwrap().status = PaymentStatus::Completed;
        Ok(())
    }

    pub fn refund(&mut self, session_id: &str) -> Result<(), PaymentError> {
        let session = self.sessions.get_mut(session_id)
            .ok_or(PaymentError::SessionNotFound(session_id.to_string()))?;
        if session.status != PaymentStatus::Completed {
            return Err(PaymentError::InvalidState {
                session: session_id.to_string(),
                expected: PaymentStatus::Completed,
                actual: session.status.clone(),
            });
        }
        let payer = session.payer.clone();
        let payee = session.payee.clone();
        let amount = session.amount;
        let denom = session.denom.clone();
        self.deduct(&payee, &denom, amount)?;
        self.credit(&payer, &denom, amount)?;
        session.status = PaymentStatus::Refunded;
        Ok(())
    }

    pub fn get_balance(&self, address: &str, denom: &str) -> Option<u128> {
        self.accounts.get(address)
            .and_then(|a| a.balances.get(denom).copied())
    }

    pub fn get_session(&self, session_id: &str) -> Option<&PaymentSession> {
        self.sessions.get(session_id)
    }

    pub fn list_sessions(&self, address: &str) -> Vec<&PaymentSession> {
        self.sessions.values()
            .filter(|s| s.payer == address || s.payee == address)
            .collect()
    }

    pub fn set_balance(&mut self, address: &str, denom: &str, amount: u128) {
        if let Some(acct) = self.accounts.get_mut(address) {
            acct.balances.insert(denom.to_string(), amount);
        }
    }

    fn deduct(&mut self, address: &str, denom: &str, amount: u128) -> Result<(), PaymentError> {
        let account = self.accounts.get_mut(address)
            .ok_or(PaymentError::AccountNotFound(address.to_string()))?;
        let balance = account.balances.entry(denom.to_string()).or_insert(0);
        if *balance < amount {
            return Err(PaymentError::InsufficientBalance {
                account: address.to_string(), required: amount, available: *balance,
            });
        }
        *balance -= amount;
        Ok(())
    }

    fn credit(&mut self, address: &str, denom: &str, amount: u128) -> Result<(), PaymentError> {
        let account = self.accounts.get_mut(address)
            .ok_or(PaymentError::AccountNotFound(address.to_string()))?;
        *account.balances.entry(denom.to_string()).or_insert(0) += amount;
        Ok(())
    }
}

#[derive(Debug)]
pub enum PaymentError {
    SessionNotFound(String),
    AccountNotFound(String),
    InsufficientBalance { account: String, required: u128, available: u128 },
    InvalidState { session: String, expected: PaymentStatus, actual: PaymentStatus },
    SessionExpired(String),
}

impl std::fmt::Display for PaymentError {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        match self {
            PaymentError::SessionNotFound(id) => write!(f, "Session {} not found", id),
            PaymentError::AccountNotFound(a) => write!(f, "Account {} not found", a),
            PaymentError::InsufficientBalance { account, required, available } => {
                write!(f, "Account {} insufficient: need {} have {}", account, required, available)
            }
            PaymentError::InvalidState { session, expected, actual } => {
                write!(f, "Session {} invalid state: expected {:?} actual {:?}", session, expected, actual)
            }
            PaymentError::SessionExpired(id) => write!(f, "Session {} expired", id),
        }
    }
}

impl std::error::Error for PaymentError {}

5.2 Python 端 AIPAY

# mock_aipay_py/payments.py
from __future__ import annotations
import time
from dataclasses import dataclass, field
from enum import Enum
from typing import Optional


class PaymentStatus(str, Enum):
    PENDING = "pending"
    AUTHORIZED = "authorized"
    COMPLETED = "completed"
    FAILED = "failed"
    REFUNDED = "refunded"
    EXPIRED = "expired"


class PaymentPurpose(str, Enum):
    AGENT_SERVICE = "agent_service"
    TASK_REWARD = "task_reward"
    SUBSCRIPTION = "subscription"
    SLASH = "slash"
    TRANSFER = "transfer"


@dataclass
class PaymentSession:
    id: str
    payer: str
    payee: str
    amount: int
    denom: str
    purpose: PaymentPurpose
    status: PaymentStatus = PaymentStatus.PENDING
    created_at: int = 0
    expires_at: int = 0
    metadata: dict[str, str] = field(default_factory=dict)
    failure_reason: Optional[str] = None


@dataclass
class AccountBalance:
    address: str
    balances: dict[str, int] = field(default_factory=dict)


class PaymentError(Exception):
    def __init__(self, message, code):
        self.message = message
        self.code = code
        super().__init__(message)


class InsufficientBalanceError(PaymentError):
    def __init__(self, account, required, available):
        super().__init__(
            f"Account {account} insufficient: need {required}, have {available}",
            code="INSUFFICIENT_BALANCE",
        )
        self.account = account
        self.required = required
        self.available = available


class MockAIPAY:

    def __init__(self):
        self._sessions: dict[str, PaymentSession] = {}
        self._accounts: dict[str, AccountBalance] = {}
        self._session_counter = 0

    def register_account(self, address, balances=None):
        if balances is None:
            balances = {"umsg": 1_000_000_000}
        self._accounts[address] = AccountBalance(address=address, balances=balances)

    def get_balance(self, address, denom="umsg"):
        account = self._accounts.get(address)
        if account is None:
            return None
        return account.balances.get(denom, 0)

    def create_session(self, payer, payee, amount, denom="umsg",
                       purpose=PaymentPurpose.AGENT_SERVICE):
        if payer not in self._accounts:
            raise PaymentError(f"Payer {payer} not registered", "PAYER_NOT_FOUND")
        balance = self.get_balance(payer, denom)
        if balance is None or balance < amount:
            raise InsufficientBalanceError(payer, amount, balance or 0)
        self._session_counter += 1
        session_id = f"pay-{self._session_counter:016x}"
        now = int(time.time() * 1000)
        session = PaymentSession(
            id=session_id, payer=payer, payee=payee,
            amount=amount, denom=denom, purpose=purpose,
            status=PaymentStatus.PENDING,
            created_at=now, expires_at=now + 300_000,
        )
        self._sessions[session_id] = session
        return session_id

    def authorize(self, session_id):
        session = self._sessions.get(session_id)
        if session is None:
            raise PaymentError(f"Session {session_id} not found", "SESSION_NOT_FOUND")
        if session.status != PaymentStatus.PENDING:
            raise PaymentError(
                f"Session {session_id} in state {session.status.value}, expected pending",
                "INVALID_STATE",
            )
        balance = self.get_balance(session.payer, session.denom)
        if balance is None or balance < session.amount:
            session.status = PaymentStatus.FAILED
            session.failure_reason = "Insufficient balance at authorization"
            raise InsufficientBalanceError(session.payer, session.amount, balance or 0)
        session.status = PaymentStatus.AUTHORIZED

    def complete(self, session_id):
        session = self._sessions.get(session_id)
        if session is None:
            raise PaymentError(f"Session {session_id} not found", "SESSION_NOT_FOUND")
        if session.status != PaymentStatus.AUTHORIZED:
            raise PaymentError(
                f"Session {session_id} in state {session.status.value}, expected authorized",
                "INVALID_STATE",
            )
        payer_bal = self._accounts[session.payer].balances
        payer_bal[session.denom] = payer_bal.get(session.denom, 0) - session.amount
        payee_bal = self._accounts[session.payee].balances
        payee_bal[session.denom] = payee_bal.get(session.denom, 0) + session.amount
        session.status = PaymentStatus.COMPLETED

    def refund(self, session_id):
        session = self._sessions.get(session_id)
        if session is None:
            raise PaymentError(f"Session {session_id} not found", "SESSION_NOT_FOUND")
        if session.status != PaymentStatus.COMPLETED:
            raise PaymentError(
                f"Session {session_id} in state {session.status.value}, expected completed",
                "INVALID_STATE",
            )
        payee_bal = self._accounts[session.payee].balances
        payee_bal[session.denom] = payee_bal.get(session.denom, 0) - session.amount
        payer_bal = self._accounts[session.payer].balances
        payer_bal[session.denom] = payer_bal.get(session.denom, 0) + session.amount
        session.status = PaymentStatus.REFUNDED

    def get_session(self, session_id):
        return self._sessions.get(session_id)

    def list_sessions(self, address):
        return [
            s for s in self._sessions.values()
            if s.payer == address or s.payee == address
        ]

    def set_balance(self, address, denom, amount):
        if address not in self._accounts:
            self.register_account(address, {denom: amount})
        else:
            self._accounts[address].balances[denom] = amount

    def transfer(self, sender, recipient, amount, denom="umsg"):
        sender_bal = self._accounts.get(sender)
        if sender_bal is None:
            raise PaymentError(f"Sender {sender} not found", "SENDER_NOT_FOUND")
        sender_umsg = sender_bal.balances.get(denom, 0)
        if sender_umsg < amount:
            raise InsufficientBalanceError(sender, amount, sender_umsg)
        sender_bal.balances[denom] = sender_umsg - amount
        recipient_bal = self._accounts.setdefault(
            recipient, AccountBalance(address=recipient),
        )
        recipient_bal.balances[denom] = recipient_bal.balances.get(denom, 0) + amount

    def reset(self):
        self._sessions.clear()
        self._accounts.clear()
        self._session_counter = 0

5.3 支付失败场景

# test_payments.py
from payments import MockAIPAY, PaymentPurpose, PaymentError, InsufficientBalanceError


def setup() -> MockAIPAY:
    pay = MockAIPAY()
    pay.register_account("alice", {"umsg": 10_000_000})
    pay.register_account("bob", {"umsg": 5_000_000})
    pay.register_account("charlie", {"umsg": 100_000_000})
    return pay


def test_successful_payment():
    pay = setup()
    sid = pay.create_session("alice", "bob", 1_000_000)
    pay.authorize(sid)
    pay.complete(sid)
    assert pay.get_session(sid).status.value == "completed"
    assert pay.get_balance("alice") == 9_000_000
    assert pay.get_balance("bob") == 6_000_000
    print("test_successful_payment passed")


def test_insufficient_balance():
    pay = setup()
    try:
        pay.create_session("bob", "alice", 10_000_000)
        assert False
    except InsufficientBalanceError:
        pass
    print("test_insufficient_balance passed")


def test_double_complete():
    pay = setup()
    sid = pay.create_session("alice", "bob", 1_000_000)
    pay.authorize(sid)
    pay.complete(sid)
    try:
        pay.complete(sid)
        assert False
    except PaymentError:
        pass
    print("test_double_complete passed")


def test_refund():
    pay = setup()
    sid = pay.create_session("alice", "bob", 2_000_000)
    pay.authorize(sid)
    pay.complete(sid)
    alice_before = pay.get_balance("alice")
    bob_before = pay.get_balance("bob")
    pay.refund(sid)
    assert pay.get_balance("alice") == alice_before + 2_000_000
    assert pay.get_balance("bob") == bob_before - 2_000_000
    print("test_refund passed")


def test_multiple_concurrent():
    pay = setup()
    sid1 = pay.create_session("alice", "bob", 6_000_000)
    sid2 = pay.create_session("alice", "charlie", 6_000_000)
    pay.authorize(sid1)
    try:
        pay.authorize(sid2)
        assert False
    except InsufficientBalanceError:
        pass
    pay.complete(sid1)
    print("test_multiple_concurrent passed")


if __name__ == "__main__":
    test_successful_payment()
    test_insufficient_balance()
    test_double_complete()
    test_refund()
    test_multiple_concurrent()
    print("All payment tests passed")

6. Agent 行为模拟框架

6.1 场景定义 DSL

# agent_simulation/scenario.py
from __future__ import annotations
from dataclasses import dataclass, field
from enum import Enum
from typing import Any, Optional


class ActionType(str, Enum):
    REGISTER = "register"
    DISCOVER = "discover"
    NEGOTIATE = "negotiate"
    ASSIGN_TASK = "assign_task"
    EXECUTE_TASK = "execute_task"
    MAKE_PAYMENT = "make_payment"
    RECEIVE_PAYMENT = "receive_payment"
    REFUND = "refund"
    UPDATE_STATUS = "update_status"
    SEND_MESSAGE = "send_message"
    WAIT = "wait"
    ASSERT = "assert"
    CUSTOM = "custom"


@dataclass
class Action:
    type: ActionType
    actor: str
    params: dict[str, Any] = field(default_factory=dict)
    delay_before_ms: int = 0
    expected_result: Optional[dict[str, Any]] = None
    timeout_ms: int = 30_000


@dataclass
class Scenario:
    name: str
    description: str
    agents: list[dict[str, Any]]
    actions: list[Action]
    seed: int = 42
    tags: list[str] = field(default_factory=list)
    expected_final_state: Optional[dict[str, Any]] = None


class ScenarioBuilder:

    def __init__(self, name: str):
        self.scenario = Scenario(name=name, description="", agents=[], actions=[])

    def describe(self, desc: str):
        self.scenario.description = desc
        return self

    def with_agent(self, name, address, capabilities,
                   initial_balance=1_000_000_000):
        self.scenario.agents.append({
            "name": name, "address": address,
            "capabilities": capabilities,
            "initial_balance": initial_balance,
        })
        return self

    def with_seed(self, seed: int):
        self.scenario.seed = seed
        return self

    def register(self, actor, capabilities=None):
        self.scenario.actions.append(Action(
            type=ActionType.REGISTER, actor=actor,
            params={"capabilities": capabilities or []},
        ))
        return self

    def discover(self, actor, required_cap, timeout_ms=5000):
        self.scenario.actions.append(Action(
            type=ActionType.DISCOVER, actor=actor,
            params={"required_capability": required_cap},
            timeout_ms=timeout_ms,
        ))
        return self

    def assign_task(self, actor, target, task_type, reward, params=None):
        self.scenario.actions.append(Action(
            type=ActionType.ASSIGN_TASK, actor=actor,
            params={
                "target": target, "task_type": task_type,
                "reward": reward, "parameters": params or {},
            },
        ))
        return self

    def execute_task(self, actor, task_id, result=None, duration_ms=0):
        self.scenario.actions.append(Action(
            type=ActionType.EXECUTE_TASK, actor=actor,
            params={"task_id": task_id, "result": result or {},
                    "duration_ms": duration_ms},
        ))
        return self

    def make_payment(self, actor, recipient, amount, purpose="agent_service"):
        self.scenario.actions.append(Action(
            type=ActionType.MAKE_PAYMENT, actor=actor,
            params={"recipient": recipient, "amount": amount, "purpose": purpose},
        ))
        return self

    def send_message(self, actor, recipient, msg_type, payload):
        self.scenario.actions.append(Action(
            type=ActionType.SEND_MESSAGE, actor=actor,
            params={"recipient": recipient, "msg_type": msg_type, "payload": payload},
        ))
        return self

    def wait(self, actor, duration_ms):
        self.scenario.actions.append(Action(
            type=ActionType.WAIT, actor=actor,
            params={"duration_ms": duration_ms},
        ))
        return self

    def build(self):
        return self.scenario


def marketplace_scenario():
    return (
        ScenarioBuilder("marketplace_flow")
        .describe("完整的 Agent 市场交易流程")
        .with_agent("alice", "msg1alice", ["task_publisher"], 10_000_000_000)
        .with_agent("bob", "msg1bob", ["data_analysis"], 5_000_000_000)
        .with_seed(42)
        .register("alice", ["task_publisher"])
        .register("bob", ["data_analysis"])
        .discover("alice", "data_analysis")
        .assign_task("alice", "bob", "data_analysis", 1_000_000, {
            "dataset": "market_data_q2",
        })
        .make_payment("alice", "bob", 1_000_000)
        .build()
    )


def negotiation_scenario():
    return (
        ScenarioBuilder("negotiation_flow")
        .describe("Agent 间价格谈判")
        .with_agent("alice", "msg1alice", ["requester"])
        .with_agent("bob", "msg1bob", ["compute"], 3_000_000_000)
        .register("alice", ["requester"])
        .register("bob", ["compute"])
        .discover("alice", "compute")
        .send_message("alice", "bob", "a2a/negotiate/start", {
            "task": "model_training", "initial_offer": 500_000,
        })
        .send_message("bob", "alice", "a2a/negotiate/counter", {
            "counter_offer": 800_000,
        })
        .send_message("alice", "bob", "a2a/negotiate/accept", {
            "final_price": 800_000,
        })
        .make_payment("alice", "bob", 800_000)
        .build()
    )

6.2 模拟引擎

# agent_simulation/engine.py
from __future__ import annotations
import json
import time
import asyncio
from dataclasses import dataclass, field
from enum import Enum
from typing import Any, Optional

from mock_registry_py.bridge import MockRegistry, AgentCapability, AgentStatus
from mock_a2a_py.bus import MockA2ABus
from mock_aipay_py.payments import MockAIPAY, PaymentPurpose

from scenario import Scenario, Action, ActionType


class SimulationStatus(str, Enum):
    PENDING = "pending"
    RUNNING = "running"
    COMPLETED = "completed"
    FAILED = "failed"
    TIMEOUT = "timeout"


@dataclass
class SimulationMetrics:
    total_actions: int = 0
    completed_actions: int = 0
    failed_actions: int = 0
    total_duration_ms: float = 0.0
    messages_sent: int = 0
    messages_delivered: int = 0
    messages_dropped: int = 0
    payments_created: int = 0
    payments_completed: int = 0
    payments_failed: int = 0
    agents_registered: int = 0
    steps: list[dict] = field(default_factory=list)


@dataclass
class SimulationResult:
    status: SimulationStatus
    metrics: SimulationMetrics
    error: Optional[str] = None
    final_state: Optional[dict] = None


class AgentSimulation:

    def __init__(self, scenario: Scenario):
        self.scenario = scenario
        self.registry = MockRegistry()
        self.a2a_bus = MockA2ABus(seed=scenario.seed)
        self.aipay = MockAIPAY()
        self.metrics = SimulationMetrics()
        self._agent_addresses: dict[str, str] = {}

    def setup(self):
        for agent_cfg in self.scenario.agents:
            addr = agent_cfg["address"]
            balance = agent_cfg.get("initial_balance", 1_000_000_000)
            self.aipay.register_account(addr, {"umsg": balance})
            self._agent_addresses[agent_cfg["name"]] = addr

    async def run(self) -> SimulationResult:
        self.metrics = SimulationMetrics()
        start_time = time.time()
        try:
            self.setup()
            for action in self.scenario.actions:
                if action.delay_before_ms > 0:
                    await asyncio.sleep(action.delay_before_ms / 1000.0)
                try:
                    await self._execute_action(action)
                    self.metrics.completed_actions += 1
                except Exception as e:
                    self.metrics.failed_actions += 1
                    raise
                self.metrics.total_actions += 1
                self.metrics.steps.append({
                    "action_index": len(self.metrics.steps),
                    "action_type": action.type.value,
                    "actor": action.actor,
                })

            final_state = self._capture_state()
            duration = (time.time() - start_time) * 1000
            self.metrics.total_duration_ms = duration
            a2a_stats = self.a2a_bus.stats()
            self.metrics.messages_sent = a2a_stats.sent
            self.metrics.messages_delivered = a2a_stats.delivered
            self.metrics.messages_dropped = a2a_stats.dropped

            return SimulationResult(
                status=SimulationStatus.COMPLETED,
                metrics=self.metrics,
                final_state=final_state,
            )
        except Exception as e:
            duration = (time.time() - start_time) * 1000
            self.metrics.total_duration_ms = duration
            return SimulationResult(
                status=SimulationStatus.FAILED,
                metrics=self.metrics,
                error=str(e),
            )

    async def _execute_action(self, action: Action):
        actor_addr = self._agent_addresses.get(action.actor)
        if not actor_addr:
            raise ValueError(f"Unknown actor: {action.actor}")

        if action.type == ActionType.REGISTER:
            caps = [
                AgentCapability(name=c, version="1.0.0")
                for c in action.params.get("capabilities", [])
            ]
            self.registry.register_agent(owner=actor_addr,
                name=f"{action.actor}_agent", capabilities=caps)
            self.metrics.agents_registered += 1

        elif action.type == ActionType.DISCOVER:
            cap = action.params.get("required_capability", "")
            agents = self.registry.find_by_capability(cap)
            if not agents and action.expected_result is None:
                raise RuntimeError(f"No agent found with capability '{cap}'")

        elif action.type == ActionType.ASSIGN_TASK:
            target = action.params.get("target", "")
            target_addr = self._agent_addresses.get(target)
            reward = action.params.get("reward", 0)
            msg = json.dumps({
                "task_type": action.params.get("task_type", ""),
                "reward": reward,
                "parameters": action.params.get("parameters", {}),
            }).encode()
            self.a2a_bus.send(sender=actor_addr, recipient=target_addr,
                             msg_type="a2a/task/assign", payload=msg)

        elif action.type == ActionType.SEND_MESSAGE:
            recipient = action.params.get("recipient", "")
            recipient_addr = self._agent_addresses.get(recipient)
            msg_type = action.params.get("msg_type", "")
            payload = action.params.get("payload", {})
            self.a2a_bus.send(sender=actor_addr, recipient=recipient_addr,
                             msg_type=msg_type,
                             payload=json.dumps(payload).encode())

        elif action.type == ActionType.MAKE_PAYMENT:
            recipient = action.params.get("recipient", "")
            recipient_addr = self._agent_addresses.get(recipient)
            amount = int(action.params.get("amount", 0))
            try:
                sid = self.aipay.create_session(
                    payer=actor_addr, payee=recipient_addr,
                    amount=amount, purpose=PaymentPurpose.AGENT_SERVICE,
                )
                self.aipay.authorize(sid)
                self.aipay.complete(sid)
                self.metrics.payments_completed += 1
            except Exception:
                self.metrics.payments_failed += 1
                raise
            finally:
                self.metrics.payments_created += 1

        elif action.type == ActionType.WAIT:
            duration_ms = action.params.get("duration_ms", 0)
            await asyncio.sleep(duration_ms / 1000.0)

    def _capture_state(self) -> dict:
        return {
            "registry": {
                "total_agents": len(self.registry.list_agents()),
                "active": self.registry.stats().active,
            },
            "a2a_bus": {
                "sent": self.a2a_bus.sent_count,
                "delivered": self.a2a_bus.delivered_count,
                "dropped": self.a2a_bus.dropped_count,
            },
            "aipay": {
                "accounts": {
                    addr: self.aipay.get_balance(addr)
                    for addr in self._agent_addresses.values()
                },
            },
        }

    def summary(self, result: SimulationResult) -> str:
        m = result.metrics
        lines = [
            f"场景: {self.scenario.name}",
            f"状态: {result.status.value}",
            f"耗时: {m.total_duration_ms:.1f}ms",
            f"动作: {m.completed_actions}/{m.total_actions} 完成",
            f"注册: {m.agents_registered} Agent",
            f"A2A: {m.messages_sent} 发送 / {m.messages_delivered} 投递",
            f"支付: {m.payments_completed} 成功 / {m.payments_failed} 失败",
        ]
        if result.error:
            lines.append(f"错误: {result.error}")
        return "\n".join(lines)

6.3 完整运行示例

# example_marketplace.py
import asyncio
import json
from scenario import marketplace_scenario, negotiation_scenario
from engine import AgentSimulation


async def run_marketplace():
    print("=" * 60)
    print("场景: Agent 市场交易流程")
    print("=" * 60)
    scenario = marketplace_scenario()
    sim = AgentSimulation(scenario)
    result = await sim.run()
    print(sim.summary(result))
    if result.status.value == "completed":
        print("场景执行成功")
    print(json.dumps(result.final_state, indent=2, ensure_ascii=False))


async def run_negotiation():
    print("=" * 60)
    print("场景: Agent 协商谈判")
    print("=" * 60)
    scenario = negotiation_scenario()
    sim = AgentSimulation(scenario)
    result = await sim.run()
    print(sim.summary(result))
    if result.status.value == "completed":
        print("场景执行成功")


async def run_batch():
    scenarios = [
        ("Marketplace", marketplace_scenario()),
        ("Negotiation", negotiation_scenario()),
    ]
    results = []
    for name, scenario in scenarios:
        sim = AgentSimulation(scenario)
        result = await sim.run()
        results.append((name, result))

    passed = 0
    for name, result in results:
        status = "PASS" if result.status.value == "completed" else "FAIL"
        print(f"  {status} {name:20s} | {result.metrics.total_duration_ms:8.1f}ms | "
              f"{result.metrics.completed_actions}/{result.metrics.total_actions}")
        if result.status.value == "completed":
            passed += 1
    print(f"通过: {passed}/{len(scenarios)}")


async def main():
    await run_marketplace()
    print()
    await run_negotiation()
    print()
    await run_batch()


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

6.4 确定性保证

# agent_simulation/determinism.py
import json
import hashlib


class DeterministicRNG:

    def __init__(self, seed: int = 42):
        self._state = seed
        self._seed = seed

    def reset(self):
        self._state = self._seed

    def next_u64(self):
        self._state = (
            self._state * 6364136223846793005 + 1442695040888963407
        ) & 0xFFFFFFFFFFFFFFFF
        return self._state

    def next_f64(self):
        return self.next_u64() / 0xFFFFFFFFFFFFFFFF

    def next_int(self, low, high):
        return low + self.next_u64() % (high - low)

    def shuffle(self, items):
        result = list(items)
        for i in range(len(result) - 1, 0, -1):
            j = self.next_int(0, i + 1)
            result[i], result[j] = result[j], result[i]
        return result


def deterministic_hash(data: dict) -> str:
    serialized = json.dumps(data, sort_keys=True, ensure_ascii=False)
    return hashlib.sha256(serialized.encode()).hexdigest()


def verify_determinism(scenario_fn, runs: int = 3):
    import asyncio
    from engine import AgentSimulation

    hashes = []
    for _ in range(runs):
        scenario = scenario_fn()
        sim = AgentSimulation(scenario)
        result = asyncio.run(sim.run())
        trace = {
            "metrics": {
                "completed": result.metrics.completed_actions,
                "failed": result.metrics.failed_actions,
                "messages_delivered": result.metrics.messages_delivered,
                "messages_dropped": result.metrics.messages_dropped,
            },
            "final_state": result.final_state,
        }
        hashes.append(deterministic_hash(trace))

    return all(h == hashes[0] for h in hashes)

7. 宪法与策略验证

7.1 宪法规则引擎

# constitution/engine.py
from __future__ import annotations
from dataclasses import dataclass, field
from enum import Enum
from typing import Any, Callable, Optional


class ViolationSeverity(str, Enum):
    WARNING = "warning"
    MINOR = "minor"
    MAJOR = "major"
    CRITICAL = "critical"


@dataclass
class Violation:
    rule_id: str
    rule_name: str
    severity: ViolationSeverity
    description: str
    actor: str
    action: str
    details: dict[str, Any] = field(default_factory=dict)


@dataclass
class ConstitutionRule:
    id: str
    name: str
    description: str
    severity: ViolationSeverity
    check_fn: Callable[['ConstitutionEngine', dict], Optional[Violation]]


class ConstitutionEngine:

    def __init__(self):
        self._rules: list[ConstitutionRule] = []
        self._violations: list[Violation] = []
        self._register_default_rules()

    def _register_default_rules(self):
        self.add_rule(ConstitutionRule(
            id="C001", name="Unique Agent Address",
            description="每个 Agent 必须有唯一的链上地址",
            severity=ViolationSeverity.MAJOR,
            check_fn=self._check_unique_address,
        ))
        self.add_rule(ConstitutionRule(
            id="C002", name="Registration Integrity",
            description="Agent 注册信息不可篡改核心字段",
            severity=ViolationSeverity.CRITICAL,
            check_fn=self._check_registration_integrity,
        ))
        self.add_rule(ConstitutionRule(
            id="C003", name="Payment Transparency",
            description="所有 Agent 间支付必须明确记录金额和用途",
            severity=ViolationSeverity.MAJOR,
            check_fn=self._check_payment_transparency,
        ))
        self.add_rule(ConstitutionRule(
            id="C004", name="No Malicious Messages",
            description="禁止发送包含恶意载荷的 A2A 消息",
            severity=ViolationSeverity.CRITICAL,
            check_fn=self._check_malicious_messages,
        ))
        self.add_rule(ConstitutionRule(
            id="C005", name="Minimum Delegation",
            description="高风险操作需要足够的委托层级",
            severity=ViolationSeverity.MINOR,
            check_fn=self._check_delegation_level,
        ))
        self.add_rule(ConstitutionRule(
            id="C006", name="Gas Fee Payment",
            description="Agent 必须有足够余额支付燃气费",
            severity=ViolationSeverity.MAJOR,
            check_fn=self._check_gas_fee,
        ))
        self.add_rule(ConstitutionRule(
            id="C007", name="Rate Limiting",
            description="单个区块内操作次数受限制",
            severity=ViolationSeverity.WARNING,
            check_fn=self._check_rate_limit,
        ))

    def add_rule(self, rule: ConstitutionRule):
        self._rules.append(rule)

    def validate(self, action: str, actor: str,
                 context: dict[str, Any]) -> list[Violation]:
        violations = []
        for rule in self._rules:
            violation = rule.check_fn(self, {
                "action": action, "actor": actor, "context": context,
            })
            if violation:
                violations.append(violation)
                self._violations.append(violation)
        return violations

    def validate_state(self, state: dict) -> list[Violation]:
        violations = []
        agents = state.get("registry", {}).get("list", [])
        addresses_seen = set()
        for agent in agents:
            addr = agent.get("address")
            if addr in addresses_seen:
                violations.append(Violation(
                    rule_id="C001", rule_name="Unique Agent Address",
                    severity=ViolationSeverity.MAJOR,
                    description=f"Duplicate agent address: {addr}",
                    actor=agent.get("owner", "unknown"),
                    action="register", details={"address": addr},
                ))
            addresses_seen.add(addr)
        return violations

    def get_violations(self, severity: Optional[ViolationSeverity] = None):
        if severity:
            return [v for v in self._violations if v.severity == severity]
        return list(self._violations)

    def clear(self):
        self._violations.clear()

    def report(self) -> str:
        if not self._violations:
            return "所有宪法规则通过,无违规"
        lines = ["宪法合规性报告", "=" * 40]
        for sev in ViolationSeverity:
            items = [v for v in self._violations if v.severity == sev]
            if items:
                lines.append(f"[{sev.value.upper()}] ({len(items)} 条)")
                for v in items:
                    lines.append(f"  [{v.rule_id}] {v.description}")
        return "\n".join(lines)

    def _check_unique_address(self, context):
        return None

    def _check_registration_integrity(self, context):
        return None

    def _check_payment_transparency(self, context):
        return None

    def _check_malicious_messages(self, context):
        return None

    def _check_delegation_level(self, context):
        return None

    def _check_gas_fee(self, context):
        return None

    def _check_rate_limit(self, context):
        return None

7.2 策略验证测试

# test_constitution.py
from constitution.engine import ConstitutionEngine, ViolationSeverity


def test_no_violations():
    engine = ConstitutionEngine()
    violations = engine.validate("register", "agent_alice", {})
    assert len(violations) == 0
    print("test_no_violations passed")


def test_violation_tracking():
    engine = ConstitutionEngine()
    engine.validate("register", "agent_alice", {})
    engine.validate("make_payment", "agent_bob", {})
    assert len(engine.get_violations()) == 0
    print("test_violation_tracking passed")


def test_state_validation():
    engine = ConstitutionEngine()
    state = {
        "registry": {
            "list": [
                {"address": "msg1...", "owner": "alice"},
                {"address": "msg1...", "owner": "bob"},
            ],
        },
    }
    violations = engine.validate_state(state)
    assert len(violations) == 0
    print("test_state_validation passed")


def test_duplicate_detection():
    engine = ConstitutionEngine()
    state = {
        "registry": {
            "list": [
                {"address": "msg1dup", "owner": "alice"},
                {"address": "msg1dup", "owner": "bob"},
            ],
        },
    }
    violations = engine.validate_state(state)
    assert len(violations) == 1
    assert violations[0].rule_id == "C001"
    print("test_duplicate_detection passed")


def test_report_generation():
    engine = ConstitutionEngine()
    state = {
        "registry": {
            "list": [
                {"address": "msg1dup", "owner": "alice"},
                {"address": "msg1dup", "owner": "bob"},
            ],
        },
    }
    engine.validate_state(state)
    report = engine.report()
    assert "C001" in report
    assert "Unique Agent Address" in report
    print("test_report_generation passed")


def test_severity_filtering():
    engine = ConstitutionEngine()
    state = {
        "registry": {
            "list": [
                {"address": "msg1x", "owner": "a"},
                {"address": "msg1x", "owner": "b"},
                {"address": "msg1x", "owner": "c"},
            ],
        },
    }
    engine.validate_state(state)
    major = engine.get_violations(ViolationSeverity.MAJOR)
    assert len(major) == 3
    critical = engine.get_violations(ViolationSeverity.CRITICAL)
    assert len(critical) == 0
    print("test_severity_filtering passed")


def test_clear():
    engine = ConstitutionEngine()
    engine.validate("test", "actor", {})
    engine.clear()
    assert len(engine.get_violations()) == 0
    print("test_clear passed")


if __name__ == "__main__":
    test_no_violations()
    test_violation_tracking()
    test_state_validation()
    test_duplicate_detection()
    test_report_generation()
    test_severity_filtering()
    test_clear()
    print("All constitution tests passed")

8. CI/CD 集成

8.1 GitHub Actions 沙箱设置

# .github/workflows/agent-sandbox-tests.yml
name: Agent Sandbox Tests

on:
  push:
    branches: [main, develop]
  pull_request:
    branches: [main]
  workflow_dispatch:

env:
  CHAIN_ID: msg-chain-1
  MSGD_VERSION: v1.2.0
  GO_VERSION: "1.21"
  PYTHON_VERSION: "3.11"

jobs:
  sandbox-test:
    runs-on: ubuntu-latest
    strategy:
      matrix:
        test-suite:
          - registry
          - a2a-communication
          - aipay-payments
          - agent-simulation
          - constitution
      fail-fast: false

    services:
      msg-sandbox:
        image: msgchain/msgd:${{ env.MSGD_VERSION }}
        ports:
          - 26657:26657
          - 9090:9090
        options: >-
          --health-cmd "curl -f http://localhost:26657/status || exit 1"
          --health-interval 5s
          --health-retries 15

    steps:
      - uses: actions/checkout@v4

      - name: Setup Go
        uses: actions/setup-go@v5
        with:
          go-version: ${{ env.GO_VERSION }}

      - name: Setup Python
        uses: actions/setup-python@v5
        with:
          python-version: ${{ env.PYTHON_VERSION }}

      - name: Setup Rust
        uses: actions-rust-lang/setup-rust-toolchain@v1
        with:
          toolchain: stable

      - name: Install msgd
        run: |
          wget -q https://github.com/msgchain/msgd/releases/download/${{ env.MSGD_VERSION }}/msgd-linux-amd64.tar.gz
          tar xzf msgd-linux-amd64.tar.gz
          sudo mv msgd /usr/local/bin/
          msgd version

      - name: Initialize Sandbox Node
        run: |
          msgd init ci-sandbox --chain-id ${{ env.CHAIN_ID }} --default-denom umsg
          msgd keys add faucet --keyring-backend test
          msgd add-genesis-account $(msgd keys show faucet -a --keyring-backend test) 1000000000000000umsg
          msgd gentx faucet 1000000000umsg --chain-id ${{ env.CHAIN_ID }} --keyring-backend test
          msgd collect-gentxs
          sed -i 's/minimum-gas-prices = ""/minimum-gas-prices = "0umsg"/g' ~/.msgd/config/app.toml
          msgd start --minimum-gas-prices 0umsg &
          sleep 5

      - name: Create Test Accounts
        run: |
          for name in agent_alice agent_bob agent_charlie trader_dave oracle_eve; do
            msgd keys add $name --keyring-backend test
            addr=$(msgd keys show $name -a --keyring-backend test)
            msgd tx bank send faucet $addr 1000000000000umsg \
              --chain-id ${{ env.CHAIN_ID }} \
              --keyring-backend test \
              -y
          done

      - name: Install Python Dependencies
        run: |
          pip install pytest pytest-asyncio

      - name: Run Rust Tests
        if: matrix.test-suite == 'registry'
        working-directory: ./mock_registry
        run: cargo test -- --nocapture

      - name: Run A2A Tests
        if: matrix.test-suite == 'a2a-communication'
        run: python -m pytest mock_a2a_py/test_a2a.py -v

      - name: Run AIPAY Tests
        if: matrix.test-suite == 'aipay-payments'
        run: python -m pytest mock_aipay_py/test_payments.py -v

      - name: Run Agent Simulation Tests
        if: matrix.test-suite == 'agent-simulation'
        run: |
          python -m pytest agent_simulation/example_marketplace.py -v

      - name: Run Constitution Tests
        if: matrix.test-suite == 'constitution'
        run: python -m pytest constitution/test_constitution.py -v

      - name: Upload Test Logs
        if: always()
        uses: actions/upload-artifact@v4
        with:
          name: test-logs-${{ matrix.test-suite }}
          path: |
            ~/.msgd/config/
            /tmp/msg-sandbox-*.log
          retention-days: 7

  full-integration:
    runs-on: ubuntu-latest
    needs: [sandbox-test]
    if: github.ref == 'refs/heads/main'

    steps:
      - uses: actions/checkout@v4

      - name: Full Integration Test
        run: |
          pip install pytest pytest-asyncio
          python -m pytest tests/integration/ -v --timeout=120

      - name: Run Stress Test
        run: |
          python tests/stress/many_agents.py --count 50 --rounds 10

8.2 Makefile 辅助

# Makefile - 沙箱测试快捷命令
.PHONY: help sandbox-up sandbox-down sandbox-reset test-all test-registry \
        test-a2a test-aipay test-simulation test-constitution ci

help:
	@echo "沙箱测试命令"
	@echo "  make sandbox-up    启动本地沙箱节点"
	@echo "  make sandbox-down  停止沙箱节点"
	@echo "  make sandbox-reset 重置沙箱状态"
	@echo "  make test-all      运行所有测试"
	@echo "  make test-registry 运行注册表测试"
	@echo "  make test-a2a      运行 A2A 通信测试"
	@echo "  make test-aipay    运行 AIPAY 支付测试"
	@echo "  make test-simulation 运行模拟测试"
	@echo "  make test-constitution 运行宪法验证测试"
	@echo "  make ci            完整 CI 流程"

SANDBOX_DATA = ~/.msgd
CHAIN_ID = msg-chain-1

sandbox-up:
	msgd init sandbox-node --chain-id $(CHAIN_ID) --default-denom umsg
	msgd keys add faucet --keyring-backend test 2>/dev/null || true
	msgd add-genesis-account $$(msgd keys show faucet -a --keyring-backend test) 1000000000000000umsg
	msgd gentx faucet 1000000000umsg --chain-id $(CHAIN_ID) --keyring-backend test
	msgd collect-gentxs
	sed -i 's/minimum-gas-prices = ""/minimum-gas-prices = "0umsg"/g' $(SANDBOX_DATA)/config/app.toml
	nohup msgd start --minimum-gas-prices 0umsg > /tmp/msgd.log 2>&1 &
	sleep 3
	@echo "Sandbox node started"

sandbox-down:
	pkill msgd 2>/dev/null || true
	@echo "Sandbox node stopped"

sandbox-reset: sandbox-down
	rm -rf $(SANDBOX_DATA)/data
	msgd unsafe-reset-all 2>/dev/null
	@echo "Sandbox reset"

test-registry:
	cd mock_registry && cargo test -- --nocapture

test-a2a:
	python -m pytest mock_a2a_py/test_a2a.py -v

test-aipay:
	python -m pytest mock_aipay_py/test_payments.py -v

test-simulation:
	python -m pytest agent_simulation/example_marketplace.py -v

test-constitution:
	python -m pytest constitution/test_constitution.py -v

test-all: test-registry test-a2a test-aipay test-simulation test-constitution
	@echo "All tests passed"

ci: sandbox-up test-all sandbox-down
	@echo "CI pipeline completed"

8.3 Docker 测试镜像

# Dockerfile.test
FROM msgchain/msgd:v1.2.0 AS msgd-base

FROM python:3.11-slim

RUN apt-get update && apt-get install -y --no-install-recommends \
    curl jq make git && \
    rm -rf /var/lib/apt/lists/*

COPY --from=msgd-base /usr/local/bin/msgd /usr/local/bin/msgd

RUN pip install pytest pytest-asyncio

WORKDIR /workspace
COPY . .

RUN msgd init ci-sandbox --chain-id msg-chain-1 --default-denom umsg && \
    msgd keys add faucet --keyring-backend test && \
    msgd add-genesis-account $$(msgd keys show faucet -a --keyring-backend test) 1000000000000000umsg && \
    msgd gentx faucet 1000000000umsg --chain-id msg-chain-1 --keyring-backend test && \
    msgd collect-gentxs && \
    sed -i 's/minimum-gas-prices = ""/minimum-gas-prices = "0umsg"/g' ~/.msgd/config/app.toml

CMD ["make", "test-all"]

8.4 测试指标收集

# collect_metrics.sh - 沙箱测试指标收集
#!/bin/bash

REPORT_DIR="/tmp/msg-sandbox-metrics/$(date +%Y%m%d-%H%M%S)"
mkdir -p "$REPORT_DIR"

echo "Collecting sandbox test metrics..."

# 1. 链状态
msgd status --node tcp://localhost:26657 > "$REPORT_DIR/chain_status.json" 2>&1

# 2. 验证者信息
msgd query staking validators --node tcp://localhost:26657 -o json > "$REPORT_DIR/validators.json" 2>&1

# 3. 区块信息
msgd query block --node tcp://localhost:26657 > "$REPORT_DIR/latest_block.json" 2>&1

# 4. 测试结果汇总
{
    echo "Test Results Summary"
    echo "===================="
    echo "Timestamp: $(date -u +%Y-%m-%dT%H:%M:%SZ)"
    echo ""
    echo "Registry Tests:"
    cargo test --manifest-path mock_registry/Cargo.toml -- --nocapture 2>&1 | tail -5
    echo ""
    echo "A2A Tests:"
    python -m pytest mock_a2a_py/test_a2a.py -v 2>&1 | tail -5
    echo ""
    echo "AIPAY Tests:"
    python -m pytest mock_aipay_py/test_payments.py -v 2>&1 | tail -5
    echo ""
    echo "Simulation Tests:"
    python -m pytest agent_simulation/example_marketplace.py -v 2>&1 | tail -5
    echo ""
    echo "Constitution Tests:"
    python -m pytest constitution/test_constitution.py -v 2>&1 | tail -5
} > "$REPORT_DIR/test_summary.txt"

echo "Metrics saved to: $REPORT_DIR"

8.5 最佳实践

# 沙箱测试最佳实践

## 1. 测试独立性
每个测试用例应创建自己的 MockRegistry / MockA2ABus / MockAIPAY 实例,
避免测试间状态干扰。

## 2. 种子选择
- 使用确定性种子 (seed=42) 用于调试
- 使用随机种子 (seed=None) 用于探索性测试
- CI 中固定种子以确保可复现

## 3. 故障注入
- 消息丢失: drop_rate=0.1 ~ 0.3 测试容错
- 支付失败: 操作 balance 模拟余额不足
- 网络延迟: base_delay_ms=500 模拟慢速环境

## 4. 断言策略
- 验证最终状态余额
- 验证消息投递计数
- 验证违规记录完整性

## 5. CI 集成
- 并行运行各测试套件以加快反馈
- 收集失败时的完整日志
- 定期清理沙箱数据目录