dApp Docs/AI Agent 链上数据索引与事件订阅指南
Development reference. Not independently verified for production.

MSG Chain AI Agent 链上数据索引与事件订阅指南

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

目录

  1. 概述
  2. 链上直接查询
  3. WebSocket 事件订阅
  4. 合约事件解析
  5. 索引器 SDK 集成
  6. 自定义索引部署
  7. 数据聚合与分析
  8. 完整 Agent 数据流示例

1. 概述

1.1 为什么 AI Agent 需要链上数据索引

MSG Chain 上的 AI Agent 并非孤立运行——它们需要感知链上状态变化才能做出智能决策。无论是监听转账、捕捉合约事件、还是查询质押状态,Agent 都需要高效的数据获取管道。

链上数据是 Agent 的「感官输入」。没有实时、准确的数据索引,Agent 就像蒙着眼睛执行任务的操作员。本指南覆盖从直接查询到开发参考级索引器的完整技术栈,帮助开发者构建数据驱动型 Agent。

1.2 MSG Chain 数据基础设施

MSG Chain 基于 Cosmos SDK + CosmWasm,提供多层数据访问能力:

层级 延迟 适用场景 资源消耗
CosmWasm 智能合约查询 ~100ms 即时状态读取 低
Tendermint RPC / LCD ~200ms 账户/交易历史 中
WebSocket 事件推送 ~50ms 实时触发 中
索引器 REST API ~20ms 复杂聚合查询 低
自定义索引器 ~100ms 专属数据管道 高

1.3 事件驱动 Agent 架构

Agent 的典型数据流:

链上交易/事件
    │
    ▼
WebSocket / RPC 轮询
    │
    ▼
事件解析器 ──→ 过滤器 ──→ 路由 ──→ 处理器
                                        │
                                        ▼
                              on-chain action / 推理 / 存储

Agent 不需要持续轮询——通过 WebSocket 订阅 + 索引器回调,可以实现响应式触发。数据管道中的每一层都可以独立扩展。

1.4 MSG Chain 网络参数

1.5 依赖安装

pip install cosmwasm-client aiohttp websockets pydantic asyncpg sqlalchemy

2. 链上直接查询

2.1 CosmWasm 客户端封装

适用于 Agent 需要即时读取合约状态的场景。无需等待索引器同步,直接查询链上最新状态。

import asyncio
from typing import Any, Dict, List, Optional
from datetime import datetime
from decimal import Decimal

from cosmwasm_client import CosmWasmClient


class MSGChainClient:
    """MSG Chain CosmWasm 客户端封装"""

    CHAIN_ID = "msg-chain-1"
    BECH32_PREFIX = "msg"
    DENOM = "umsg"

    def __init__(
        self,
        rpc_endpoint: str = "https://rpc.msgchain.org",
        lcd_endpoint: str = "https://lcd.msgchain.org",
    ):
        self.rpc_endpoint = rpc_endpoint
        self.lcd_endpoint = lcd_endpoint
        self.client: Optional[CosmWasmClient] = None
        self._session: Optional[aiohttp.ClientSession] = None

    async def connect(self):
        self.client = await CosmWasmClient.connect(self.rpc_endpoint)
        self._session = aiohttp.ClientSession()

    async def close(self):
        if self.client:
            await self.client.close()
        if self._session:
            await self._session.close()

2.2 链上查询器


class OnChainQuerier:
    """链上直接查询器 —— 通过 CosmWasm 客户端查询链上状态"""

    def __init__(self, client: CosmWasmClient):
        self.client = client

    async def query_contract_smart(
        self, contract_addr: str, query: Dict[str, Any]
    ) -> Dict[str, Any]:
        return await self.client.query_contract_smart(contract_addr, query)

    async def query_bank_balance(
        self, address: str, denom: str = "umsg"
    ):
        return await self.client.get_balance(address, denom)

    async def query_all_balances(self, address: str):
        return await self.client.get_all_balances(address)

    async def query_staking_validators(self):
        return await self.client.get_validators()

    async def query_delegation(self, delegator: str, validator: str):
        return await self.client.get_delegation(delegator, validator)

    async def query_delegations(self, delegator: str):
        return await self.client.get_delegations(delegator)

    async def query_rewards(self, delegator: str):
        return await self.client.get_total_rewards(delegator)

    async def query_governance_proposals(self, status: Optional[str] = None):
        proposals = await self.client.get_proposals()
        if status:
            return [p for p in proposals if p["status"] == status]
        return proposals

    async def query_proposal(self, proposal_id: int):
        return await self.client.get_proposal(proposal_id)

    async def query_staking_pool(self):
        return await self.client.get_staking_pool()

    async def query_inflation(self):
        return await self.client.get_inflation()

    async def query_supply(self, denom: str = "umsg"):
        return await self.client.get_supply(denom)

2.3 交易历史查询


class TxQuerier:
    """交易历史查询器"""

    def __init__(self, client: CosmWasmClient, lcd_endpoint: str):
        self.client = client
        self.lcd_endpoint = lcd_endpoint

    async def search_txs(
        self, events: List[Dict[str, str]], limit: int = 100, page: int = 1,
    ):
        return await self.client.search_tx(events, limit=limit, page=page)

    async def search_by_sender(self, sender: str, limit: int = 100):
        return await self.search_txs([{"transfer.sender": sender}], limit=limit)

    async def search_by_recipient(self, recipient: str, limit: int = 100):
        return await self.search_txs([{"transfer.recipient": recipient}], limit=limit)

    async def search_by_contract(self, contract_addr: str, limit: int = 100):
        return await self.search_txs(
            [{"wasm._contract_address": contract_addr}], limit=limit
        )

    async def search_agent_actions(self, agent_id: str, action: Optional[str] = None, limit: int = 100):
        events = [{"wasm.agent": agent_id}]
        if action:
            events.append({"wasm.action": action})
        return await self.search_txs(events, limit=limit)

    async def get_tx_by_hash(self, tx_hash: str):
        return await self.client.get_tx(tx_hash)

    async def get_block(self, height: int):
        return await self.client.get_block(height)

    async def get_latest_block(self):
        return await self.client.get_latest_block()

2.4 缓存与重试


class CachedOnChainQuerier:
    """带缓存的查询器 —— 减少重复 RPC 调用"""

    def __init__(self, querier: OnChainQuerier, ttl_seconds: int = 30):
        self.querier = querier
        self._cache: Dict[str, tuple] = {}
        self.ttl = ttl_seconds

    def _key(self, method: str, *args) -> str:
        return f"{method}:{':'.join(str(a) for a in args)}"

    def _fresh(self, key: str) -> bool:
        if key not in self._cache:
            return False
        _, ts = self._cache[key]
        return (datetime.now().timestamp() - ts) < self.ttl

    async def query_contract_smart(self, contract: str, query: Dict):
        key = self._key("query", contract, str(query))
        if self._fresh(key):
            return self._cache[key][0]
        result = await self.querier.query_contract_smart(contract, query)
        self._cache[key] = (result, datetime.now().timestamp())
        return result

    async def query_bank_balance(self, address: str, denom: str = "umsg"):
        key = self._key("balance", address, denom)
        if self._fresh(key):
            return self._cache[key][0]
        result = await self.querier.query_bank_balance(address, denom)
        self._cache[key] = (result, datetime.now().timestamp())
        return result

    def invalidate(self, method: str, *args):
        self._cache.pop(self._key(method, *args), None)

    def clear(self):
        self._cache.clear()


async def safe_query(
    querier: OnChainQuerier,
    contract: str,
    query: Dict[str, Any],
    retries: int = 3,
    timeout: float = 10.0,
):
    """带重试和超时的安全查询"""
    for attempt in range(retries):
        try:
            return await asyncio.wait_for(
                querier.query_contract_smart(contract, query),
                timeout=timeout,
            )
        except asyncio.TimeoutError:
            if attempt == retries - 1:
                raise
            await asyncio.sleep(2 ** attempt)
    return None

2.5 使用示例


async def onchain_example():
    client = MSGChainClient()
    await client.connect()
    querier = OnChainQuerier(client.client)

    balance = await querier.query_bank_balance("msg1...")
    print(f"余额: {balance.amount} umsg")

    state = await querier.query_contract_smart(
        "msg1contract...", {"config": {}}
    )
    print(f"合约配置: {state}")

    proposals = await querier.query_governance_proposals("PROPOSAL_STATUS_VOTING")
    print(f"投票中提案数: {len(proposals)}")

    await client.close()

3. WebSocket 事件订阅

3.1 Tendermint WebSocket 基础

MSG Chain 的 RPC 节点暴露 WebSocket 端点,支持 Tendermint 的事件订阅协议。Agent 通过 JSON-RPC 订阅链上事件,实现实时响应。

WebSocket 端点:wss://rpc.msgchain.org/websocket

订阅格式:

{
    "jsonrpc": "2.0",
    "method": "subscribe",
    "params": ["tm.event='Tx' AND wasm.action='agent_action'"],
    "id": 1
}

3.2 事件订阅器

import json
import logging
from typing import Any, Callable, Dict, List, Optional, Set
from dataclasses import dataclass, field
from datetime import datetime
from enum import Enum
import websockets

logger = logging.getLogger(__name__)


class EventType(Enum):
    TX = "Tx"
    NEW_BLOCK = "NewBlock"
    NEW_BLOCK_HEADER = "NewBlockHeader"
    VALIDATOR_SET_UPDATES = "ValidatorSetUpdates"


@dataclass
class TendermintEvent:
    type: str
    height: int
    tx_hash: Optional[str]
    attributes: Dict[str, str]
    raw: Dict[str, Any]
    timestamp: datetime = field(default_factory=datetime.now)


class Subscription:
    def __init__(self, query: str, handler: Optional[Callable] = None):
        self.query = query
        self.handler = handler
        self.id: int = 0
        self.active: bool = False


class EventSubscriber:
    """实时事件订阅器 —— Tendermint WebSocket 监听链上事件"""

    MAX_RECONNECT_DELAY = 60
    INITIAL_RECONNECT_DELAY = 1
    PING_INTERVAL = 30

    def __init__(
        self,
        agent_id: str,
        ws_endpoint: str = "wss://rpc.msgchain.org/websocket",
    ):
        self.agent_id = agent_id
        self.ws_endpoint = ws_endpoint
        self.ws: Optional[websockets.WebSocketClientProtocol] = None
        self.subscriptions: Dict[str, Subscription] = {}
        self.handlers: Dict[str, List[Callable]] = {}
        self._running = False
        self._next_id = 1
        self._reconnect_delay = self.INITIAL_RECONNECT_DELAY
        self._event_queue: asyncio.Queue = asyncio.Queue()
        self._tasks: Set[asyncio.Task] = set()

    def _next_sid(self) -> int:
        sid = self._next_id
        self._next_id += 1
        return sid

    async def connect(self):
        self.ws = await websockets.connect(
            self.ws_endpoint,
            ping_interval=self.PING_INTERVAL,
            max_size=10 * 1024 * 1024,
        )
        self._reconnect_delay = self.INITIAL_RECONNECT_DELAY

    async def disconnect(self):
        if self.ws:
            await self.ws.close()
            self.ws = None

    async def subscribe(self, query: str, handler: Optional[Callable] = None) -> Subscription:
        """订阅事件

        Args:
            query: Tendermint 查询表达式,如 "tm.event='Tx' AND wasm.action='transfer'"
            handler: 事件处理回调
        """
        sub = Subscription(query, handler)
        sub.id = self._next_sid()
        if self.ws:
            await self.ws.send(json.dumps({
                "jsonrpc": "2.0",
                "method": "subscribe",
                "params": [query],
                "id": sub.id,
            }))
        self.subscriptions[query] = sub
        if handler:
            self.handlers.setdefault(query, []).append(handler)
        return sub

    async def subscribe_to_events(self, topics: List[str], handler: Optional[Callable] = None):
        """批量订阅多个事件主题"""
        for topic in topics:
            await self.subscribe(f"tm.event='Tx' AND wasm.{topic}", handler)

    async def subscribe_to_agent_actions(self, agent_id: Optional[str] = None, handler: Optional[Callable] = None):
        aid = agent_id or self.agent_id
        await self.subscribe(f"tm.event='Tx' AND wasm.agent='{aid}'", handler)

    async def subscribe_to_new_blocks(self, handler: Optional[Callable] = None):
        await self.subscribe("tm.event='NewBlock'", handler)

    async def subscribe_to_transfers(
        self, from_addr: Optional[str] = None, to_addr: Optional[str] = None, handler: Optional[Callable] = None,
    ):
        conditions = ["tm.event='Tx'", "transfer.amount EXISTS"]
        if from_addr:
            conditions.append(f"transfer.sender='{from_addr}'")
        if to_addr:
            conditions.append(f"transfer.recipient='{to_addr}'")
        await self.subscribe(" AND ".join(conditions), handler)

    async def subscribe_to_contract_events(
        self, contract_addr: str, action: Optional[str] = None, handler: Optional[Callable] = None,
    ):
        conditions = ["tm.event='Tx'", f"wasm._contract_address='{contract_addr}'"]
        if action:
            conditions.append(f"wasm.action='{action}'")
        await self.subscribe(" AND ".join(conditions), handler)

    def on(self, event_type: str):
        """装饰器方式注册事件处理器"""
        def decorator(handler):
            self.handlers.setdefault(event_type, []).append(handler)
            return handler
        return decorator

    async def unsubscribe(self, query: str):
        if query in self.subscriptions:
            sub = self.subscriptions[query]
            if self.ws and sub.active:
                await self.ws.send(json.dumps({
                    "jsonrpc": "2.0",
                    "method": "unsubscribe",
                    "params": [query],
                    "id": sub.id,
                }))
            del self.subscriptions[query]
            self.handlers.pop(query, None)

    async def _process_message(self, raw_msg: str):
        try:
            msg = json.loads(raw_msg)
        except json.JSONDecodeError:
            return
        if "error" in msg:
            logger.error(f"Subscription error: {msg['error']}")
            return
        result = msg.get("result", {}).get("data", {})
        if not result:
            return
        value = result.get("value", {})
        event = TendermintEvent(
            type=result.get("type", ""),
            height=int(value.get("height", 0)),
            tx_hash=value.get("TxHash") or value.get("tx_hash"),
            attributes=self._extract_attributes(value),
            raw=value,
        )
        event_key = self._resolve_key(event, raw_msg)
        if event_key in self.handlers:
            for handler in self.handlers[event_key]:
                try:
                    if asyncio.iscoroutinefunction(handler):
                        await handler(event)
                    else:
                        handler(event)
                except Exception as e:
                    logger.error(f"Handler error: {e}")
        await self._event_queue.put(event)

    def _extract_attributes(self, value: Dict) -> Dict[str, str]:
        attrs = {}
        for event in value.get("events", []):
            for attr in event.get("attributes", []):
                key = attr.get("key", "")
                val = attr.get("value", "")
                if isinstance(key, bytes):
                    key = key.decode()
                if isinstance(val, bytes):
                    val = val.decode()
                attrs[f"{event.get('type')}.{key}"] = val
        return attrs

    def _resolve_key(self, event: TendermintEvent, raw: str) -> str:
        for key in self.handlers:
            if key in raw:
                return key
        wasm_events = [k for k in event.attributes if k.startswith("wasm.")]
        return wasm_events[0] if wasm_events else event.type

    async def _reconnect(self):
        logger.info(f"Reconnecting in {self._reconnect_delay}s...")
        await asyncio.sleep(self._reconnect_delay)
        self._reconnect_delay = min(self._reconnect_delay * 2, self.MAX_RECONNECT_DELAY)
        try:
            await self.connect()
            for sub in self.subscriptions.values():
                await self.ws.send(json.dumps({
                    "jsonrpc": "2.0", "method": "subscribe",
                    "params": [sub.query], "id": sub.id,
                }))
                sub.active = True
            self._reconnect_delay = self.INITIAL_RECONNECT_DELAY
        except Exception as e:
            logger.error(f"Reconnect failed: {e}")
            self._tasks.add(asyncio.create_task(self._reconnect()))

    async def _listen(self):
        while self._running:
            try:
                raw = await asyncio.wait_for(self.ws.recv(), timeout=self.PING_INTERVAL + 10)
                self._reconnect_delay = self.INITIAL_RECONNECT_DELAY
                await self._process_message(raw)
            except (asyncio.TimeoutError, websockets.exceptions.ConnectionClosed):
                if self._running:
                    await self._reconnect()
                    break
            except Exception as e:
                logger.error(f"Listen error: {e}")
                if self._running:
                    await self._reconnect()
                    break

    async def event_loop(self):
        self._running = True
        await self.connect()
        for sub in self.subscriptions.values():
            await self.ws.send(json.dumps({
                "jsonrpc": "2.0", "method": "subscribe",
                "params": [sub.query], "id": sub.id,
            }))
            sub.active = True
        self._tasks.add(asyncio.create_task(self._listen()))

    async def get_event(self, timeout: Optional[float] = None) -> Optional[TendermintEvent]:
        try:
            return await asyncio.wait_for(self._event_queue.get(), timeout=timeout)
        except asyncio.TimeoutError:
            return None

    async def stop(self):
        self._running = False
        await self.disconnect()
        for task in self._tasks:
            task.cancel()
        self._tasks.clear()

    @property
    def is_connected(self) -> bool:
        return self.ws is not None and self.ws.open

3.3 事件过滤器与路由器


class EventFilter:
    """事件过滤链"""

    def __init__(self):
        self._filters: List[Callable] = []

    def by_contract(self, addr: str):
        return self._add(lambda e: e.attributes.get("wasm._contract_address", "") == addr)

    def by_action(self, action: str):
        return self._add(lambda e: e.attributes.get("wasm.action", "") == action)

    def by_agent(self, agent_id: str):
        return self._add(lambda e: e.attributes.get("wasm.agent", "") == agent_id)

    def by_amount_gt(self, min_amount: int):
        return self._add(lambda e: self._parse_amount(e) > min_amount)

    def _add(self, fn):
        self._filters.append(fn)
        return self

    def _parse_amount(self, e: TendermintEvent) -> int:
        try:
            return int(e.attributes.get("transfer.amount", "0").replace("umsg", "").strip())
        except (ValueError, AttributeError):
            return 0

    def apply(self, event: TendermintEvent) -> bool:
        return all(f(event) for f in self._filters)


class PriorityEventRouter:
    """优先级事件路由器"""

    HIGH, MEDIUM, LOW = 0, 1, 2

    def __init__(self):
        self._routes: Dict[str, Dict[int, List[Callable]]] = {}

    def route(self, event_type: str, priority: int = MEDIUM):
        def decorator(handler):
            self._routes.setdefault(event_type, {}).setdefault(priority, []).append(handler)
            return handler
        return decorator

    async def dispatch(self, event: TendermintEvent):
        if event.type not in self._routes:
            return
        for priority in sorted(self._routes[event.type]):
            for handler in self._routes[event.type][priority]:
                try:
                    if asyncio.iscoroutinefunction(handler):
                        await handler(event)
                    else:
                        handler(event)
                except Exception as e:
                    logger.error(f"Router error (pri={priority}): {e}")

3.4 使用示例


async def websocket_example():
    subscriber = EventSubscriber("agent1")

    @subscriber.on("wasm.agent_action")
    async def on_action(event: TendermintEvent):
        action = event.attributes.get("wasm.action", "unknown")
        agent = event.attributes.get("wasm.agent", "unknown")
        print(f"Agent {agent} -> {action}")

    async def on_transfer(event: TendermintEvent):
        print(f"转账: {event.attributes.get('transfer.sender', '')[:12]} -> "
              f"{event.attributes.get('transfer.recipient', '')[:12]}: "
              f"{event.attributes.get('transfer.amount', '')}")

    await subscriber.subscribe_to_agent_actions()
    await subscriber.subscribe_to_transfers(handler=on_transfer)
    await subscriber.event_loop()

    try:
        while True:
            event = await subscriber.get_event(timeout=5.0)
            await asyncio.sleep(0.1)
    except KeyboardInterrupt:
        await subscriber.stop()

4. 合约事件解析

4.1 CosmWasm 事件结构

CosmWasm 合约为每个执行返回事件(events),结构如下:

{
    "type": "wasm",
    "attributes": [
        {"key": "_contract_address", "value": "msg1..."},
        {"key": "action", "value": "transfer"},
        {"key": "from", "value": "msg1..."},
        {"key": "to", "value": "msg1..."},
        {"key": "amount", "value": "1000000umsg"}
    ]
}

Agent 合约可自定义事件:

let event = Event::new("agent_action")
    .add_attribute("agent", env.contract.address.to_string())
    .add_attribute("action", action)
    .add_attribute("params", params);
deps.api.emit_event(event);

4.2 事件数据模型

from dataclasses import dataclass, field
from datetime import datetime
from typing import Any, Dict, List, Optional, Type, TypeVar
from abc import ABC, abstractmethod
import json


T = TypeVar("T", bound="BaseEvent")


@dataclass
class BaseEvent(ABC):
    contract: str
    height: int
    timestamp: datetime
    raw: Dict[str, Any]

    @classmethod
    @abstractmethod
    def from_tendermint(cls: Type[T], event: TendermintEvent) -> Optional[T]:
        pass


@dataclass
class TransferEvent(BaseEvent):
    from_addr: str
    to_addr: str
    amount: str
    denom: str

    @classmethod
    def from_tendermint(cls, event: TendermintEvent) -> Optional["TransferEvent"]:
        if "transfer.amount" not in event.attributes:
            return None
        amt = event.attributes.get("transfer.amount", "0")
        denom = "umsg"
        if amt.endswith("umsg"):
            amt = amt.replace("umsg", "")
        return cls(
            contract=event.attributes.get("transfer.recipient", ""),
            height=event.height, timestamp=event.timestamp, raw=event.raw,
            from_addr=event.attributes.get("transfer.sender", ""),
            to_addr=event.attributes.get("transfer.recipient", ""),
            amount=amt, denom=denom,
        )


@dataclass
class AgentActionEvent(BaseEvent):
    agent_id: str
    action: str
    params: Dict[str, Any]
    caller: str

    @classmethod
    def from_tendermint(cls, event: TendermintEvent) -> Optional["AgentActionEvent"]:
        attrs = event.attributes
        if "wasm.agent" not in attrs:
            return None
        try:
            params = json.loads(attrs.get("wasm.params", "{}"))
        except json.JSONDecodeError:
            params = {}
        return cls(
            contract=attrs.get("wasm._contract_address", ""),
            height=event.height, timestamp=event.timestamp, raw=event.raw,
            agent_id=attrs.get("wasm.agent", ""),
            action=attrs.get("wasm.action", ""),
            params=params,
            caller=attrs.get("wasm.caller", ""),
        )


@dataclass
class AgentRegistrationEvent(BaseEvent):
    agent_id: str
    owner: str
    name: str
    model: str
    capabilities: List[str]

    @classmethod
    def from_tendermint(cls, event: TendermintEvent) -> Optional["AgentRegistrationEvent"]:
        attrs = event.attributes
        if attrs.get("wasm.action") != "register_agent":
            return None
        return cls(
            contract=attrs.get("wasm._contract_address", ""),
            height=event.height, timestamp=event.timestamp, raw=event.raw,
            agent_id=attrs.get("wasm.agent_id", ""),
            owner=attrs.get("wasm.owner", ""),
            name=attrs.get("wasm.name", ""),
            model=attrs.get("wasm.model", ""),
            capabilities=attrs.get("wasm.capability", "").split(","),
        )


@dataclass
class A2AMessageEvent(BaseEvent):
    sender: str
    recipient: str
    message_type: str
    payload: Dict[str, Any]

    @classmethod
    def from_tendermint(cls, event: TendermintEvent) -> Optional["A2AMessageEvent"]:
        attrs = event.attributes
        if attrs.get("wasm.action") != "a2a_message":
            return None
        try:
            payload = json.loads(attrs.get("wasm.payload", "{}"))
        except json.JSONDecodeError:
            payload = {}
        return cls(
            contract=attrs.get("wasm._contract_address", ""),
            height=event.height, timestamp=event.timestamp, raw=event.raw,
            sender=attrs.get("wasm.sender", ""),
            recipient=attrs.get("wasm.recipient", ""),
            message_type=attrs.get("wasm.message_type", ""),
            payload=payload,
        )


@dataclass
class PaymentEvent(BaseEvent):
    agent_id: str
    payer: str
    payee: str
    amount: str
    purpose: str

    @classmethod
    def from_tendermint(cls, event: TendermintEvent) -> Optional["PaymentEvent"]:
        attrs = event.attributes
        if attrs.get("wasm.action") not in ("agent_payment", "payment"):
            return None
        return cls(
            contract=attrs.get("wasm._contract_address", ""),
            height=event.height, timestamp=event.timestamp, raw=event.raw,
            agent_id=attrs.get("wasm.agent_id", ""),
            payer=attrs.get("wasm.payer", ""), payee=attrs.get("wasm.payee", ""),
            amount=attrs.get("wasm.amount", "0"), purpose=attrs.get("wasm.purpose", ""),
        )


@dataclass
class StakingEvent(BaseEvent):
    delegator: str
    validator: str
    amount: str
    action_type: str

    @classmethod
    def from_tendermint(cls, event: TendermintEvent) -> Optional["StakingEvent"]:
        attrs = event.attributes
        if "staking.amount" not in attrs:
            return None
        raw = str(event.raw).lower()
        action = "delegate"
        if "undelegate" in raw:
            action = "undelegate"
        elif "redelegate" in raw:
            action = "redelegate"
        return cls(
            contract="staking", height=event.height, timestamp=event.timestamp, raw=event.raw,
            delegator=attrs.get("staking.delegator", ""),
            validator=attrs.get("staking.validator", ""),
            amount=attrs.get("staking.amount", "0"),
            action_type=action,
        )


@dataclass
class GovernanceEvent(BaseEvent):
    proposal_id: int
    proposer: str
    title: str
    action_type: str

    @classmethod
    def from_tendermint(cls, event: TendermintEvent) -> Optional["GovernanceEvent"]:
        attrs = event.attributes
        if "governance.proposal_id" not in attrs:
            return None
        return cls(
            contract="governance", height=event.height, timestamp=event.timestamp, raw=event.raw,
            proposal_id=int(attrs.get("proposal_id", 0)),
            proposer=attrs.get("proposer", ""),
            title=attrs.get("title", ""),
            action_type=attrs.get("governance.action", "submit"),
        )

4.3 事件解析器与路由器


class ContractEventParser:
    """合约事件解析器 —— 将原始 Tendermint 事件解析为结构化对象"""

    PARSERS = {
        "wasm.transfer": TransferEvent,
        "wasm.agent_action": AgentActionEvent,
        "wasm.register_agent": AgentRegistrationEvent,
        "wasm.a2a_message": A2AMessageEvent,
        "wasm.agent_payment": PaymentEvent,
        "staking": StakingEvent,
        "governance": GovernanceEvent,
    }

    def __init__(self):
        self._custom: Dict[str, Type[BaseEvent]] = {}

    def register(self, event_type: str, parser: Type[BaseEvent]):
        self._custom[event_type] = parser

    def parse(self, event: TendermintEvent) -> Optional[BaseEvent]:
        for key, cls in self._custom.items():
            if key in event.attributes:
                parsed = cls.from_tendermint(event)
                if parsed:
                    return parsed
        for key, cls in self.PARSERS.items():
            if any(key in k for k in event.attributes):
                parsed = cls.from_tendermint(event)
                if parsed:
                    return parsed
        return None

    def parse_multiple(self, events: List[TendermintEvent]) -> List[BaseEvent]:
        return [p for e in events if (p := self.parse(e)) is not None]


class EventRouter:
    """事件路由分发器"""

    def __init__(self, parser: ContractEventParser):
        self.parser = parser
        self._routes: Dict[Type[BaseEvent], List[Callable]] = {}

    def on(self, event_type: Type[BaseEvent]):
        def decorator(handler):
            self._routes.setdefault(event_type, []).append(handler)
            return handler
        return decorator

    async def route(self, event: TendermintEvent):
        parsed = self.parser.parse(event)
        if parsed is None:
            return
        for handler in self._routes.get(type(parsed), []):
            try:
                if asyncio.iscoroutinefunction(handler):
                    await handler(parsed)
                else:
                    handler(parsed)
            except Exception as e:
                logger.error(f"Route error for {type(parsed).__name__}: {e}")

4.4 使用示例


async def parsing_example():
    parser = ContractEventParser()
    router = EventRouter(parser)

    @router.on(TransferEvent)
    async def handle_transfer(event: TransferEvent):
        print(f"[转账] {event.from_addr[:12]} -> {event.to_addr[:12]}: {event.amount}{event.denom}")

    @router.on(AgentActionEvent)
    async def handle_action(event: AgentActionEvent):
        print(f"[Agent] {event.agent_id[:12]}: {event.action} params={event.params}")

    @router.on(A2AMessageEvent)
    async def handle_a2a(event: A2AMessageEvent):
        print(f"[A2A] {event.sender[:12]} -> {event.recipient[:12]}: {event.message_type}")

    mock = TendermintEvent(
        type="Tx", height=123456, tx_hash="ABCD",
        attributes={
            "wasm._contract_address": "msg1agent...",
            "wasm.agent": "agent-001",
            "wasm.action": "analyze_market",
            "wasm.params": '{"symbol": "MSG"}',
            "wasm.caller": "msg1user...",
        },
        raw={},
    )
    await router.route(mock)

5. 索引器 SDK 集成

5.1 RESTful 索引器客户端

MSG Chain 提供标准化的 RESTful 索引器 API,支持复杂聚合查询、全文搜索和统计分析。

import aiohttp
from dataclasses import dataclass
from datetime import datetime


@dataclass
class PaginatedResponse:
    items: List[Dict[str, Any]]
    total: int
    page: int
    per_page: int
    has_more: bool


class IndexerClient:
    """索引器 REST API 客户端"""

    def __init__(self, base_url: str = "https://api.msgchain.org"):
        self.base_url = base_url
        self._session: Optional[aiohttp.ClientSession] = None
        self._rps = 10
        self._tokens = 10
        self._last_refill = datetime.now()

    async def _session_get(self) -> aiohttp.ClientSession:
        if self._session is None:
            self._session = aiohttp.ClientSession(
                base_url=self.base_url,
                timeout=aiohttp.ClientTimeout(total=30),
            )
        return self._session

    async def _rate_limit(self):
        now = datetime.now()
        elapsed = (now - self._last_refill).total_seconds()
        self._tokens = min(self._rps, self._tokens + elapsed * self._rps)
        self._last_refill = now
        if self._tokens < 1:
            await asyncio.sleep(1 / self._rps)
            self._tokens = 0
        else:
            self._tokens -= 1

    async def _get(self, path: str, params: dict = None) -> dict:
        await self._rate_limit()
        session = await self._session_get()
        async with session.get(f"/indexer/v1{path}", params=params) as resp:
            resp.raise_for_status()
            return await resp.json()

    async def _post(self, path: str, data: dict = None) -> dict:
        await self._rate_limit()
        session = await self._session_get()
        async with session.post(f"/indexer/v1{path}", json=data) as resp:
            resp.raise_for_status()
            return await resp.json()

    # ─── Agent 查询 ──────────────────────────────

    async def query_agents(
        self, status: str = None, capability: str = None,
        owner: str = None, limit: int = 20, offset: int = 0,
    ) -> PaginatedResponse:
        params = {"limit": limit, "offset": offset}
        if status:
            params["status"] = status
        if capability:
            params["capability"] = capability
        if owner:
            params["owner"] = owner
        result = await self._get("/agents", params)
        return PaginatedResponse(
            items=result.get("agents", []),
            total=result.get("total", 0),
            page=result.get("page", 1),
            per_page=limit,
            has_more=(offset + limit) < result.get("total", 0),
        )

    async def query_agents_by_capability(self, capability: str, limit: int = 20) -> List[Dict]:
        result = await self._get("/agents", {"capability": capability, "limit": limit})
        return result.get("agents", [])

    async def query_agent_detail(self, agent_id: str) -> Dict:
        return await self._get(f"/agents/{agent_id}")

    async def query_agent_stats(self, agent_id: str) -> Dict:
        return await self._get(f"/agents/{agent_id}/stats")

    async def query_agent_interactions(self, agent_id: str, time_range: str = "7d") -> List[Dict]:
        result = await self._get(f"/agents/{agent_id}/interactions", {"range": time_range})
        return result.get("interactions", [])

    async def query_agent_earnings(self, agent_id: str, time_range: str = "30d") -> Dict:
        return await self._get(f"/agents/{agent_id}/earnings", {"range": time_range})

    # ─── 市场数据 ──────────────────────────────

    async def query_market_activity(self, time_range: str = "24h", granularity: str = "1h") -> Dict:
        return await self._get("/market/activity", {"range": time_range, "granularity": granularity})

    async def query_top_agents(self, metric: str = "transactions", limit: int = 10, time_range: str = "7d") -> List[Dict]:
        result = await self._get("/market/top-agents", {"metric": metric, "limit": limit, "range": time_range})
        return result.get("agents", [])

    async def query_network_stats(self) -> Dict:
        return await self._get("/network/stats")

    # ─── 交易搜索 ──────────────────────────────

    async def search_transactions(self, query: Dict, limit: int = 50, offset: int = 0) -> PaginatedResponse:
        payload = {**query, "limit": limit, "offset": offset}
        result = await self._post("/tx/search", payload)
        return PaginatedResponse(
            items=result.get("transactions", []),
            total=result.get("total", 0),
            page=result.get("page", 1),
            per_page=limit,
            has_more=(offset + limit) < result.get("total", 0),
        )

    async def search_txs_by_contract(
        self, contract_addr: str, action: str = None, from_height: int = None, to_height: int = None, limit: int = 50,
    ) -> List[Dict]:
        query = {"contract": contract_addr}
        if action:
            query["action"] = action
        if from_height:
            query["from_height"] = from_height
        if to_height:
            query["to_height"] = to_height
        result = await self.search_transactions(query, limit=limit)
        return result.items

    async def search_txs_by_agent(self, agent_id: str, action: str = None, limit: int = 50) -> List[Dict]:
        query = {"agent_id": agent_id}
        if action:
            query["action"] = action
        result = await self.search_transactions(query, limit=limit)
        return result.items

    async def full_text_search(self, text: str, fields: List[str] = None, limit: int = 20) -> List[Dict]:
        params = {"q": text, "limit": limit}
        if fields:
            params["fields"] = ",".join(fields)
        result = await self._get("/search", params)
        return result.get("results", [])

    # ─── 统计与分析 ──────────────────────────────

    async def query_analytics(self, metric: str, dimensions: List[str] = None, filters: Dict = None, time_range: str = "7d") -> Dict:
        params = {"metric": metric, "range": time_range}
        if dimensions:
            params["dimensions"] = ",".join(dimensions)
        if filters:
            params["filters"] = json.dumps(filters)
        return await self._get("/analytics", params)

    async def close(self):
        if self._session and not self._session.closed:
            await self._session.close()

5.2 GraphQL 接口


class GraphQLIndexerClient:
    """GraphQL 索引器客户端"""

    def __init__(self, endpoint: str = "https://api.msgchain.org/indexer/v1/graphql"):
        self.endpoint = endpoint
        self._session: Optional[aiohttp.ClientSession] = None

    async def query(self, query_str: str, variables: dict = None) -> Dict:
        if not self._session:
            self._session = aiohttp.ClientSession()
        payload = {"query": query_str}
        if variables:
            payload["variables"] = variables
        async with self._session.post(self.endpoint, json=payload) as resp:
            resp.raise_for_status()
            return await resp.json()

    async def get_agents_with_stats(self, limit: int = 10):
        q = """
        query GetAgents($limit: Int!) {
            agents(limit: $limit) {
                id name owner capabilities status
                stats { totalTransactions totalEarnings avgResponseTime successRate }
            }
        }"""
        return await self.query(q, {"limit": limit})

    async def get_network_overview(self):
        q = """
        query { network {
            activeAgents totalTransactions totalStaked
            avgBlockTime activeValidators latestHeight
        }}"""
        return await self.query(q)

    async def close(self):
        if self._session:
            await self._session.close()

5.3 使用示例


async def indexer_example():
    client = IndexerClient()

    agents = await client.query_agents(status="active", limit=5)
    print(f"Active agents: {agents.total}")

    market = await client.query_market_activity(time_range="7d")
    print(f"7d volume: {market.get('total_volume', 'N/A')}")

    stats = await client.query_agent_stats("agent-001")
    print(f"TX count: {stats.get('total_transactions', 0)}")

    results = await client.search_txs_by_contract("msg1contract...", action="agent_action")
    print(f"Matching txs: {len(results)}")

    await client.close()

6. 自定义索引部署

6.1 自定义索引器架构

当标准索引器无法满足 Agent 的特定需求时,可以部署自定义索引器。典型架构:

MSG Chain RPC 节点
    │
    ▼
区块流 ──→ 事件提取 ──→ 过滤 ──→ 转换 ──→ 存储层
                                                │
                                                ▼
                                     Agent 查询 API / 回调

6.2 区块获取与事件提取

import asyncpg
from dataclasses import dataclass, field
from collections import defaultdict


@dataclass
class ParsedBlock:
    height: int
    hash: str
    timestamp: datetime
    txs: List[Dict]
    events: List[Dict]
    proposer: str
    num_txs: int


@dataclass
class IndexerState:
    current_height: int = 0
    target_height: int = 0
    blocks_processed: int = 0
    events_indexed: int = 0
    errors: List[str] = field(default_factory=list)


class BlockFetcher:
    """从 RPC 节点拉取区块"""

    def __init__(self, rpc: str = "https://rpc.msgchain.org"):
        self.rpc = rpc
        self._session: Optional[aiohttp.ClientSession] = None

    async def fetch_block(self, height: int) -> Optional[ParsedBlock]:
        session = await self._get_session()
        try:
            async with session.get(f"{self.rpc}/block?height={height}") as resp:
                if resp.status != 200:
                    return None
                return self._parse(await resp.json(), height)
        except Exception:
            return None

    async def fetch_block_results(self, height: int) -> Optional[Dict]:
        session = await self._get_session()
        try:
            async with session.get(f"{self.rpc}/block_results?height={height}") as resp:
                return await resp.json() if resp.status == 200 else None
        except Exception:
            return None

    async def get_latest_height(self) -> int:
        session = await self._get_session()
        async with session.get(f"{self.rpc}/status") as resp:
            data = await resp.json()
            return int(data["result"]["sync_info"]["latest_block_height"])

    async def _get_session(self):
        if not self._session:
            self._session = aiohttp.ClientSession()
        return self._session

    def _parse(self, data: Dict, height: int) -> ParsedBlock:
        header = data.get("result", {}).get("block", {}).get("header", {})
        return ParsedBlock(
            height=height,
            hash=header.get("last_block_id", {}).get("hash", ""),
            timestamp=datetime.now(),
            txs=[], events=[], proposer=header.get("proposer_address", ""),
            num_txs=0,
        )

    async def close(self):
        if self._session:
            await self._session.close()


class EventExtractor:
    """从区块结果中提取结构化事件"""

    def __init__(self, filters: Optional[List[Callable]] = None):
        self.filters = filters or []

    async def extract(self, block_results: Dict, height: int) -> List[Dict]:
        events = []
        for tx_result in block_results.get("result", {}).get("txs_results", []):
            tx_hash = tx_result.get("hash", "")
            for event in tx_result.get("events", []):
                parsed = self._parse_event(event, height, tx_hash)
                if parsed and self._accept(parsed):
                    events.append(parsed)
        return events

    def _parse_event(self, event: Dict, height: int, tx_hash: str) -> Optional[Dict]:
        attrs = {}
        for attr in event.get("attributes", []):
            k, v = attr.get("key", ""), attr.get("value", "")
            if isinstance(k, bytes): k = k.decode()
            if isinstance(v, bytes): v = v.decode()
            attrs[k] = v
        if not event.get("type"):
            return None
        return {"type": event["type"], "height": height, "tx_hash": tx_hash, "attributes": attrs}

    def _accept(self, event: Dict) -> bool:
        return all(f(event) for f in self.filters) if self.filters else True

6.3 自定义索引器


class AgentIndexer:
    """自定义 Agent 索引器 —— 过滤并存储 Agent 相关事件到 PostgreSQL"""

    def __init__(
        self,
        agent_id: str,
        db_dsn: str = "postgresql://localhost:5432/msg_indexer",
        rpc_endpoint: str = "https://rpc.msgchain.org",
        start_height: Optional[int] = None,
    ):
        self.agent_id = agent_id
        self.db_dsn = db_dsn
        self.fetcher = BlockFetcher(rpc_endpoint)
        self.extractor = EventExtractor()
        self.state = IndexerState()
        self.start_height = start_height or 0
        self._pool: Optional[asyncpg.Pool] = None
        self._running = False

    async def initialize(self):
        self._pool = await asyncpg.create_pool(self.db_dsn, min_size=2, max_size=10)
        async with self._pool.acquire() as conn:
            await conn.execute("""
                CREATE TABLE IF NOT EXISTS agent_events (
                    id BIGSERIAL PRIMARY KEY,
                    agent_id VARCHAR(255) NOT NULL,
                    event_type VARCHAR(255) NOT NULL,
                    tx_hash VARCHAR(255),
                    height BIGINT NOT NULL,
                    data JSONB NOT NULL DEFAULT '{}',
                    indexed_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
                )
            """)
            await conn.execute("CREATE INDEX IF NOT EXISTS idx_ae_agent ON agent_events(agent_id)")
            await conn.execute("CREATE INDEX IF NOT EXISTS idx_ae_height ON agent_events(height)")
            await conn.execute("""
                CREATE TABLE IF NOT EXISTS indexer_state (
                    id INTEGER PRIMARY KEY DEFAULT 1,
                    last_indexed_height BIGINT NOT NULL DEFAULT 0
                )
            """)
            await conn.execute(
                "INSERT INTO indexer_state (id, last_indexed_height) VALUES (1, $1) ON CONFLICT DO NOTHING",
                self.start_height,
            )

    async def process_block(self, height: int) -> int:
        block_results = await self.fetcher.fetch_block_results(height)
        if not block_results:
            return 0
        events = await self.extractor.extract(block_results, height)
        agent_events = [e for e in events if self._is_relevant(e)]
        if agent_events:
            async with self._pool.acquire() as conn:
                for e in agent_events:
                    await conn.execute(
                        "INSERT INTO agent_events (agent_id, event_type, tx_hash, height, data) VALUES ($1,$2,$3,$4,$5)",
                        self.agent_id, e["type"], e.get("tx_hash", ""), e["height"], json.dumps(e),
                    )
                    self.state.events_indexed += 1
        self.state.blocks_processed += 1
        self.state.current_height = height
        return len(agent_events)

    def _is_relevant(self, event: Dict) -> bool:
        attrs = event.get("attributes", {})
        agent_attr = attrs.get("agent", "") or attrs.get("wasm.agent", "")
        if agent_attr:
            return True
        contract = attrs.get("_contract_address", "")
        return bool(contract and self.agent_id in contract)

    async def run(self, start_height: Optional[int] = None):
        self._running = True
        current = start_height or self.start_height
        while self._running:
            latest = await self.fetcher.get_latest_height()
            self.state.target_height = latest
            if current > latest:
                await asyncio.sleep(1)
                continue
            try:
                await self.process_block(current)
                async with self._pool.acquire() as conn:
                    await conn.execute(
                        "UPDATE indexer_state SET last_indexed_height = $1 WHERE id = 1", current,
                    )
                current += 1
            except Exception as e:
                logger.error(f"Error at height {current}: {e}")
                self.state.errors.append(f"Height {current}: {e}")
                current += 1
                await asyncio.sleep(2)

    async def run_catchup(self, from_height: int, to_height: Optional[int] = None):
        if to_height is None:
            to_height = await self.fetcher.get_latest_height()
        logger.info(f"Catchup {from_height} -> {to_height} ({to_height - from_height} blocks)")
        for h in range(from_height, to_height + 1):
            if not self._running:
                break
            await self.process_block(h)
            if h % 100 == 0:
                pct = (h - from_height) / (to_height - from_height) * 100
                logger.info(f"Catchup {pct:.1f}% (height {h})")

    async def stop(self):
        self._running = False
        if self._pool:
            await self._pool.close()
        await self.fetcher.close()

6.4 ClickHouse 存储适配


class ClickHouseEventStore:
    """ClickHouse 事件存储 —— 适用于大规模时序分析"""

    CREATE_SQL = """
    CREATE TABLE IF NOT EXISTS agent_events (
        event_id UInt64, agent_id String, event_type String,
        tx_hash String, height UInt64, data String, indexed_at DateTime
    ) ENGINE = MergeTree()
    PARTITION BY toYYYYMM(indexed_at)
    ORDER BY (agent_id, height, event_type)
    """

    def __init__(self, dsn: str = "clickhouse://localhost:9000/msg"):
        self.dsn = dsn
        self._client = None

    async def connect(self):
        from clickhouse_driver import AsyncClient
        self._client = AsyncClient.from_url(self.dsn)
        await self._client.execute(self.CREATE_SQL)

    async def insert(self, events: List[Dict]):
        if not events:
            return
        rows = [(e.get("id", i), e["agent_id"], e["event_type"], e.get("tx_hash", ""),
                 e["height"], json.dumps(e["data"]), e.get("indexed_at", datetime.now()))
                for i, e in enumerate(events)]
        await self._client.execute("INSERT INTO agent_events VALUES", rows)

    async def close(self):
        if self._client:
            await self._client.disconnect()

6.5 部署示例


async def deploy_indexer():
    indexer = AgentIndexer(
        agent_id="agent-001",
        db_dsn="postgresql://user:pass@localhost:5432/msg_indexer",
    )
    await indexer.initialize()

    await indexer.run_catchup(from_height=1_000_000)
    await indexer.run(start_height=await indexer.get_last_indexed_height())

    try:
        while True:
            await asyncio.sleep(60)
    except KeyboardInterrupt:
        await indexer.stop()

7. 数据聚合与分析

7.1 跨合约数据聚合


class DataAggregator:
    """跨合约数据聚合器"""

    def __init__(self, onchain: OnChainQuerier, indexer: IndexerClient):
        self.onchain = onchain
        self.indexer = indexer

    async def get_agent_profile(self, agent_id: str) -> Dict:
        onchain_data, indexer_data = await asyncio.gather(
            self._get_onchain_data(agent_id),
            self.indexer.query_agent_detail(agent_id),
            return_exceptions=True,
        )
        profile = {"agent_id": agent_id, "aggregated_at": datetime.now().isoformat()}
        if not isinstance(onchain_data, Exception):
            profile.update(onchain_data)
        if not isinstance(indexer_data, Exception):
            profile.update(indexer_data)
        return profile

    async def _get_onchain_data(self, agent_id: str) -> Dict:
        try:
            balance = await self.onchain.query_bank_balance(agent_id)
            return {"balance": str(balance.amount), "denom": "umsg"}
        except Exception:
            return {}

    async def get_market_snapshot(self, time_ranges: List[str] = ("1h", "24h", "7d")) -> Dict:
        results = await asyncio.gather(
            *[self.indexer.query_market_activity(r) for r in time_ranges],
            return_exceptions=True,
        )
        return {r: res for r, res in zip(time_ranges, results) if not isinstance(res, Exception)}

    async def get_agent_rankings(self, metrics: List[str] = ("transactions", "earnings")) -> Dict:
        results = await asyncio.gather(
            *[self.indexer.query_top_agents(metric=m) for m in metrics],
            return_exceptions=True,
        )
        return {m: res for m, res in zip(metrics, results) if not isinstance(res, Exception)}

7.2 Agent 性能分析


class AgentPerformanceAnalyzer:
    """Agent 性能分析器"""

    def __init__(self, indexer: IndexerClient):
        self.indexer = indexer

    async def compute_score(self, agent_id: str) -> Dict:
        stats, interactions = await asyncio.gather(
            self.indexer.query_agent_stats(agent_id),
            self.indexer.query_agent_interactions(agent_id),
            return_exceptions=True,
        )
        if isinstance(stats, Exception):
            return {"agent_id": agent_id, "error": str(stats)}
        total = stats.get("total_transactions", 0)
        success = stats.get("successful_transactions", 0) or stats.get("total_transactions", 0) - stats.get("failed_transactions", 0)
        response_time = stats.get("avg_response_time_ms", 0)
        success_rate = success / max(total, 1)
        reliability = min(success_rate * 100, 100)
        efficiency = max(0, 100 - (response_time / 100))
        engagement = min(len(interactions) if not isinstance(interactions, Exception) else 0 * 10, 100)
        overall = reliability * 0.4 + efficiency * 0.3 + engagement * 0.3
        return {
            "agent_id": agent_id,
            "overall_score": round(overall, 2),
            "dimensions": {"reliability": round(reliability, 2), "efficiency": round(efficiency, 2), "engagement": round(engagement, 2)},
            "evaluated_at": datetime.now().isoformat(),
        }

    async def compare(self, agent_ids: List[str]) -> List[Dict]:
        scores = await asyncio.gather(*[self.compute_score(a) for a in agent_ids], return_exceptions=True)
        return sorted([s for s in scores if not isinstance(s, Exception)], key=lambda x: x["overall_score"], reverse=True)

7.3 时序趋势分析


class TimeSeriesAnalyzer:
    """时序趋势分析"""

    def __init__(self, indexer: IndexerClient):
        self.indexer = indexer

    async def get_trend(self, agent_id: str, metric: str = "transactions", time_range: str = "7d") -> Dict:
        return await self.indexer.query_analytics(
            metric=metric, dimensions=["hour"], filters={"agent_id": agent_id}, time_range=time_range,
        )

    async def detect_anomalies(self, agent_id: str, metric: str = "transactions", window: int = 24, threshold: float = 2.0) -> List[Dict]:
        data = await self.get_trend(agent_id, metric, f"{window}h")
        values = [v for v in data.get("series", []) if isinstance(v, (int, float))]
        if not values:
            return []
        mean = sum(values) / len(values)
        std = (sum((v - mean) ** 2 for v in values) / len(values)) ** 0.5
        anomalies = []
        for point in zip(data.get("timestamps", []), values):
            if abs(point[1] - mean) > threshold * std:
                anomalies.append({"timestamp": point[0], "value": point[1], "z_score": (point[1] - mean) / max(std, 0.01)})
        return anomalies

8. 完整 Agent 数据流示例

8.1 数据驱动 Agent 框架


@dataclass
class DataDrivenAgentConfig:
    agent_id: str
    owner: str = ""
    capabilities: List[str] = field(default_factory=list)
    rpc_endpoint: str = "https://rpc.msgchain.org"
    lcd_endpoint: str = "https://lcd.msgchain.org"
    ws_endpoint: str = "wss://rpc.msgchain.org/websocket"
    indexer_url: str = "https://api.msgchain.org"
    poll_interval: float = 5.0


class DataDrivenAgent:
    """基于链上数据的 AI Agent —— 整合查询、事件订阅、索引器"""

    def __init__(self, config: DataDrivenAgentConfig):
        self.config = config
        self.agent_id = config.agent_id
        self._msg_client: Optional[MSGChainClient] = None
        self._querier: Optional[OnChainQuerier] = None
        self._cached_querier: Optional[CachedOnChainQuerier] = None
        self._event_subscriber: Optional[EventSubscriber] = None
        self._event_parser: Optional[ContractEventParser] = None
        self._event_router: Optional[EventRouter] = None
        self._indexer: Optional[IndexerClient] = None
        self._aggregator: Optional[DataAggregator] = None
        self._analyzer: Optional[TimeSeriesAnalyzer] = None
        self._performance: Optional[AgentPerformanceAnalyzer] = None
        self._running = False
        self._tasks: Set[asyncio.Task] = set()
        self._memory: Dict[str, Any] = {}
        self._event_log: List[BaseEvent] = []

    async def initialize(self):
        logger.info(f"Initializing agent: {self.agent_id}")

        self._msg_client = MSGChainClient(self.config.rpc_endpoint, self.config.lcd_endpoint)
        await self._msg_client.connect()
        self._querier = OnChainQuerier(self._msg_client.client)
        self._cached_querier = CachedOnChainQuerier(self._querier, ttl_seconds=30)

        self._event_subscriber = EventSubscriber(self.agent_id, self.config.ws_endpoint)
        self._event_parser = ContractEventParser()
        self._event_router = EventRouter(self._event_parser)
        self._indexer = IndexerClient(self.config.indexer_url)

        self._aggregator = DataAggregator(self._querier, self._indexer)
        self._analyzer = TimeSeriesAnalyzer(self._indexer)
        self._performance = AgentPerformanceAnalyzer(self._indexer)

        self._register_default_handlers()
        logger.info(f"Agent {self.agent_id} initialized")
        return self

    def _register_default_handlers(self):
        @self._event_router.on(TransferEvent)
        async def _(e: TransferEvent):
            self._event_log.append(e)
            self._memory.setdefault("transfers", []).append(e)
            if self._is_relevant(e):
                await self._on_relevant_transfer(e)

        @self._event_router.on(AgentActionEvent)
        async def _(e: AgentActionEvent):
            self._event_log.append(e)
            self._memory.setdefault("actions", []).append(e)

        @self._event_router.on(A2AMessageEvent)
        async def _(e: A2AMessageEvent):
            self._event_log.append(e)
            if e.recipient == self.agent_id:
                await self._process_a2a_message(e)

        @self._event_router.on(PaymentEvent)
        async def _(e: PaymentEvent):
            self._event_log.append(e)
            self._memory.setdefault("payments", []).append(e)

    def _is_relevant(self, event: TransferEvent) -> bool:
        return self.agent_id in event.to_addr or self.agent_id in event.from_addr

    async def _on_relevant_transfer(self, event: TransferEvent):
        pass

    async def _process_a2a_message(self, event: A2AMessageEvent):
        payload = event.payload
        if event.message_type == "task_request":
            result = await self.execute_task(payload.get("task", ""), payload.get("params", {}))
            logger.info(f"Task result for {event.sender}: {result}")
        elif event.message_type == "query":
            pass

    async def execute_task(self, task: str, params: Dict) -> Dict:
        if task == "query_market":
            return await self._aggregator.get_market_snapshot()
        elif task == "analyze_performance":
            return await self._performance.compute_score(params.get("agent_id", self.agent_id))
        elif task == "get_balance":
            balance = await self._querier.query_bank_balance(params.get("address", self.agent_id))
            return {"balance": str(balance.amount), "denom": "umsg"}
        return {"status": "unknown_task"}

    async def subscribe_to_events(self):
        sub = self._event_subscriber
        await sub.subscribe_to_agent_actions(
            handler=lambda e: asyncio.create_task(self._event_router.route(e))
        )
        await sub.subscribe_to_transfers(
            to_addr=self.agent_id,
            handler=lambda e: asyncio.create_task(self._event_router.route(e))
        )
        await sub.subscribe(
            "tm.event='Tx' AND wasm.action='a2a_message'",
            handler=lambda e: asyncio.create_task(self._event_router.route(e))
        )
        await sub.subscribe_to_new_blocks(
            handler=lambda e: self._memory.update({"last_height": e.height})
        )

    async def start(self):
        self._running = True
        await self.subscribe_to_events()
        self._tasks.add(asyncio.create_task(self._event_subscriber.event_loop()))
        self._tasks.add(asyncio.create_task(self._main_loop()))
        logger.info(f"Agent {self.agent_id} running")

    async def _main_loop(self):
        while self._running:
            try:
                await self.execute_strategy()
            except Exception as e:
                logger.error(f"Strategy error: {e}")
            await asyncio.sleep(self.config.poll_interval)

    async def execute_strategy(self):
        """子类覆盖此方法实现具体策略逻辑"""
        pass

    async def query_onchain(self, contract: str, query: dict) -> dict:
        return await self._cached_querier.query_contract_smart(contract, query)

    async def query_balance(self):
        return await self._querier.query_bank_balance(self.agent_id)

    async def get_stats(self) -> Dict:
        return {
            "agent_id": self.agent_id,
            "running": self._running,
            "events_logged": len(self._event_log),
            "connected": self._event_subscriber.is_connected if self._event_subscriber else False,
        }

    async def stop(self):
        self._running = False
        if self._event_subscriber:
            await self._event_subscriber.stop()
        if self._msg_client:
            await self._msg_client.close()
        if self._indexer:
            await self._indexer.close()
        for t in self._tasks:
            t.cancel()
        self._tasks.clear()
        logger.info(f"Agent {self.agent_id} stopped")

8.2 市场分析 Agent 示例


class MarketAnalysisAgent(DataDrivenAgent):
    """市场分析 Agent —— 监听链上交易,定期分析市场趋势"""

    async def execute_strategy(self):
        stats = await self._indexer.query_network_stats()
        last = await self._memory.get("last_height", 0)
        current = stats.get("latest_height", 0)
        if current > last + 10:
            snapshot = await self._aggregator.get_market_snapshot()
            self._memory["last_snapshot"] = snapshot
            self._memory["last_height"] = current
            logger.debug(f"Snapshot updated at height {current}")

    async def _on_relevant_transfer(self, event: TransferEvent):
        amount = int(event.amount)
        if amount >= 10_000_000:
            logger.info(f"Large transfer: {amount}umsg")
            score = await self._performance.compute_score(self.agent_id)
            logger.info(f"Performance: {score['overall_score']}")

8.3 启动脚本


async def main():
    config = DataDrivenAgentConfig(
        agent_id="msg1agentdemo...",
        capabilities=["market_analysis"],
    )
    agent = MarketAnalysisAgent(config)
    await agent.initialize()
    try:
        await agent.start()
        print(f"Agent {agent.agent_id} running. Press Ctrl+C to stop.")
        while True:
            await asyncio.sleep(10)
            balance = await agent.query_balance()
            print(f"[Heartbeat] Balance: {balance}")
    except KeyboardInterrupt:
        print("\nShutting down...")
    finally:
        await agent.stop()


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

8.4 多 Agent 协调


class AgentCoordinator:
    """多 Agent 协调器"""

    def __init__(self):
        self.agents: Dict[str, DataDrivenAgent] = {}

    async def register(self, agent_id: str, agent: DataDrivenAgent):
        self.agents[agent_id] = agent
        await agent.initialize()

    async def start_all(self):
        await asyncio.gather(*[a.start() for a in self.agents.values()], return_exceptions=True)

    async def stop_all(self):
        for a in self.agents.values():
            await a.stop()

    async def network_status(self) -> Dict:
        statuses = {}
        for aid, agent in self.agents.items():
            try:
                statuses[aid] = await agent.get_stats()
            except Exception as e:
                statuses[aid] = {"error": str(e)}
        return {"total": len(self.agents), "agents": statuses, "timestamp": datetime.now().isoformat()}

附录

A. 常见问题

Q: 直接查询链上和通过索引器查询有什么区别?

直接查询通过 CosmWasm RPC 获取实时状态,适合小规模、即时性高的查询。索引器查询通过预处理的数据库,适合复杂聚合、历史分析和全文搜索。

Q: WebSocket 断线后如何处理?

EventSubscriber 实现了自动重连机制,指数退避策略(1s → 2s → 4s → ... → 60s max),重连后自动恢复所有订阅。

Q: 自定义索引器需要什么基础设施?

需要 PostgreSQL 或 ClickHouse 数据库,以及能够访问 MSG Chain RPC 节点的服务。建议在独立服务器或容器中运行。

B. 相关资源

C. 类清单

类名 章节 功能
MSGChainClient 2.1 MSG Chain 客户端封装
OnChainQuerier 2.2 链上直接查询
TxQuerier 2.3 交易历史查询
CachedOnChainQuerier 2.4 带缓存的查询器
EventSubscriber 3.2 WebSocket 事件订阅
EventFilter 3.3 事件过滤链
PriorityEventRouter 3.3 优先级事件路由
TransferEvent 4.2 转账事件模型
AgentActionEvent 4.2 Agent 动作事件模型
A2AMessageEvent 4.2 Agent 间消息模型
ContractEventParser 4.3 合约事件解析器
EventRouter 4.3 事件路由分发
IndexerClient 5.1 索引器 REST API 客户端
GraphQLIndexerClient 5.2 GraphQL 索引器客户端
BlockFetcher 6.2 区块获取器
EventExtractor 6.2 事件提取器
AgentIndexer 6.3 自定义索引器
ClickHouseEventStore 6.4 ClickHouse 存储
DataAggregator 7.1 数据聚合器
AgentPerformanceAnalyzer 7.2 性能分析器
TimeSeriesAnalyzer 7.3 时序分析器
DataDrivenAgent 8.1 完整数据驱动 Agent
MarketAnalysisAgent 8.2 市场分析 Agent 示例
AgentCoordinator 8.4 多 Agent 协调器

所有代码均使用 msg 前缀(Bech32),Chain ID 为 msg-chain-1,最小单位 umsg。