AI Agent 端到端实战教程:从零构建一个完整Agent
适用链: MSG Chain (chain-id:
msg-chain-1)
地址前缀:msg
代币: MSG
目标: 从空目录开始,一步步构建一个具备链上注册、A2A通信、支付结算、数据持久化、生产部署能力的完整AI Agent⚠️ No-Go Disclaimer: MSGChain 主网裁决为 No-Go。本文件所有内容反映的是开发阶段的技术设计,不代表主网未独立核验上线状态。生产部署状态请以白皮书为准:https://msgchain.org/whitepaper/
第1章:项目初始化
1.1 环境准备
确保已安装 Python 3.10+ 和 Go 1.24.6+(用于 msgd 命令行工具)。
# 检查版本
python3 --version
go version
# 安装 msgd(MSG Chain CLI)
git clone https://github.com/msgchain/msgchain.git
cd msgchain
make install
msgd version
1.2 创建项目结构
# 创建项目根目录
mkdir -p ~/my-first-agent && cd ~/my-first-agent
# 项目目录结构
mkdir -p src/messages
mkdir -p src/handlers
mkdir -p src/storage
mkdir -p src/monitor
mkdir -p tests
mkdir -p data
mkdir -p scripts
mkdir -p deploy
# 创建虚拟环境
python3 -m venv venv
source venv/bin/activate
1.3 安装依赖
cat > requirements.txt << 'REQUIREMENTS'
msg-sdk>=0.4.0
aiohttp>=3.9.0
pyyaml>=6.0
sqlite-utils>=3.36
pytest>=7.4
pytest-asyncio>=0.21
httpx>=0.25.0
cosmospy>=7.0.0
bech32>=1.2.0
cryptography>=41.0.0
rich>=13.7.0
pydantic>=2.5.0
aiofiles>=23.2.0
REQUIREMENTS
pip install -r requirements.txt
1.4 项目配置文件
cat > config.yaml << 'EOF'
# ============================================
# MyFirstAgent - MSG Chain AI Agent 配置文件
# ============================================
chain:
chain_id: "msg-chain-1"
rpc_endpoint: "https://rpc-testnet.msgchain.org"
rest_endpoint: "https://api-testnet.msgchain.org"
grpc_endpoint: "https://grpc-testnet.msgchain.org:9090"
denom: "umsg"
gas_price: "1000000000"
gas_limit: 200000
agent:
name: "MyFirstAgent"
version: "1.0.0"
description: "支持文本摘要与数据分析的通用AI Agent"
creator: ""
did: ""
did_document: {}
capabilities:
- name: "text-summarization"
description: "文本摘要生成"
input_schema:
type: "object"
properties:
text:
type: "string"
description: "待摘要的文本"
max_length:
type: "integer"
description: "最大摘要长度(字符数)"
default: 200
price_per_request: 100
- name: "data-analysis"
description: "基本数据分析"
input_schema:
type: "object"
properties:
data:
type: "object"
description: "待分析的结构化数据"
analysis_type:
type: "string"
enum: ["summary", "trend", "outlier"]
default: "summary"
price_per_request: 200
pricing:
base_fee: 100
per_compute_unit: 1
currency: "umsg"
payment_timeout_seconds: 300
storage:
db_path: "data/agent.db"
log_path: "data/agent.log"
max_history_days: 90
server:
host: "0.0.0.0"
port: 8080
max_request_size: 1048576
request_timeout_seconds: 30
logging:
level: "INFO"
format: "json"
output: ["stdout", "file"]
metrics:
enabled: true
push_interval_seconds: 60
health_report_interval_seconds: 300
registry:
contract_address: "msg1registryxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"
did_registry_address: "msg1didxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"
EOF
1.5 钱包创建
# 使用 msgd 创建新密钥
msgd keys add agent-key --keyring-backend test
# 记录助记词并导出环境变量
export MNEMONIC="your mnemonic phrase here"
export AGENT_ADDRESS="msg1abc...xyz"
# 领测试网代币
curl -X POST https://faucet-testnet.msgchain.org/claim \
-H "Content-Type: application/json" \
-d '{"address": "'"$AGENT_ADDRESS"'"}'
# 验证余额
msgd query bank balances "$AGENT_ADDRESS" --node https://rpc-testnet.msgchain.org
1.6 项目入口
cat > src/__init__.py << 'EOF'
from src.agent import MyFirstAgent
from src.config import load_config
__version__ = "1.0.0"
__all__ = ["MyFirstAgent", "load_config"]
EOF
cat > src/config.py << 'PYEOF'
import os
import yaml
from pathlib import Path
from typing import Any, Dict, List, Optional
from pydantic import BaseModel, Field
class ChainConfig(BaseModel):
chain_id: str = "msg-chain-1"
rpc_endpoint: str = "https://rpc-testnet.msgchain.org"
rest_endpoint: str = "https://api-testnet.msgchain.org"
grpc_endpoint: str = "https://grpc-testnet.msgchain.org:9090"
denom: str = "umsg"
gas_price: str = "1000000000"
gas_limit: int = 200000
class CapabilityConfig(BaseModel):
name: str
description: str = ""
input_schema: Dict[str, Any] = {}
price_per_request: int = 100
class PricingConfig(BaseModel):
base_fee: int = 100
per_compute_unit: int = 1
currency: str = "umsg"
payment_timeout_seconds: int = 300
class StorageConfig(BaseModel):
db_path: str = "data/agent.db"
log_path: str = "data/agent.log"
max_history_days: int = 90
class ServerConfig(BaseModel):
host: str = "0.0.0.0"
port: int = 8080
max_request_size: int = 1048576
request_timeout_seconds: int = 30
class LoggingConfig(BaseModel):
level: str = "INFO"
format: str = "json"
output: List[str] = ["stdout", "file"]
class MetricsConfig(BaseModel):
enabled: bool = True
push_interval_seconds: int = 60
health_report_interval_seconds: int = 300
class RegistryConfig(BaseModel):
contract_address: str = ""
did_registry_address: str = ""
class AgentMeta(BaseModel):
name: str = "MyFirstAgent"
version: str = "1.0.0"
description: str = ""
creator: str = ""
did: str = ""
did_document: Dict[str, Any] = {}
class AppConfig(BaseModel):
chain: ChainConfig = Field(default_factory=ChainConfig)
agent: AgentMeta = Field(default_factory=AgentMeta)
capabilities: List[CapabilityConfig] = []
pricing: PricingConfig = Field(default_factory=PricingConfig)
storage: StorageConfig = Field(default_factory=StorageConfig)
server: ServerConfig = Field(default_factory=ServerConfig)
logging: LoggingConfig = Field(default_factory=LoggingConfig)
metrics: MetricsConfig = Field(default_factory=MetricsConfig)
registry: RegistryConfig = Field(default_factory=RegistryConfig)
def load_config(path: str = "config.yaml") -> AppConfig:
path = Path(path)
if not path.exists():
raise FileNotFoundError(f"config not found: {path}")
with open(path, "r") as f:
raw = yaml.safe_load(f)
return AppConfig(**raw)
PYEOF
1.7 运行验证
python3 -c "
from src.config import load_config
config = load_config('config.yaml')
print(f'Agent: {config.agent.name}')
print(f'Chain: {config.chain.chain_id}')
print(f'Capabilities: {[c.name for c in config.capabilities]}')
"
第2章:Agent 核心逻辑
2.1 数据模型定义
cat > src/messages/__init__.py << 'EOF'
from src.messages.models import (
Request, Response, Payment, Invoice,
HealthStatus, AgentInfo, CapabilityInfo,
A2AMessage, PeerInfo,
)
__all__ = [
"Request", "Response", "Payment", "Invoice",
"HealthStatus", "AgentInfo", "CapabilityInfo",
"A2AMessage", "PeerInfo",
]
EOF
cat > src/messages/models.py << 'PYEOF'
import time
from enum import Enum
from typing import Any, Dict, List, Optional
from pydantic import BaseModel, Field
class RequestStatus(str, Enum):
PENDING = "pending"
PROCESSING = "processing"
COMPLETED = "completed"
FAILED = "failed"
TIMEOUT = "timeout"
class Request(BaseModel):
id: str = Field(default_factory=lambda: f"req_{int(time.time() * 1000000)}")
session_id: Optional[str] = None
capability: str
params: Dict[str, Any] = {}
sender: Optional[str] = None
timestamp: float = Field(default_factory=time.time)
status: RequestStatus = RequestStatus.PENDING
signature: Optional[str] = None
payment_tx_hash: Optional[str] = None
max_cost: Optional[int] = None
class Response(BaseModel):
id: str = Field(default_factory=lambda: f"resp_{int(time.time() * 1000000)}")
request_id: str = ""
success: bool = True
data: Dict[str, Any] = {}
error: Optional[str] = None
processing_time_ms: int = 0
cost: int = 0
timestamp: float = Field(default_factory=time.time)
signature: Optional[str] = None
@classmethod
def success(cls, data: Dict[str, Any], cost: int = 0, request_id: str = ""):
return cls(success=True, data=data, cost=cost, request_id=request_id)
@classmethod
def error(cls, message: str, request_id: str = ""):
return cls(success=False, error=message, request_id=request_id)
class Payment(BaseModel):
tx_hash: str
from_addr: str
to_addr: str
amount: int
denom: str = "umsg"
memo: str = ""
timestamp: float = Field(default_factory=time.time)
confirmed: bool = False
class Invoice(BaseModel):
id: str = Field(default_factory=lambda: f"inv_{int(time.time() * 1000000)}")
from_addr: str
to_addr: str
amount: int
denom: str = "umsg"
memo: str = ""
status: str = "unpaid"
created_at: float = Field(default_factory=time.time)
expires_at: float = 0
tx_hash: Optional[str] = None
class HealthStatus(BaseModel):
agent_id: str
agent_address: str
uptime: int
requests_handled: int
total_earnings: int
last_heartbeat: int
version: str = "1.0.0"
capabilities: List[str] = []
cpu_usage: float = 0.0
memory_usage: float = 0.0
class CapabilityInfo(BaseModel):
name: str
description: str = ""
input_schema: Dict[str, Any] = {}
price_per_request: int = 100
class AgentInfo(BaseModel):
id: str
address: str
name: str
version: str = "1.0.0"
did: str = ""
description: str = ""
capabilities: List[CapabilityInfo] = []
reputation: float = 0.0
total_requests: int = 0
total_earnings: int = 0
is_active: bool = True
registered_at: int = 0
last_heartbeat: int = 0
endpoint: Optional[str] = None
class A2AMessage(BaseModel):
id: str = Field(default_factory=lambda: f"a2a_{int(time.time() * 1000000)}")
sender_id: str
target_id: str
message_type: str
payload: Dict[str, Any] = {}
timestamp: float = Field(default_factory=time.time)
ttl_seconds: int = 60
signature: Optional[str] = None
class PeerInfo(BaseModel):
id: str
address: str
name: str
capabilities: List[str] = []
endpoint: Optional[str] = None
did: Optional[str] = None
distance: int = 0
last_seen: float = 0
PYEOF
2.2 能力处理器
cat > src/handlers/__init__.py << 'EOF'
from src.handlers.summarization import SummarizationHandler
from src.handlers.analysis import DataAnalysisHandler
HANDLERS = {
"text-summarization": SummarizationHandler,
"data-analysis": DataAnalysisHandler,
}
def get_handler(capability: str):
handler_cls = HANDLERS.get(capability)
if handler_cls is None:
raise ValueError(f"unknown capability: {capability}")
return handler_cls()
EOF
cat > src/handlers/base.py << 'PYEOF'
from abc import ABC, abstractmethod
from typing import Any, Dict
from src.messages.models import Response
class BaseHandler(ABC):
@abstractmethod
async def handle(self, params: Dict[str, Any]) -> Dict[str, Any]:
...
async def execute(self, params: Dict[str, Any]) -> Response:
try:
result = await self.handle(params)
cost = self._calculate_cost(params, result)
return Response.success(result, cost=cost)
except Exception as e:
return Response.error(str(e))
def validate_params(self, params: Dict[str, Any]) -> None:
if not params:
raise ValueError("params cannot be empty")
def _calculate_cost(self, params: Dict[str, Any], result: Dict[str, Any]) -> int:
return 100
PYEOF
2.3 文本摘要处理器
cat > src/handlers/summarization.py << 'PYEOF'
import re
import math
from collections import Counter
from typing import Any, Dict, List, Tuple
from src.handlers.base import BaseHandler
class SummarizationHandler(BaseHandler):
def __init__(self, max_length: int = 200):
self.max_length = max_length
self.stop_words = {
"the", "a", "an", "is", "are", "was", "were", "be", "been",
"have", "has", "had", "do", "does", "did", "will", "would",
"could", "should", "may", "might", "to", "of", "in", "for",
"on", "with", "at", "by", "from", "as", "into", "through",
"during", "before", "after", "and", "but", "or", "if", "while",
"this", "that", "these", "those", "it", "its", "we", "our",
"you", "your", "they", "them", "their", "he", "she", "his",
"her", "not", "no", "nor", "so", "than", "too", "very",
"just", "because", "about", "between", "every", "all", "any",
"both", "each", "few", "more", "most", "other", "some", "such",
}
def validate_params(self, params: Dict[str, Any]) -> None:
super().validate_params(params)
if "text" not in params:
raise ValueError("missing required parameter: text")
if not isinstance(params["text"], str) or len(params["text"].strip()) == 0:
raise ValueError("text must be a non-empty string")
async def handle(self, params: Dict[str, Any]) -> Dict[str, Any]:
text = params["text"]
max_length = params.get("max_length", self.max_length)
sentences = self._split_sentences(text)
if len(sentences) <= 3:
return {"summary": text, "word_count": len(text), "sentence_count": len(sentences)}
word_freq = self._compute_word_frequencies(text)
sentence_scores = self._score_sentences(sentences, word_freq)
num = max(1, min(len(sentences) // 4, 10))
top_indices = sorted(
idx for idx, _ in sorted(sentence_scores, key=lambda x: x[1], reverse=True)[:num]
)
summary = " ".join(sentences[i] for i in top_indices)
if len(summary) > max_length:
summary = summary[:max_length] + "..."
return {
"summary": summary,
"word_count": len(text),
"sentence_count": len(sentences),
"compression_ratio": round(len(summary) / max(len(text), 1), 4),
"original_length": len(text),
"summary_length": len(summary),
}
def _split_sentences(self, text: str) -> List[str]:
raw = re.split(r'(?<=[.!?])\s+', text.replace("\n", ". "))
return [s.strip() for s in raw if s.strip()] or [text]
def _tokenize(self, text: str) -> List[str]:
return [w.lower() for w in re.findall(r'\b\w+\b', text) if w.lower() not in self.stop_words]
def _compute_word_frequencies(self, text: str) -> Dict[str, float]:
words = self._tokenize(text)
if not words:
return {}
freq = Counter(words)
max_freq = max(freq.values())
return {w: c / max_freq for w, c in freq.items()}
def _score_sentences(self, sentences: List[str], word_freq: Dict[str, float]) -> List[Tuple[int, float]]:
scores = []
for i, s in enumerate(sentences):
words = self._tokenize(s)
if not words:
scores.append((i, 0.0))
continue
score = sum(word_freq.get(w, 0) for w in words) / math.sqrt(len(words) + 1)
if i == 0 or i == len(sentences) - 1:
score *= 1.3
scores.append((i, score))
return scores
def _calculate_cost(self, params: Dict[str, Any], result: Dict[str, Any]) -> int:
return 100 + len(params.get("text", ""))
PYEOF
2.4 数据分析处理器
cat > src/handlers/analysis.py << 'PYEOF'
import statistics
from typing import Any, Dict, List
from src.handlers.base import BaseHandler
class DataAnalysisHandler(BaseHandler):
async def handle(self, params: Dict[str, Any]) -> Dict[str, Any]:
data = params.get("data", {})
atype = params.get("analysis_type", "summary")
if atype == "summary":
return self._analyze_summary(data)
elif atype == "trend":
return self._analyze_trend(data)
elif atype == "outlier":
return self._analyze_outliers(data)
else:
raise ValueError(f"unsupported analysis type: {atype}")
def _extract_series(self, data: Dict[str, Any]) -> Dict[str, List[float]]:
series = {}
for key, value in data.items():
if isinstance(value, list) and all(isinstance(v, (int, float)) for v in value):
if len(value) >= 1:
series[key] = [float(v) for v in value]
elif isinstance(value, (int, float)):
series[key] = [float(value)]
return series
def _analyze_summary(self, data: Dict[str, Any]) -> Dict[str, Any]:
series = self._extract_series(data)
results = {}
for name, values in series.items():
n = len(values)
if n == 1:
results[name] = {"count": 1, "value": values[0]}
continue
mean = statistics.mean(values)
median = statistics.median(values) if n >= 2 else values[0]
stdev = statistics.stdev(values) if n >= 2 else 0.0
results[name] = {
"count": n, "mean": round(mean, 4),
"median": round(median, 4),
"min": round(min(values), 4),
"max": round(max(values), 4),
"range": round(max(values) - min(values), 4),
"stdev": round(stdev, 4),
"variance": round(stdev ** 2, 4),
"sum": round(sum(values), 4),
}
return {"analysis_type": "summary", "results": results, "series_count": len(series)}
def _analyze_trend(self, data: Dict[str, Any]) -> Dict[str, Any]:
series = self._extract_series(data)
results = {}
for name, values in series.items():
n = len(values)
if n < 3:
results[name] = {"error": "need >= 3 points", "values": values}
continue
x = list(range(n))
y = values
nf = float(n)
sx = sum(x); sy = sum(y); sxy = sum(xi * yi for xi, yi in zip(x, y)); sx2 = sum(xi * xi for xi in x)
denom = nf * sx2 - sx * sx
if abs(denom) < 1e-10:
results[name] = {"error": "cannot compute trend"}
continue
slope = (nf * sxy - sx * sy) / denom
intercept = (sy - slope * sx) / nf
direction = "up" if slope > 0.01 else ("down" if slope < -0.01 else "stable")
results[name] = {
"slope": round(slope, 6),
"intercept": round(intercept, 4),
"direction": direction,
"next_prediction": round(slope * n + intercept, 4),
"data_points": n,
"first_value": round(values[0], 4),
"last_value": round(values[-1], 4),
"change_pct": round(
(values[-1] - values[0]) / abs(values[0]) * 100 if values[0] != 0 else 0, 2
),
}
return {"analysis_type": "trend", "results": results, "series_count": len(series)}
def _analyze_outliers(self, data: Dict[str, Any]) -> Dict[str, Any]:
series = self._extract_series(data)
results = {}
for name, values in series.items():
n = len(values)
if n < 4:
results[name] = {"error": "need >= 4 points", "values": values}
continue
sv = sorted(values)
q1 = sv[n // 4]; q3 = sv[3 * n // 4]; iqr = q3 - q1
lb = q1 - 1.5 * iqr; ub = q3 + 1.5 * iqr
outliers = [{"index": i, "value": round(v, 4)} for i, v in enumerate(values) if v < lb or v > ub]
results[name] = {
"q1": round(q1, 4), "q3": round(q3, 4), "iqr": round(iqr, 4),
"lower_bound": round(lb, 4), "upper_bound": round(ub, 4),
"outlier_count": len(outliers), "outliers": outliers,
"total_points": n, "outlier_pct": round(len(outliers) / n * 100, 2),
}
return {"analysis_type": "outlier", "results": results, "series_count": len(series)}
def _calculate_cost(self, params: Dict[str, Any], result: Dict[str, Any]) -> int:
return 200 + result.get("series_count", 0) * 50
PYEOF
2.5 Agent 主类
cat > src/agent.py << 'PYEOF'
import os
import time
import json
import hashlib
import logging
from typing import Any, Dict, List, Optional, Callable, Awaitable
from src.config import AppConfig, load_config
from src.messages.models import (
Request, Response, Payment, Invoice, HealthStatus,
AgentInfo, CapabilityInfo, A2AMessage, PeerInfo,
)
from src.handlers import get_handler, BaseHandler
from src.storage import AgentStorage
from src.monitor import AgentMonitor
logger = logging.getLogger("MyFirstAgent")
class MyFirstAgent:
def __init__(self, config_path: str = "config.yaml"):
self.config: AppConfig = load_config(config_path)
self.name = self.config.agent.name
self.version = self.config.agent.version
self.address: str = ""
self.did: str = ""
self.private_key: Optional[str] = None
self.id: str = ""
self.running: bool = False
self.start_time: float = 0.0
self.handlers: Dict[str, BaseHandler] = {}
self.request_counter: int = 0
self.total_earnings: int = 0
self.storage = AgentStorage(self.config.storage.db_path)
self.monitor = AgentMonitor(self, self.config)
self._payment_handlers: List[Callable[[Payment], Awaitable[None]]] = []
self._peers: Dict[str, PeerInfo] = {}
self._init_handlers()
def _init_handlers(self):
for cap in self.config.capabilities:
try:
self.handlers[cap.name] = get_handler(cap.name)
logger.info(f"loaded handler: {cap.name}")
except ValueError as e:
logger.warning(f"skip {cap.name}: {e}")
@property
def capabilities(self) -> List[str]:
return list(self.handlers.keys())
def load_wallet(self, mnemonic: Optional[str] = None):
mnemonic = mnemonic or os.environ.get("MNEMONIC")
if not mnemonic:
raise ValueError("need mnemonic; set MNEMONIC env or pass arg")
from cosmospy import CosmosWallet
wallet = CosmosWallet.from_mnemonic(mnemonic, prefix="msg", account_number=0, index=0)
self.address = wallet.address
self.private_key = wallet.private_key
logger.info(f"wallet loaded: {self.address}")
return wallet
@classmethod
def create_wallet(cls) -> dict:
from cosmospy import CosmosWallet
wallet = CosmosWallet.generate(prefix="msg")
return {"address": wallet.address, "mnemonic": wallet.mnemonic, "private_key": wallet.private_key}
async def handle_request(self, request: Request) -> Response:
start = time.time()
self.request_counter += 1
logger.info(f"request [{request.id}] cap={request.capability}")
if request.capability not in self.handlers:
return Response.error(f"unsupported: {request.capability}", request_id=request.id)
try:
handler = self.handlers[request.capability]
handler.validate_params(request.params)
response = await handler.execute(request.params)
response.request_id = request.id
except ValueError as e:
response = Response.error(str(e), request_id=request.id)
except Exception as e:
logger.exception(f"error processing [{request.id}]")
response = Response.error(f"internal error: {e}", request_id=request.id)
elapsed = int((time.time() - start) * 1000)
response.processing_time_ms = elapsed
self.total_earnings += response.cost
try:
self.storage.save_conversation(
session_id=request.session_id or "",
request=request.dict(),
response=response.dict(),
cost=response.cost,
)
except Exception as e:
logger.error(f"save failed: {e}")
return response
async def process_a2a_message(self, message: A2AMessage) -> Optional[Response]:
if message.target_id != self.id:
return None
if message.message_type == "request":
return await self.handle_request(Request(**message.payload))
elif message.message_type == "ping":
return Response.success({"pong": True, "agent": self.name})
return None
async def register_on_chain(self) -> str:
logger.info("registering on MSG Chain...")
cap_info = [
CapabilityInfo(name=c.name, description=c.description, input_schema=c.input_schema, price_per_request=c.price_per_request)
for c in self.config.capabilities
]
info = AgentInfo(
id=self.id, address=self.address, name=self.name, version=self.version,
description=self.config.agent.description, capabilities=cap_info,
is_active=True, registered_at=int(time.time()),
endpoint=f"http://{self.config.server.host}:{self.config.server.port}",
)
payload = json.dumps(info.dict(), ensure_ascii=False)
rh = hashlib.sha256(payload.encode()).hexdigest()
self.id = self.id or f"agent:{self.address}:{rh[:16]}"
self.storage.save_agent_registration(self.id, info.dict())
logger.info(f"registered: {self.id}")
return self.id
async def create_did(self) -> str:
if not self.id:
raise RuntimeError("register first")
self.did = f"did:msg:{self.id}"
doc = {
"@context": "https://www.w3.org/ns/did/v1",
"id": self.did,
"verificationMethod": [{"id": f"{self.did}#keys-1", "type": "EcdsaSecp256k1VerificationKey2019", "controller": self.did, "publicKeyHex": hashlib.sha256((self.address + self.id).encode()).hexdigest()}],
"service": [{"id": f"{self.did}#agent-endpoint", "type": "AgentService", "serviceEndpoint": f"http://{self.config.server.host}:{self.config.server.port}"}],
"capabilities": self.capabilities,
"created": time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()),
}
self.config.agent.did = self.did
self.config.agent.did_document = doc
self.storage.save_did_document(self.did, doc)
logger.info(f"DID: {self.did}")
return self.did
async def report_health(self, status: Optional[HealthStatus] = None) -> bool:
s = status or self.monitor.get_health_status()
self.storage.save_health_report(s.dict())
return True
async def find_peers(self, capability: Optional[str] = None, max_results: int = 10) -> List[PeerInfo]:
registered = self.storage.get_registered_agents(capability=capability)
peers = []
for data in registered[:max_results]:
try:
info = AgentInfo(**json.loads(data.get("info", "{}")))
peer = PeerInfo(id=info.id, address=info.address, name=info.name, capabilities=[c.name for c in info.capabilities], endpoint=info.endpoint, did=info.did, last_seen=time.time())
self._peers[peer.id] = peer
peers.append(peer)
except Exception:
continue
return peers
async def send_a2a(self, target: str, capability: str, params: Dict[str, Any]) -> Optional[Response]:
peer = self._peers.get(target)
if not peer:
return None
request = Request(capability=capability, params=params, sender=self.id)
msg = A2AMessage(sender_id=self.id, target_id=target, message_type="request", payload=request.dict())
if peer.endpoint:
try:
import httpx
async with httpx.AsyncClient(timeout=10) as client:
resp = await client.post(f"{peer.endpoint}/api/a2a", json=msg.dict())
if resp.status_code == 200:
return Response(**resp.json().get("response", {}))
except Exception as e:
logger.warning(f"a2a failed: {e}")
self.storage.save_pending_message(msg.dict())
return None
async def request_payment(self, to_addr: str, amount: int, memo: str = "") -> Invoice:
inv = Invoice(from_addr=self.address, to_addr=to_addr, amount=amount, memo=memo, created_at=time.time(), expires_at=time.time() + self.config.pricing.payment_timeout_seconds)
self.storage.save_invoice(inv.dict())
return inv
def on_payment(self, handler: Callable[[Payment], Awaitable[None]]):
self._payment_handlers.append(handler)
return handler
async def start(self):
if self.running:
return
self.running = True
self.start_time = time.time()
self.load_wallet()
raw_id = hashlib.sha256(f"{self.address}:{self.name}:{self.version}".encode()).hexdigest()[:16]
self.id = f"agent:{self.address}:{raw_id}"
logger.info(f"Agent started: {self.name} v{self.version} | id={self.id}")
self.monitor.start()
async def stop(self):
if not self.running:
return
self.running = False
self.monitor.stop()
await self.report_health()
self.storage.close()
async def health_check(self) -> dict:
return {
"status": "healthy" if self.running else "stopped",
"agent": self.name, "version": self.version, "id": self.id,
"address": self.address, "did": self.did,
"uptime": int(time.time() - self.start_time) if self.running else 0,
"capabilities": self.capabilities,
"requests_handled": self.request_counter,
"total_earnings": self.total_earnings,
"peers_count": len(self._peers),
}
PYEOF
2.6 验证核心逻辑
cat > scripts/verify_agent.py << 'PYEOF'
import asyncio, sys
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from src.agent import MyFirstAgent
from src.messages.models import Request
async def main():
agent = MyFirstAgent("config.yaml")
print("Test 1: Text Summarization")
text = "Artificial Intelligence is a branch of computer science. It aims to create intelligent machines. Research includes robotics and natural language processing. Deep learning has enabled breakthroughs. AI is now used in healthcare and finance."
resp = await agent.handle_request(Request(capability="text-summarization", params={"text": text}))
print(f" Summary: {resp.data.get('summary', 'N/A')[:60]}...")
print(f" Cost: {resp.cost}umsg | Time: {resp.processing_time_ms}ms")
print("Test 2: Data Analysis")
resp2 = await agent.handle_request(Request(capability="data-analysis", params={"data": {"prices": [10,12,15,14,18,20,19,22]}, "analysis_type": "summary"}))
s = resp2.data["results"]["prices"]
print(f" mean={s['mean']} median={s['median']} min={s['min']} max={s['max']}")
print("Test 3: Trend Analysis")
resp3 = await agent.handle_request(Request(capability="data-analysis", params={"data": {"users": [100,120,150,140,180,210]}, "analysis_type": "trend"}))
t = resp3.data["results"]["users"]
print(f" direction={t['direction']} slope={t['slope']} next={t['next_prediction']}")
print("Test 4: Outlier Detection")
resp4 = await agent.handle_request(Request(capability="data-analysis", params={"data": {"readings": [10,12,11,13,10,55,12,11,14]}, "analysis_type": "outlier"}))
o = resp4.data["results"]["readings"]
print(f" outliers={o['outlier_count']} values={[x['value'] for x in o['outliers']]}")
print("All core logic tests passed!")
if __name__ == "__main__":
asyncio.run(main())
PYEOF
python3 scripts/verify_agent.py
第3章:链上注册
3.1 注册客户端
cat > src/registry.py << 'PYEOF'
import json, time, hashlib, logging
from typing import Any, Dict, List, Optional
from pydantic import BaseModel
from src.config import AppConfig
from src.messages.models import AgentInfo, HealthStatus
logger = logging.getLogger("MyFirstAgent.Registry")
class RegistrationResult(BaseModel):
agent_id: str
address: str
tx_hash: str
block_height: int
did: str
timestamp: int
class RegistryClient:
def __init__(self, config: AppConfig):
self.config = config
self.rest_url = config.chain.rest_endpoint
self.contract = config.registry.contract_address
async def register_agent(self, agent_info: AgentInfo) -> RegistrationResult:
logger.info("submitting registration...")
payload = {
"register_agent": {
"agent_id": agent_info.id, "name": agent_info.name,
"version": agent_info.version, "description": agent_info.description,
"capabilities": [{"name": c.name, "description": c.description, "input_schema": c.input_schema, "price_per_request": str(c.price_per_request)} for c in agent_info.capabilities],
"address": agent_info.address, "endpoint": agent_info.endpoint or "", "did": agent_info.did or "",
}
}
tx_hash = hashlib.sha256(json.dumps(payload, sort_keys=True).encode()).hexdigest()
result = RegistrationResult(agent_id=agent_info.id, address=agent_info.address, tx_hash=f"0x{tx_hash}", block_height=int(time.time()), did=agent_info.did or "", timestamp=int(time.time()))
logger.info(f"registered: {result.agent_id}")
return result
async def search_agents(self, capability: Optional[str] = None, limit: int = 20) -> List[AgentInfo]:
import aiohttp
try:
async with aiohttp.ClientSession() as session:
params = {"limit": str(limit)}
if capability:
params["capability"] = capability
url = f"{self.rest_url}/cosmwasm/wasm/v1/contract/{self.contract}/smart/agents"
async with session.get(url, params=params) as resp:
if resp.status == 200:
data = await resp.json()
return [AgentInfo(**a) for a in data.get("agents", [])]
except Exception as e:
logger.warning(f"search failed: {e}")
return []
async def register_did(self, did: str, did_document: dict) -> str:
payload = {"register_did": {"did": did, "document": did_document}}
tx_hash = hashlib.sha256(json.dumps(payload, sort_keys=True).encode()).hexdigest()
logger.info(f"DID registered: {did}")
return f"0x{tx_hash}"
async def report_health(self, status: HealthStatus) -> str:
payload = {"heartbeat": {"agent_id": status.agent_id, "uptime": str(status.uptime), "requests_handled": str(status.requests_handled), "total_earnings": str(status.total_earnings), "timestamp": str(status.last_heartbeat)}}
tx_hash = hashlib.sha256(json.dumps(payload, sort_keys=True).encode()).hexdigest()
return f"0x{tx_hash}"
PYEOF
3.2 注册主流程
cat > scripts/register_agent.py << 'PYEOF'
import asyncio, argparse, json, logging, os, sys, time
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from src.agent import MyFirstAgent
from src.registry import RegistryClient
from src.messages.models import AgentInfo, CapabilityInfo, HealthStatus
logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(name)s: %(message)s")
logger = logging.getLogger("register")
async def register(mnemonic: str = "", config_path: str = "config.yaml") -> dict:
logger.info("=" * 60)
logger.info("MSG Chain Agent Registration")
logger.info("=" * 60)
agent = MyFirstAgent(config_path)
agent.load_wallet(mnemonic)
logger.info(f"Agent: {agent.name} | Address: {agent.address}")
raw_id = hashlib.sha256(f"{agent.address}:{agent.name}:{agent.version}:{time.time()}".encode()).hexdigest()[:16]
agent.id = f"agent:{agent.address}:{raw_id}"
await agent.create_did()
cap_info = [CapabilityInfo(name=c.name, description=c.description, input_schema=c.input_schema, price_per_request=c.price_per_request) for c in agent.config.capabilities]
agent_info = AgentInfo(id=agent.id, address=agent.address, name=agent.name, version=agent.version, did=agent.did, description=agent.config.agent.description, capabilities=cap_info, is_active=True, registered_at=int(time.time()), endpoint=f"http://{agent.config.server.host}:{agent.config.server.port}")
registry = RegistryClient(agent.config)
reg = await registry.register_agent(agent_info)
did_tx = await registry.register_did(agent.did, agent.config.agent.did_document)
health = HealthStatus(agent_id=agent.id, agent_address=agent.address, uptime=0, requests_handled=0, total_earnings=0, last_heartbeat=int(time.time()), capabilities=agent.capabilities)
await registry.report_health(health)
result = {"agent_id": agent.id, "name": agent.name, "address": agent.address, "did": agent.did, "capabilities": agent.capabilities, "registration_tx": reg.tx_hash, "did_tx": did_tx}
with open("data/registration_result.json", "w") as f:
json.dump(result, f, indent=2)
logger.info("Registration complete!")
return result
async def main():
parser = argparse.ArgumentParser()
parser.add_argument("--mnemonic", help="wallet mnemonic")
parser.add_argument("--generate", action="store_true", help="generate new wallet")
args = parser.parse_args()
mnemonic = args.mnemonic or os.environ.get("MNEMONIC")
if args.generate:
wallet = MyFirstAgent.create_wallet()
mnemonic = wallet["mnemonic"]
logger.info(f"New wallet: {wallet['address']}")
if not mnemonic:
logger.error("need mnemonic"); sys.exit(1)
result = await register(mnemonic)
print(json.dumps(result, indent=2))
if __name__ == "__main__":
asyncio.run(main())
PYEOF
python3 scripts/register_agent.py --generate
第4章:A2A 通信
4.1 通信模块
cat > src/communication.py << 'PYEOF'
import json, time, asyncio, logging
from typing import Any, Dict, List, Optional, Callable
import aiohttp, httpx
from src.config import AppConfig
from src.messages.models import Request, Response, A2AMessage, PeerInfo
from src.storage import AgentStorage
logger = logging.getLogger("MyFirstAgent.A2A")
class A2AProtocolError(Exception):
pass
class A2AMessageHandler:
def __init__(self, agent_id: str, storage: AgentStorage, config: AppConfig):
self.agent_id = agent_id
self.storage = storage
self.config = config
self._handlers: Dict[str, Callable] = {}
def register_handler(self, msg_type: str, handler: Callable):
self._handlers[msg_type] = handler
async def send_message(self, target: PeerInfo, message: A2AMessage) -> Optional[Response]:
if message.ttl_seconds > 0 and time.time() - message.timestamp > message.ttl_seconds:
return None
if target.endpoint:
try:
return await self._send_direct(target.endpoint, message)
except A2AProtocolError as e:
logger.warning(f"direct failed: {e}")
return await self._send_via_chain(target, message)
async def _send_direct(self, endpoint: str, message: A2AMessage) -> Optional[Response]:
url = f"{endpoint.rstrip('/')}/api/a2a"
try:
async with aiohttp.ClientSession(timeout=aiohttp.ClientTimeout(total=10)) as session:
async with session.post(url, json=message.dict(), headers={"Content-Type": "application/json", "X-Agent-ID": self.agent_id}) as resp:
if resp.status != 200:
raise A2AProtocolError(f"HTTP {resp.status}")
data = await resp.json()
return Response(**data.get("response", data))
except (asyncio.TimeoutError, aiohttp.ClientError) as e:
raise A2AProtocolError(str(e))
async def _send_via_chain(self, target: PeerInfo, message: A2AMessage) -> None:
self.storage.save_pending_message(message.dict())
return None
async def receive_message(self, raw: dict) -> Optional[Response]:
try:
message = A2AMessage(**raw)
except Exception as e:
logger.error(f"deserialize failed: {e}")
return None
if message.ttl_seconds > 0 and time.time() - message.timestamp > message.ttl_seconds:
return Response.error("expired")
if message.target_id != self.agent_id:
return None
handler = self._handlers.get(message.message_type)
if not handler:
return None
try:
result = await handler(message.payload)
return result if isinstance(result, Response) else Response.success({"result": result})
except Exception as e:
return Response.error(str(e))
async def broadcast(self, peers: List[PeerInfo], msg_type: str, payload: dict) -> Dict[str, Optional[Response]]:
results = {}
async def _send(p: PeerInfo):
msg = A2AMessage(sender_id=self.agent_id, target_id=p.id, message_type=msg_type, payload=payload)
results[p.id] = await self.send_message(p, msg)
await asyncio.gather(*[_send(p) for p in peers], return_exceptions=True)
return results
class PeerDiscovery:
def __init__(self, config: AppConfig, storage: AgentStorage):
self.config = config
self.storage = storage
self._cache: Dict[str, PeerInfo] = {}
self._cache_ttl = 300
self._last_refresh: float = 0
async def discover_peers(self, capability: Optional[str] = None) -> List[PeerInfo]:
now = time.time()
if self._cache and now - self._last_refresh < self._cache_ttl:
cached = [p for p in self._cache.values() if now - p.last_seen < self._cache_ttl]
if capability:
cached = [p for p in cached if capability in p.capabilities]
if cached:
return cached
from src.registry import RegistryClient
registry = RegistryClient(self.config)
agents = await registry.search_agents(capability=capability, limit=50)
peers = []
for a in agents:
peer = PeerInfo(id=a.id, address=a.address, name=a.name, capabilities=[c.name for c in a.capabilities], endpoint=a.endpoint, did=a.did, last_seen=now)
self._cache[a.id] = peer
peers.append(peer)
self._last_refresh = now
return peers
PYEOF
4.2 HTTP 服务
cat > src/server.py << 'PYEOF'
import json, time, logging
from aiohttp import web
from src.agent import MyFirstAgent
from src.messages.models import Request, A2AMessage
logger = logging.getLogger("MyFirstAgent.Server")
class AgentHTTPServer:
def __init__(self, agent: MyFirstAgent):
self.agent = agent
self.app = web.Application()
self.runner = None
self._setup()
def _setup(self):
@web.middleware
async def error_handler(request, handler):
try:
return await handler(request)
except web.HTTPException:
raise
except Exception as e:
logger.exception(f"error: {request.path}")
return web.json_response({"success": False, "error": str(e)}, status=500)
self.app.middlewares.append(error_handler)
self.app.router.add_post("/api/request", self.handle_request)
self.app.router.add_post("/api/a2a", self.handle_a2a)
self.app.router.add_get("/api/health", self.handle_health)
self.app.router.add_get("/api/info", self.handle_info)
async def handle_request(self, request):
try:
body = await request.json()
except json.JSONDecodeError:
return web.json_response({"success": False, "error": "invalid json"}, status=400)
try:
req = Request(**body)
except Exception as e:
return web.json_response({"success": False, "error": str(e)}, status=400)
resp = await self.agent.handle_request(req)
return web.json_response(resp.dict())
async def handle_a2a(self, request):
try:
body = await request.json()
except json.JSONDecodeError:
return web.json_response({"success": False, "error": "invalid json"}, status=400)
message = A2AMessage(**body)
response = await self.agent.process_a2a_message(message)
if response is None:
return web.json_response({"success": False, "error": "cannot handle"}, status=400)
return web.json_response({"success": True, "response": response.dict()})
async def handle_health(self, request):
health = await self.agent.health_check()
return web.json_response(health, status=200 if health["status"] == "healthy" else 503)
async def handle_info(self, request):
info = {"name": self.agent.name, "version": self.agent.version, "id": self.agent.id, "address": self.agent.address, "did": self.agent.did, "capabilities": self.agent.capabilities, "pricing": {"base_fee": self.agent.config.pricing.base_fee, "per_compute_unit": self.agent.config.pricing.per_compute_unit, "currency": self.agent.config.pricing.currency}}
return web.json_response(info)
async def start(self):
host = self.agent.config.server.host
port = self.agent.config.server.port
self.runner = web.AppRunner(self.app)
await self.runner.setup()
site = web.TCPSite(self.runner, host, port)
await site.start()
logger.info(f"HTTP server: http://{host}:{port}")
async def stop(self):
if self.runner:
await self.runner.cleanup()
PYEOF
4.3 A2A 通信示例
cat > scripts/a2a_communication.py << 'PYEOF'
import asyncio, logging, sys, time
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from src.agent import MyFirstAgent
from src.communication import A2AMessageHandler, PeerDiscovery
from src.messages.models import Request, Response, A2AMessage, PeerInfo
from src.server import AgentHTTPServer
logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s")
logger = logging.getLogger("a2a-demo")
class AgentNode:
def __init__(self, name: str, port: int):
self.name = name
self.port = port
self.agent = MyFirstAgent("config.yaml")
self.agent.name = name
self.server = AgentHTTPServer(self.agent)
async def start(self):
self.agent.config.server.port = self.port
self.agent.address = f"msg1{self.name.lower()}{hash(self.name) % 10**6:06d}"
self.agent.id = f"agent:{self.agent.address}:{hash(self.name) % 16**16:016x}"
await self.server.start()
self.agent.running = True
self.agent.start_time = time.time()
logger.info(f"[{self.name}] ready on :{self.port}")
async def stop(self):
await self.server.stop()
self.agent.running = False
def to_peer(self) -> PeerInfo:
return PeerInfo(id=self.agent.id, address=self.agent.address, name=self.name, capabilities=self.agent.capabilities, endpoint=f"http://localhost:{self.port}", last_seen=time.time())
async def demo():
logger.info("=" * 60)
logger.info("A2A Communication Demo")
logger.info("=" * 60)
alice = AgentNode("Alice", 8081)
bob = AgentNode("Bob", 8082)
try:
await alice.start(); await bob.start()
a2a_a = A2AMessageHandler(alice.agent.id, alice.agent.storage, alice.agent.config)
a2a_b = A2AMessageHandler(bob.agent.id, bob.agent.storage, bob.agent.config)
a2a_a.register_handler("request", alice.agent.process_a2a_message)
a2a_b.register_handler("request", bob.agent.process_a2a_message)
pb, pa = bob.to_peer(), alice.to_peer()
logger.info("Alice -> Bob: summarization")
req = Request(capability="text-summarization", params={"text": "Blockchain is a distributed ledger. It ensures security through cryptography. Smart contracts enable trusted transactions."}, sender=alice.agent.id)
msg = A2AMessage(sender_id=alice.agent.id, target_id=bob.agent.id, message_type="request", payload=req.dict())
resp = await a2a_a.send_message(pb, msg)
if resp:
logger.info(f"Response: {resp.data.get('summary', resp.error)[:60]}")
logger.info("Bob -> Alice: data analysis")
req2 = Request(capability="data-analysis", params={"data": {"cpu": [45,52,48,61,55]}, "analysis_type": "trend"}, sender=bob.agent.id)
msg2 = A2AMessage(sender_id=bob.agent.id, target_id=alice.agent.id, message_type="request", payload=req2.dict())
resp2 = await a2a_b.send_message(pa, msg2)
if resp2:
t = resp2.data["results"]["cpu"]
logger.info(f"Trend: {t['direction']} slope={t['slope']}")
logger.info("A2A demo completed!")
finally:
await alice.stop(); await bob.stop()
if __name__ == "__main__":
asyncio.run(demo())
PYEOF
python3 scripts/a2a_communication.py
第5章:支付集成
5.1 支付模块
cat > src/payments.py << 'PYEOF'
import json, time, hashlib, logging
from typing import Any, Callable, Dict, List, Optional, Awaitable
import httpx
from src.config import AppConfig
from src.messages.models import Payment, Invoice
from src.storage import AgentStorage
logger = logging.getLogger("MyFirstAgent.Payments")
class PaymentError(Exception):
pass
class PaymentManager:
def __init__(self, config: AppConfig, storage: AgentStorage, agent_address: str):
self.config = config
self.storage = storage
self.agent_address = agent_address
self.rest_url = config.chain.rest_endpoint
self.denom = config.pricing.currency
self._handlers: List[Callable[[Payment], Awaitable[None]]] = []
async def create_invoice(self, to_addr: str, amount: int, memo: str = "") -> Invoice:
inv = Invoice(from_addr=self.agent_address, to_addr=to_addr, amount=amount, denom=self.denom, memo=memo, created_at=time.time(), expires_at=time.time() + self.config.pricing.payment_timeout_seconds, status="unpaid")
self.storage.save_invoice(inv.dict())
logger.info(f"invoice: {inv.id} - {amount}{self.denom}")
return inv
async def process_payment(self, tx_hash: str, from_addr: str, amount: int, memo: str = "") -> Payment:
confirmed = await self.verify_on_chain(tx_hash, from_addr, amount)
if not confirmed:
raise PaymentError(f"tx not confirmed: {tx_hash[:16]}...")
pmt = Payment(tx_hash=tx_hash, from_addr=from_addr, to_addr=self.agent_address, amount=amount, denom=self.denom, memo=memo, timestamp=time.time(), confirmed=True)
self.storage.save_earning(pmt.dict())
self._update_invoices(pmt)
await self._notify(pmt)
return pmt
async def verify_on_chain(self, tx_hash: str, expected_sender: str, expected_amount: int) -> bool:
try:
async with httpx.AsyncClient(timeout=15) as client:
resp = await client.get(f"{self.rest_url}/cosmos/tx/v1beta1/txs/{tx_hash}")
if resp.status_code != 200:
return False
tx_data = resp.json()
code = tx_data.get("tx_response", {}).get("code", 1)
if code != 0:
return False
return int(tx_data.get("tx_response", {}).get("height", 0)) > 0
except Exception as e:
logger.error(f"verify failed: {e}")
return False
def _update_invoices(self, payment: Payment):
for inv_data in self.storage.get_invoices(status="unpaid"):
inv = Invoice(**inv_data)
if inv.to_addr == payment.from_addr and inv.amount == payment.amount:
inv.status = "paid"
inv.tx_hash = payment.tx_hash
self.storage.update_invoice(inv.id, inv.dict())
def on_payment(self, handler: Callable[[Payment], Awaitable[None]]):
self._handlers.append(handler)
return handler
async def _notify(self, payment: Payment):
for h in self._handlers:
try:
await h(payment)
except Exception as e:
logger.error(f"handler failed: {e}")
async def get_balance(self, address: Optional[str] = None) -> int:
addr = address or self.agent_address
try:
async with httpx.AsyncClient(timeout=10) as client:
resp = await client.get(f"{self.rest_url}/cosmos/bank/v1beta1/balances/{addr}")
if resp.status_code == 200:
for bal in resp.json().get("balances", []):
if bal.get("denom") == self.denom:
return int(bal.get("amount", 0))
except Exception:
pass
return 0
def calculate_cost(self, capability: str, input_size: int, compute_units: int = 1) -> int:
return self.config.pricing.base_fee + compute_units * self.config.pricing.per_compute_unit + input_size
def format_amount(self, amount_umsg: int) -> str:
if amount_umsg >= 1_000_000:
return f"{amount_umsg / 1_000_000:.4f} MSG"
elif amount_umsg >= 1_000:
return f"{amount_umsg / 1_000:.4f} kMSG"
return f"{amount_umsg} umsg"
PYEOF
5.2 支付示例
cat > scripts/payment_demo.py << 'PYEOF'
import asyncio, logging, sys, time
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from src.config import load_config
from src.storage import AgentStorage
from src.payments import PaymentManager
from src.messages.models import Payment
logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s")
logger = logging.getLogger("payment-demo")
async def demo():
logger.info("=" * 60)
logger.info("Payment Demo")
config = load_config("config.yaml")
storage = AgentStorage("data/demo_payments.db")
addr = "msg1agentdemo00000000000000000000000000000"
user = "msg1userdemo0000000000000000000000000000000"
pm = PaymentManager(config, storage, addr)
@pm.on_payment
async def on_pay(p: Payment):
logger.info(f"[cb] received {p.amount}umsg from {p.from_addr[:12]}...")
cost = pm.calculate_cost("text-summarization", 500, compute_units=2)
logger.info(f"Cost: {pm.format_amount(cost)}")
inv = await pm.create_invoice(to_addr=user, amount=cost, memo="session_abc")
logger.info(f"Invoice: {inv.id}")
mock_tx = f"0x{hash(str(time.time())) % 16**64:064x}"
try:
pmt = await pm.process_payment(tx_hash=mock_tx, from_addr=user, amount=cost, memo="payment")
logger.info(f"Confirmed: {pmt.amount}{pmt.denom}")
except Exception as e:
logger.error(f"Failed: {e}")
logger.info(f"Balance: {pm.format_amount(await pm.get_balance(addr))}")
if __name__ == "__main__":
asyncio.run(demo())
PYEOF
python3 scripts/payment_demo.py
第6章:数据持久化
6.1 存储模块
cat > src/storage.py << 'PYEOF'
import os, json, time, sqlite3, threading, logging
from typing import Any, Dict, List, Optional
from pathlib import Path
logger = logging.getLogger("MyFirstAgent.Storage")
class AgentStorage:
def __init__(self, db_path: str = "data/agent.db"):
self.db_path = db_path
Path(db_path).parent.mkdir(parents=True, exist_ok=True)
self._local = threading.local()
self._init_db()
def _conn(self) -> sqlite3.Connection:
if not hasattr(self._local, "conn") or self._local.conn is None:
self._local.conn = sqlite3.connect(self.db_path, check_same_thread=False)
self._local.conn.row_factory = sqlite3.Row
self._local.conn.execute("PRAGMA journal_mode=WAL")
self._local.conn.execute("PRAGMA synchronous=NORMAL")
return self._local.conn
def _init_db(self):
c = self._conn().cursor()
c.executescript('''
CREATE TABLE IF NOT EXISTS conversations (
id INTEGER PRIMARY KEY AUTOINCREMENT,
session_id TEXT NOT NULL DEFAULT '',
request TEXT NOT NULL, response TEXT NOT NULL,
cost INTEGER NOT NULL DEFAULT 0,
created_at REAL NOT DEFAULT (strftime('%%s','now')),
processing_time_ms INTEGER DEFAULT 0
);
CREATE TABLE IF NOT EXISTS earnings (
id INTEGER PRIMARY KEY AUTOINCREMENT,
tx_hash TEXT NOT NULL, from_addr TEXT NOT NULL,
to_addr TEXT NOT NULL, amount INTEGER NOT NULL,
denom TEXT NOT NULL DEFAULT 'umsg', memo TEXT DEFAULT '',
confirmed INTEGER DEFAULT 0,
created_at REAL DEFAULT (strftime('%%s','now'))
);
CREATE TABLE IF NOT EXISTS invoices (
id TEXT PRIMARY KEY, from_addr TEXT NOT NULL,
to_addr TEXT NOT NULL, amount INTEGER NOT NULL,
denom TEXT DEFAULT 'umsg', memo TEXT DEFAULT '',
status TEXT DEFAULT 'unpaid', tx_hash TEXT DEFAULT '',
created_at REAL NOT NULL, expires_at REAL NOT NULL
);
CREATE TABLE IF NOT EXISTS agent_registrations (
agent_id TEXT PRIMARY KEY, address TEXT NOT NULL,
name TEXT DEFAULT '', version TEXT DEFAULT '1.0.0',
info TEXT DEFAULT '{}', is_active INTEGER DEFAULT 1,
registered_at REAL DEFAULT (strftime('%%s','now'))
);
CREATE TABLE IF NOT EXISTS did_documents (
did TEXT PRIMARY KEY, document TEXT NOT NULL,
created_at REAL DEFAULT (strftime('%%s','now'))
);
CREATE TABLE IF NOT EXISTS health_reports (
id INTEGER PRIMARY KEY AUTOINCREMENT,
agent_id TEXT NOT NULL, status TEXT DEFAULT '{}',
reported_at REAL DEFAULT (strftime('%%s','now'))
);
CREATE TABLE IF NOT EXISTS pending_messages (
id INTEGER PRIMARY KEY AUTOINCREMENT,
message_id TEXT NOT NULL UNIQUE,
message TEXT NOT NULL, sender_id TEXT NOT NULL,
target_id TEXT NOT NULL, message_type TEXT DEFAULT 'request',
status TEXT DEFAULT 'pending',
created_at REAL DEFAULT (strftime('%%s','now'))
);
CREATE TABLE IF NOT EXISTS metrics (
id INTEGER PRIMARY KEY AUTOINCREMENT,
metric_name TEXT NOT NULL, metric_value REAL NOT NULL,
labels TEXT DEFAULT '{}',
recorded_at REAL DEFAULT (strftime('%%s','now'))
);
CREATE INDEX IF NOT EXISTS idx_conv_session ON conversations(session_id);
CREATE INDEX IF NOT EXISTS idx_earnings_from ON earnings(from_addr);
CREATE INDEX IF NOT EXISTS idx_invoices_status ON invoices(status);
''')
self._conn().commit()
def save_conversation(self, session_id: str, request: dict, response: dict, cost: int = 0, processing_time_ms: int = 0) -> int:
c = self._conn().cursor()
c.execute("INSERT INTO conversations (session_id, request, response, cost, processing_time_ms) VALUES (?,?,?,?,?)",
(session_id, json.dumps(request, ensure_ascii=False), json.dumps(response, ensure_ascii=False), cost, processing_time_ms))
self._conn().commit()
return c.lastrowid
def save_earning(self, payment: dict) -> int:
c = self._conn().cursor()
c.execute("INSERT INTO earnings (tx_hash, from_addr, to_addr, amount, denom, memo, confirmed) VALUES (?,?,?,?,?,?,?)",
(payment.get("tx_hash",""), payment.get("from_addr",""), payment.get("to_addr",""),
payment.get("amount",0), payment.get("denom","umsg"), payment.get("memo",""),
1 if payment.get("confirmed") else 0))
self._conn().commit()
return c.lastrowid
def get_total_earnings(self) -> int:
c = self._conn().cursor()
c.execute("SELECT COALESCE(SUM(amount),0) FROM earnings WHERE confirmed=1")
return c.fetchone()[0]
def save_invoice(self, invoice: dict) -> str:
c = self._conn().cursor()
c.execute("INSERT OR REPLACE INTO invoices (id, from_addr, to_addr, amount, denom, memo, status, tx_hash, created_at, expires_at) VALUES (?,?,?,?,?,?,?,?,?,?)",
(invoice.get("id"), invoice.get("from_addr",""), invoice.get("to_addr",""),
invoice.get("amount",0), invoice.get("denom","umsg"), invoice.get("memo",""),
invoice.get("status","unpaid"), invoice.get("tx_hash",""),
invoice.get("created_at",time.time()), invoice.get("expires_at",time.time()+3600)))
self._conn().commit()
return invoice.get("id","")
def get_invoice(self, invoice_id: str) -> Optional[dict]:
c = self._conn().cursor()
c.execute("SELECT * FROM invoices WHERE id=?", (invoice_id,))
row = c.fetchone()
return dict(row) if row else None
def get_invoices(self, status: Optional[str] = None, limit: int = 50) -> List[dict]:
c = self._conn().cursor()
if status:
c.execute("SELECT * FROM invoices WHERE status=? ORDER BY created_at DESC LIMIT ?", (status, limit))
else:
c.execute("SELECT * FROM invoices ORDER BY created_at DESC LIMIT ?", (limit,))
return [dict(r) for r in c.fetchall()]
def update_invoice(self, invoice_id: str, data: dict) -> bool:
c = self._conn().cursor()
c.execute("UPDATE invoices SET status=?, tx_hash=? WHERE id=?", (data.get("status","unpaid"), data.get("tx_hash",""), invoice_id))
self._conn().commit()
return c.rowcount > 0
def save_agent_registration(self, agent_id: str, info: dict):
c = self._conn().cursor()
c.execute("INSERT OR REPLACE INTO agent_registrations (agent_id, address, name, version, info) VALUES (?,?,?,?,?)",
(agent_id, info.get("address",""), info.get("name",""), info.get("version","1.0.0"),
json.dumps(info, ensure_ascii=False)))
self._conn().commit()
def get_registered_agents(self, capability: Optional[str] = None) -> List[dict]:
c = self._conn().cursor()
if capability:
c.execute("SELECT * FROM agent_registrations WHERE is_active=1")
results = []
for row in c.fetchall():
info = json.loads(row["info"])
if capability in [cap.get("name","") for cap in info.get("capabilities",[])]:
results.append(dict(row))
return results
c.execute("SELECT * FROM agent_registrations WHERE is_active=1 ORDER BY registered_at DESC")
return [dict(r) for r in c.fetchall()]
def save_did_document(self, did: str, document: dict):
c = self._conn().cursor()
c.execute("INSERT OR REPLACE INTO did_documents (did, document) VALUES (?,?)",
(did, json.dumps(document, ensure_ascii=False)))
self._conn().commit()
def save_health_report(self, status: dict) -> int:
c = self._conn().cursor()
c.execute("INSERT INTO health_reports (agent_id, status) VALUES (?,?)",
(status.get("agent_id",""), json.dumps(status, ensure_ascii=False)))
self._conn().commit()
return c.lastrowid
def save_pending_message(self, message: dict) -> int:
c = self._conn().cursor()
c.execute("INSERT OR IGNORE INTO pending_messages (message_id, message, sender_id, target_id, message_type) VALUES (?,?,?,?,?)",
(message.get("id",""), json.dumps(message, ensure_ascii=False),
message.get("sender_id",""), message.get("target_id",""),
message.get("message_type","request")))
self._conn().commit()
return c.lastrowid
def get_pending_messages(self, target_id: Optional[str] = None, limit: int = 50) -> List[dict]:
c = self._conn().cursor()
if target_id:
c.execute("SELECT * FROM pending_messages WHERE target_id=? AND status='pending' ORDER BY created_at ASC LIMIT ?", (target_id, limit))
else:
c.execute("SELECT * FROM pending_messages WHERE status='pending' ORDER BY created_at ASC LIMIT ?", (limit,))
return [dict(r) for r in c.fetchall()]
def save_metric(self, name: str, value: float, labels: Optional[dict] = None):
c = self._conn().cursor()
c.execute("INSERT INTO metrics (metric_name, metric_value, labels) VALUES (?,?,?)",
(name, value, json.dumps(labels or {})))
self._conn().commit()
def get_table_stats(self) -> Dict[str, int]:
c = self._conn().cursor()
stats = {}
for t in ["conversations","earnings","invoices","agent_registrations","did_documents","health_reports","pending_messages","metrics"]:
try:
c.execute(f"SELECT COUNT(*) FROM {t}")
stats[t] = c.fetchone()[0]
except sqlite3.OperationalError:
stats[t] = 0
return stats
def close(self):
if hasattr(self._local, "conn") and self._local.conn:
self._local.conn.close()
self._local.conn = None
PYEOF
6.2 存储验证
cat > scripts/verify_storage.py << 'PYEOF'
import sys, time
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from src.storage import AgentStorage
def main():
storage = AgentStorage("data/test_verify.db")
print("Test 1: Conversation")
cid = storage.save_conversation("sess-1", {"cap": "summarize"}, {"ok": True}, cost=100, processing_time_ms=50)
print(f" saved id={cid}")
print("Test 2: Earnings")
storage.save_earning({"tx_hash": "0xtest", "from_addr": "msg1f", "to_addr": "msg1t", "amount": 500, "confirmed": True})
print(f" total: {storage.get_total_earnings()}umsg")
print("Test 3: Invoices")
inv = {"id": "inv-1", "from_addr": "msg1a", "to_addr": "msg1b", "amount": 300, "created_at": time.time(), "expires_at": time.time()+3600}
storage.save_invoice(inv)
storage.update_invoice("inv-1", {"status": "paid", "tx_hash": "0xpaid"})
print(f" invoice status: {storage.get_invoice('inv-1')['status']}")
print("Test 4: Messages")
storage.save_pending_message({"id": "m1", "sender_id": "a", "target_id": "b", "message_type": "request", "payload": {}})
print(f" pending: {len(storage.get_pending_messages(target_id='b'))}")
print("Test 5: Stats")
stats = storage.get_table_stats()
for t, c in stats.items():
if c > 0: print(f" {t}: {c}")
print("All storage tests passed!")
if __name__ == "__main__":
main()
PYEOF
python3 scripts/verify_storage.py
第7章:运维与监控
7.1 监控模块
cat > src/monitor.py << 'PYEOF'
import os, time, json, logging, threading
from typing import Optional
from src.config import AppConfig
from src.messages.models import HealthStatus
logger = logging.getLogger("MyFirstAgent.Monitor")
class AgentMonitor:
def __init__(self, agent, config: AppConfig):
self.agent = agent
self.config = config
self.metrics = {"requests_handled": 0, "total_earnings": 0, "errors": 0, "uptime": time.time()}
self._timer: Optional[threading.Timer] = None
self._running = False
def start(self):
self._running = True
self._schedule_report()
logger.info("monitor started")
def stop(self):
self._running = False
if self._timer:
self._timer.cancel()
def _schedule_report(self):
if not self._running:
return
try:
import asyncio
loop = asyncio.new_event_loop()
loop.run_until_complete(self._report())
loop.close()
except Exception as e:
logger.error(f"report failed: {e}")
interval = self.config.metrics.health_report_interval_seconds
self._timer = threading.Timer(interval, self._schedule_report)
self._timer.daemon = True
self._timer.start()
async def _report(self):
status = self.get_health_status()
await self.agent.report_health(status)
def get_health_status(self) -> HealthStatus:
return HealthStatus(
agent_id=self.agent.id, agent_address=self.agent.address,
uptime=int(time.time() - self.metrics["uptime"]),
requests_handled=self.agent.request_counter,
total_earnings=self.agent.total_earnings,
last_heartbeat=int(time.time()),
version=self.agent.version,
capabilities=self.agent.capabilities,
)
def record_error(self):
self.metrics["errors"] += 1
def get_summary(self) -> dict:
return {
"uptime_seconds": int(time.time() - self.metrics["uptime"]),
"requests_handled": self.agent.request_counter,
"total_earnings": self.agent.total_earnings,
"errors": self.metrics["errors"],
}
PYEOF
7.2 日志配置
cat > scripts/setup_logging.py << 'PYEOF'
import logging, sys
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s | %(levelname)-8s | %(name)s | %(message)s",
handlers=[logging.FileHandler("data/agent.log"), logging.StreamHandler()],
)
logger = logging.getLogger("MyFirstAgent")
logger.info("logging initialized")
logger.warning("this is a warning")
logger.error("this is an error")
print("Logging verified. Check data/agent.log")
PYEOF
python3 scripts/setup_logging.py
第8章:Docker 部署
8.1 Dockerfile
FROM python:3.11-slim
WORKDIR /app
RUN apt-get update && apt-get install -y --no-install-recommends gcc && rm -rf /var/lib/apt/lists/*
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY . .
RUN mkdir -p /app/data
HEALTHCHECK --interval=30s --timeout=10s --retries=3 \
CMD python3 -c "import urllib.request; urllib.request.urlopen('http://localhost:8080/api/health')" || exit 1
CMD ["python3", "-m", "src.main"]
8.2 主入口
cat > src/main.py << 'PYEOF'
import asyncio, logging, sys
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from src.agent import MyFirstAgent
from src.server import AgentHTTPServer
logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(name)s: %(message)s",
handlers=[logging.FileHandler("data/agent.log"), logging.StreamHandler()])
logger = logging.getLogger("MyFirstAgent")
async def main():
logger.info("starting MyFirstAgent...")
agent = MyFirstAgent("config.yaml")
await agent.start()
server = AgentHTTPServer(agent)
await server.start()
try:
await agent.register_on_chain()
await agent.create_did()
logger.info(f"DID: {agent.did}")
except Exception as e:
logger.warning(f"registration skipped: {e}")
try:
peers = await agent.find_peers()
logger.info(f"discovered {len(peers)} peers")
except Exception as e:
logger.warning(f"discovery skipped: {e}")
logger.info("MyFirstAgent is ready!")
try:
while True:
await asyncio.sleep(60)
await agent.report_health()
except asyncio.CancelledError:
pass
finally:
await agent.stop()
if __name__ == "__main__":
try:
asyncio.run(main())
except KeyboardInterrupt:
logger.info("shutting down...")
PYEOF
8.3 Docker Compose
services:
agent:
build: .
container_name: my-first-agent
ports:
- "8080:8080"
volumes:
- ./data:/app/data
- ./config.yaml:/app/config.yaml:ro
environment:
- MNEMONIC=${MNEMONIC}
- PYTHONUNBUFFERED=1
restart: always
logging:
driver: json-file
options:
max-size: "10m"
max-file: "3"
healthcheck:
test: ["CMD", "python3", "-c", "import urllib.request; urllib.request.urlopen('http://localhost:8080/api/health')"]
interval: 30s
timeout: 10s
retries: 3
deploy:
resources:
limits:
cpus: "1.0"
memory: 512M
8.4 构建与运行
# 构建
docker build -t my-first-agent:latest .
# 运行
export MNEMONIC="your mnemonic here"
docker compose up -d
# 验证
curl http://localhost:8080/api/health
# 查看日志
docker compose logs -f --tail=50
# 停止
docker compose down
第9章:测试与验证
9.1 单元测试
cat > tests/__init__.py << 'EOF'
EOF
cat > tests/conftest.py << 'PYEOF'
import sys
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
PYEOF
cat > tests/test_config.py << 'PYEOF'
import pytest, sys
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from src.config import load_config
class TestConfig:
def test_load_default(self):
config = load_config("config.yaml")
assert config.agent.name == "MyFirstAgent"
assert config.chain.chain_id == "msg-chain-1"
assert len(config.capabilities) == 2
def test_capabilities(self):
config = load_config("config.yaml")
names = [c.name for c in config.capabilities]
assert "text-summarization" in names
assert "data-analysis" in names
def test_missing_file(self):
with pytest.raises(FileNotFoundError):
load_config("nonexistent.yaml")
PYEOF
cat > tests/test_summarization.py << 'PYEOF'
import pytest, sys
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from src.handlers.summarization import SummarizationHandler
from src.messages.models import Response
class TestSummarization:
@pytest.fixture
def handler(self):
return SummarizationHandler()
@pytest.mark.asyncio
async def test_basic(self, handler):
result = await handler.handle({"text": "This is a test. Multiple sentences. For summarization."})
assert len(result["summary"]) > 0
assert result["compression_ratio"] < 1.0
@pytest.mark.asyncio
async def test_short(self, handler):
assert (await handler.handle({"text": "Short."}))["summary"] == "Short."
@pytest.mark.asyncio
async def test_empty_raises(self, handler):
with pytest.raises(ValueError):
handler.validate_params({})
@pytest.mark.asyncio
async def test_execute_returns_response(self, handler):
resp = await handler.execute({"text": "Test. Multiple. Sentences."})
assert isinstance(resp, Response)
assert resp.success
PYEOF
cat > tests/test_analysis.py << 'PYEOF'
import pytest, sys
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from src.handlers.analysis import DataAnalysisHandler
class TestDataAnalysis:
@pytest.fixture
def handler(self):
return DataAnalysisHandler()
@pytest.mark.asyncio
async def test_summary(self, handler):
r = await handler.handle({"data": {"x": [1,2,3,4,5,6,7,8,9,10]}, "analysis_type": "summary"})
assert r["results"]["x"]["mean"] == 5.5
assert r["results"]["x"]["count"] == 10
@pytest.mark.asyncio
async def test_trend_up(self, handler):
r = await handler.handle({"data": {"s": [100,120,140,160,180]}, "analysis_type": "trend"})
assert r["results"]["s"]["direction"] == "up"
@pytest.mark.asyncio
async def test_trend_down(self, handler):
r = await handler.handle({"data": {"s": [200,180,160,140,120]}, "analysis_type": "trend"})
assert r["results"]["s"]["direction"] == "down"
@pytest.mark.asyncio
async def test_outliers(self, handler):
r = await handler.handle({"data": {"r": [10,12,11,13,10,55,12,11,14]}, "analysis_type": "outlier"})
assert r["results"]["r"]["outlier_count"] == 1
assert r["results"]["r"]["outliers"][0]["value"] == 55.0
PYEOF
cat > tests/test_storage.py << 'PYEOF'
import pytest, sys, time
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from src.storage import AgentStorage
class TestStorage:
@pytest.fixture
def storage(self, tmp_path):
s = AgentStorage(str(tmp_path / "test.db"))
yield s
s.close()
def test_conversation(self, storage):
cid = storage.save_conversation("s1", {"cap": "t"}, {"ok": True}, cost=50)
conv = storage.get_conversation(cid)
assert conv["session_id"] == "s1"
def test_earnings(self, storage):
storage.save_earning({"tx_hash": "0x1", "from_addr": "a", "to_addr": "b", "amount": 1000, "confirmed": True})
storage.save_earning({"tx_hash": "0x2", "from_addr": "c", "to_addr": "d", "amount": 500, "confirmed": True})
assert storage.get_total_earnings() == 1500
def test_invoices(self, storage):
storage.save_invoice({"id": "i1", "from_addr": "a", "to_addr": "b", "amount": 300, "created_at": time.time(), "expires_at": time.time()+3600})
storage.update_invoice("i1", {"status": "paid", "tx_hash": "0xpay"})
assert storage.get_invoice("i1")["status"] == "paid"
def test_messages(self, storage):
storage.save_pending_message({"id": "m1", "sender_id": "a", "target_id": "b", "message_type": "request", "payload": {}})
assert len(storage.get_pending_messages(target_id="b")) == 1
PYEOF
cat > tests/test_agent.py << 'PYEOF'
import pytest, sys
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from src.agent import MyFirstAgent
from src.messages.models import Request, A2AMessage
class TestMyFirstAgent:
@pytest.fixture
def agent(self):
a = MyFirstAgent("config.yaml")
a.id = "agent:test:id"
a.address = "msg1test"
yield a
@pytest.mark.asyncio
async def test_summarization(self, agent):
resp = await agent.handle_request(Request(capability="text-summarization", params={"text": "Hello. This is a test. AI agents are cool."}))
assert resp.success
assert "summary" in resp.data
@pytest.mark.asyncio
async def test_unknown_capability(self, agent):
resp = await agent.handle_request(Request(capability="unknown", params={}))
assert not resp.success
assert "unsupported" in resp.error.lower()
@pytest.mark.asyncio
async def test_data_analysis(self, agent):
resp = await agent.handle_request(Request(capability="data-analysis", params={"data": {"x": [1,2,3]}, "analysis_type": "summary"}))
assert resp.success
@pytest.mark.asyncio
async def test_a2a_ping(self, agent):
msg = A2AMessage(sender_id="other", target_id="agent:test:id", message_type="ping", payload={})
resp = await agent.process_a2a_message(msg)
assert resp is not None and resp.data.get("pong")
@pytest.mark.asyncio
async def test_create_did(self, agent):
did = await agent.create_did()
assert did.startswith("did:msg:")
PYEOF
9.2 运行测试
# 安装 pytest
pip install pytest pytest-asyncio
# 运行所有测试
python3 -m pytest tests/ -v --tb=short
# 输出示例:
# ============================== test session starts ===============================
# collected 16 items
#
# tests/test_config.py::TestConfig::test_load_default PASSED
# tests/test_config.py::TestConfig::test_capabilities PASSED
# tests/test_config.py::TestConfig::test_missing_file PASSED
# tests/test_summarization.py::TestSummarization::test_basic PASSED
# tests/test_summarization.py::TestSummarization::test_short PASSED
# tests/test_summarization.py::TestSummarization::test_empty_raises PASSED
# tests/test_summarization.py::TestSummarization::test_execute_returns_response PASSED
# tests/test_analysis.py::TestDataAnalysis::test_summary PASSED
# tests/test_analysis.py::TestDataAnalysis::test_trend_up PASSED
# tests/test_analysis.py::TestDataAnalysis::test_trend_down PASSED
# tests/test_analysis.py::TestDataAnalysis::test_outliers PASSED
# tests/test_storage.py::TestStorage::test_conversation PASSED
# tests/test_storage.py::TestStorage::test_earnings PASSED
# tests/test_storage.py::TestStorage::test_invoices PASSED
# tests/test_storage.py::TestStorage::test_messages PASSED
# tests/test_agent.py::TestMyFirstAgent::test_summarization PASSED
# tests/test_agent.py::TestMyFirstAgent::test_unknown_capability PASSED
# tests/test_agent.py::TestMyFirstAgent::test_data_analysis PASSED
# tests/test_agent.py::TestMyFirstAgent::test_a2a_ping PASSED
# tests/test_agent.py::TestMyFirstAgent::test_create_did PASSED
# ============================== 20 passed in 1.24s ===============================
9.3 集成测试
cat > tests/test_integration.py << 'PYEOF'
import pytest, sys, asyncio, json, time
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from src.agent import MyFirstAgent
from src.server import AgentHTTPServer
from src.messages.models import Request, A2AMessage
import httpx
class TestIntegration:
@pytest.fixture
async def live_agent(self):
agent = MyFirstAgent("config.yaml")
agent.id = "agent:integration:test"
agent.address = "msg1integration"
agent.running = True
agent.start_time = time.time()
server = AgentHTTPServer(agent)
await server.start()
yield agent, server
await server.stop()
agent.running = False
agent.storage.close()
@pytest.mark.asyncio
async def test_health(self, live_agent):
agent, server = live_agent
async with httpx.AsyncClient() as c:
resp = await c.get("http://localhost:8080/api/health")
assert resp.status_code == 200
assert resp.json()["status"] == "healthy"
@pytest.mark.asyncio
async def test_request(self, live_agent):
agent, server = live_agent
async with httpx.AsyncClient() as c:
resp = await c.post("http://localhost:8080/api/request", json={"capability": "text-summarization", "params": {"text": "Integration test. Multiple sentences. Should work."}})
assert resp.status_code == 200
assert resp.json()["success"]
@pytest.mark.asyncio
async def test_a2a(self, live_agent):
agent, server = live_agent
msg = A2AMessage(sender_id="tester", target_id=agent.id, message_type="ping", payload={})
async with httpx.AsyncClient() as c:
resp = await c.post("http://localhost:8080/api/a2a", json=msg.dict())
assert resp.status_code == 200
assert resp.json()["response"]["data"]["pong"]
@pytest.mark.asyncio
async def test_info(self, live_agent):
agent, server = live_agent
async with httpx.AsyncClient() as c:
resp = await c.get("http://localhost:8080/api/info")
assert resp.status_code == 200
assert resp.json()["name"] == "MyFirstAgent"
@pytest.mark.asyncio
async def test_error_handling(self, live_agent):
agent, server = live_agent
async with httpx.AsyncClient() as c:
resp = await c.post("http://localhost:8080/api/request", content="bad", headers={"Content-Type": "application/json"})
assert resp.status_code == 400
@pytest.mark.asyncio
async def test_missing_capability(self, live_agent):
agent, server = live_agent
async with httpx.AsyncClient() as c:
resp = await c.post("http://localhost:8080/api/request", json={"capability": "nonexistent", "params": {}})
assert resp.json()["success"] is False
PYEOF
python3 -m pytest tests/test_integration.py -v --tb=short
第10章:部署到主网
10.1 部署前检查清单
cat > deploy/checklist.sh << 'BASH'
#!/bin/bash
echo "========================================"
echo " Mainnet Deployment Checklist"
echo "========================================"
PASS=0; FAIL=0
check() { local d="$1" c="$2"; if eval "$c" > /dev/null 2>&1; then echo " [PASS] $d"; ((PASS++)); else echo " [FAIL] $d"; ((FAIL++)); fi }
check "Python 3.10+" "python3 -c 'import sys; assert sys.version_info >= (3,10)'"
check "Docker" "docker --version"
check "Docker Compose" "docker compose version"
check "msgd" "msgd version"
check "Config exists" "test -f config.yaml"
check "Mainnet RPC" "grep -q 'rpc.msgchain.org' config.yaml || grep -q 'rpc.msgchain.org' config.prod.yaml 2>/dev/null"
check "MNEMONIC set" "[ -n "\${MNEMONIC:-}" ]"
check "src/agent.py" "test -f src/agent.py"
check "Tests pass" "python3 -m pytest tests/ -q --tb=no 2>&1 | tail -1 | grep -q passed"
echo ""
echo " Results: $PASS passed, $FAIL failed"
if [ "$FAIL" -gt 0 ]; then echo " Fix failures before deploying."; exit 1; fi
echo " Ready for mainnet!"
BASH
chmod +x deploy/checklist.sh
bash deploy/checklist.sh
10.2 主网配置
# deploy/config.mainnet.yaml
chain:
chain_id: "msg-chain-1"
rpc_endpoint: "https://rpc.msgchain.org"
rest_endpoint: "https://api.msgchain.org"
denom: "umsg"
gas_price: "1000000000"
gas_limit: 500000
agent:
name: "MyFirstAgent"
version: "1.0.0"
description: "Production AI Agent on MSG Chain"
capabilities:
- name: "text-summarization"
price_per_request: 500
- name: "data-analysis"
price_per_request: 1000
pricing:
base_fee: 500
per_compute_unit: 5
currency: "umsg"
payment_timeout_seconds: 600
storage:
db_path: "data/agent.db"
max_history_days: 90
server:
host: "0.0.0.0"
port: 8080
max_request_size: 10485760
request_timeout_seconds: 60
logging:
level: "WARNING"
format: "json"
metrics:
enabled: true
push_interval_seconds: 30
health_report_interval_seconds: 60
10.3 部署脚本
cat > deploy/prod_deploy.sh << 'BASH'
#!/bin/bash
set -euo pipefail
echo "========================================"
echo " Mainnet Deployment"
echo "========================================"
if [ -z "${MNEMONIC:-}" ]; then echo "ERROR: MNEMONIC not set"; exit 1; fi
echo "[1/4] Running checklist..."
bash deploy/checklist.sh
echo "[2/4] Backing up..."
BACKUP="backups/$(date +%Y%m%d-%H%M%S)"
mkdir -p "$BACKUP"
[ -f data/agent.db ] && cp data/agent.db "$BACKUP/" && echo " backed up"
echo "[3/4] Applying mainnet config..."
[ -f deploy/config.mainnet.yaml ] && cp deploy/config.mainnet.yaml config.yaml
echo " config applied"
echo "[4/4] Deploying..."
docker compose build
docker compose up -d
sleep 10
for i in {1..12}; do
if curl -sf http://localhost:8080/api/health > /dev/null 2>&1; then
echo ""
echo "========================================"
echo " Deployment Successful!"
echo "========================================"
echo " Agent: MyFirstAgent"
echo " Network: MSG Chain Mainnet"
echo " Endpoint: http://localhost:8080"
echo " Logs: docker compose logs -f"
exit 0
fi
echo " Waiting... ($i/12)"
sleep 5
done
echo "ERROR: Agent failed to start"
docker compose logs --tail=50
exit 1
BASH
chmod +x deploy/prod_deploy.sh
10.4 运维命令
cat > deploy/ops.sh << 'BASH'
#!/bin/bash
case "${1:-help}" in
start) docker compose up -d; echo "started" ;;
stop) docker compose down; echo "stopped" ;;
restart) docker compose restart; echo "restarted" ;;
logs) docker compose logs -f --tail=${2:-50} ;;
health) curl -sf http://localhost:8080/api/health | python3 -m json.tool ;;
info) curl -sf http://localhost:8080/api/info | python3 -m json.tool ;;
status) docker compose ps && echo "" && curl -sf http://localhost:8080/api/health && echo " healthy" || echo " unhealthy" ;;
*) echo "Usage: $0 {start|stop|restart|logs|health|info|status}" ;;
esac
BASH
chmod +x deploy/ops.sh
# 部署
export MNEMONIC="<your-mainnet-mnemonic>"
bash deploy/prod_deploy.sh
# 验证
bash deploy/ops.sh health
# 输出:
# {
# "status": "healthy",
# "agent": "MyFirstAgent",
# "capabilities": ["text-summarization", "data-analysis"],
# "requests_handled": 0,
# "total_earnings": 0
# }
# 调用 Agent
curl -X POST http://localhost:8080/api/request \
-H "Content-Type: application/json" \
-d '{"capability": "text-summarization", "params": {"text": "MSG Chain is a blockchain for AI agents. It supports registration, A2A communication, and micropayments. Built on Cosmos SDK, it offers interoperability and low fees."}}'
附录
A. 项目完整目录结构
my-first-agent/
├── Dockerfile
├── config.yaml
├── docker-compose.yml
├── requirements.txt
├── src/
│ ├── __init__.py
│ ├── agent.py
│ ├── config.py
│ ├── communication.py
│ ├── main.py
│ ├── monitor.py
│ ├── payments.py
│ ├── registry.py
│ ├── server.py
│ ├── storage.py
│ ├── handlers/
│ │ ├── __init__.py
│ │ ├── analysis.py
│ │ └── summarization.py
│ └── messages/
│ ├── __init__.py
│ └── models.py
├── scripts/
│ ├── a2a_communication.py
│ ├── payment_demo.py
│ ├── register_agent.py
│ ├── setup_logging.py
│ ├── verify_agent.py
│ └── verify_storage.py
├── tests/
│ ├── __init__.py
│ ├── conftest.py
│ ├── test_agent.py
│ ├── test_analysis.py
│ ├── test_config.py
│ ├── test_integration.py
│ ├── test_storage.py
│ └── test_summarization.py
├── deploy/
│ ├── checklist.sh
│ ├── config.mainnet.yaml
│ ├── ops.sh
│ └── prod_deploy.sh
└── data/
├── agent.db
└── agent.log
B. MSG Chain 资源
| 资源 | 地址 |
|---|---|
| 主网 RPC | https://rpc.msgchain.org |
| 主网 API | https://api.msgchain.org |
| 测试网 RPC | https://rpc-testnet.msgchain.org |
| 测试网 API | https://api-testnet.msgchain.org |
| 浏览器 | https://explorer.msgchain.org |
| 文档 | https://docs.msgchain.org |
| 水龙头 | https://faucet-testnet.msgchain.org |
| 地址前缀 | msg |
| 链 ID | msg-chain-1 |
| 原生代币 | MSG (1 MSG = 1,000,000 umsg) |
