AI Agent 沙箱与模拟测试环境指南
适用链: msg-chain-1 | 地址前缀:
msg| 版本: v1.0⚠️ No-Go Disclaimer: MSGChain 主网裁决为 No-Go。本文件所有内容反映的是开发阶段的技术设计,不代表主网未独立核验上线状态。生产部署状态请以白皮书为准:https://msgchain.org/whitepaper/
目录
1. 概述
1.1 为什么需要沙箱测试
AI Agent 在区块链上运行涉及不可逆的交易和真实资产。沙箱测试环境允许开发者在无风险条件下验证 Agent 行为、发现漏洞并优化策略。在 MSG Chain 上构建 Agent 时,需要测试以下关键维度:
- 链交互正确性: Agent 发起的交易是否按预期执行
- Agent 间通信: A2A 消息能否可靠送达、正确解析
- 支付与结算: AIPAY 协议在 Agent 场景下的支付流是否完整
- 策略健壮性: 面对网络延迟、消息丢失等异常时的表现
- 宪法合规性: 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 | | |
| | +----------+ +------------------+ | |
| +----------------------------------------+ |
+------------------------------------------------------+
- 底层:
msgd本地链节点,提供完整的 Cosmos SDK / CosmWasm 执行环境 - 中间层: Mock 服务,模拟链上合约的轻量替代品,支持确定性测试
- 上层: Agent 运行时,加载配置并执行场景定义
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 集成
- 并行运行各测试套件以加快反馈
- 收集失败时的完整日志
- 定期清理沙箱数据目录
