dApp Docs/AI Agent 端到端实战教程:从零构建一个完整Agent
Development reference. Not independently verified for production.

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)