MSG Chain AI Agent 链上数据索引与事件订阅指南
⚠️ No-Go Disclaimer: MSGChain 主网裁决为 No-Go。本文件所有内容反映的是开发阶段的技术设计,不代表主网未独立核验上线状态。生产部署状态请以白皮书为准:https://msgchain.org/whitepaper/
目录
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 网络参数
- Chain ID:
msg-chain-1 - Bech32 前缀:
msg - Denom:
umsg(最小单位),msg(1 msg = 1,000,000 umsg) - RPC 端点:
https://rpc.msgchain.org - LCD 端点:
https://lcd.msgchain.org - WebSocket 端点:
wss://rpc.msgchain.org/websocket - 索引器 API:
https://api.msgchain.org/indexer/v1/
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. 相关资源
- MSG Chain RPC:
https://rpc.msgchain.org - MSG Chain LCD:
https://lcd.msgchain.org - MSG Chain 浏览器:
https://explorer.msgchain.org - CosmWasm 文档:
https://docs.cosmwasm.com - Cosmos SDK 事件:
https://docs.cosmos.network/main/core/events
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。
