AI Agent 链上事件驱动架构与实时响应指南
⚠️ No-Go Disclaimer: MSGChain 主网裁决为 No-Go。本文件所有内容反映的是开发阶段的技术设计,不代表主网未独立核验上线状态。生产部署状态请以白皮书为准:https://msgchain.org/whitepaper/
1. 概述
1.1 什么是事件驱动架构
事件驱动架构(EDA)是一种软件范式,其中系统组件通过生产和消费事件通信。对于链上 AI Agent,这意味着 Agent 不再被动等待用户指令,而是主动监听区块链上的状态变更并实时响应。
在 MSG Chain 生态中,AI Agent 需处理四类事件:
- 交易事件:合约调用、代币转账、质押等
- Agent 生命周期事件:注册、支付、A2A 通信
- 治理事件:提案、投票、参数变更
- 自定义合约事件:用户部署合约触发的自定义事件
1.2 MSG Chain 事件体系
基于 Cosmos SDK 构建,事件分三个层级:
| 层级 | 说明 | 示例 |
|---|---|---|
| Tendermint 层 | 共识层事件(区块、交易) | tm.event='Tx' |
| Cosmos SDK 层 | 模块事件(转账、质押、治理) | transfer.sender |
| MSG Chain Agent 层 | Agent 专属事件 | agent.register, agent.a2a |
1.3 三种事件获取模式
WebSocket 订阅(实时推送):Agent 通过 WSS 连接到事件推送节点,毫秒级实时推送。推荐为主通道。
REST API 轮询(定时拉取):定期调用历史事件 API。适合低频场景和事件恢复。
Webhook 回调(被动接收):注册 Webhook URL,匹配事件时链节点主动回调。需公网可访问端点。
生产环境建议以 WebSocket 为主、REST API 兜底。
1.4 端点速查
| 端点 | 用途 | 协议 |
|---|---|---|
wss://rpc.msgchain.org/websocket |
Tendermint RPC WebSocket | WSS |
wss://api.msgchain.org/api/v1/agent/events/subscribe |
Agent 事件订阅 | WSS |
/agent/v1/events/history |
事件历史查询 | HTTPS |
/agent/v1/events/subscribe |
WebSocket 升级入口 | WSS/HTTPS |
1.5 链参数
- 链 ID:
msg-chain-1 - 地址前缀:
msg - 原生代币:
MSG(最小单位umsg, 1 MSG = 10^6 umsg)
2. 链上事件模型
2.1 事件类型总览
EVENT_TYPES = {
'tx': 'Transaction events (wasm, transfer, staking, governance)',
'wasm': 'Contract execution events',
'transfer': 'Bank transfer events',
'staking': 'Staking/delegation events',
'governance': 'Governance proposal/vote events',
'agent': 'Agent-specific events (register, payment, A2A)',
}
2.2 Tendermint 层事件
TENDERMINT_EVENTS = {
'Tx': '交易执行事件(包含所有模块事件)',
'NewBlock': '新区块产生事件',
'NewBlockHeader': '新区块头事件',
'ValidatorSetUpdates': '验证人集合更新事件',
}
# 原始事件结构
tendermint_event = {
'query': "tm.event='Tx' AND agent.action='register'",
'data': {
'type': 'tendermint/event/Tx',
'value': {
'TxResult': {
'height': '12345',
'txhash': '0xABCDEF1234567890...',
'events': [{
'type': 'transfer',
'attributes': [
{'key': 'sender', 'value': 'msg1qwerty...'},
{'key': 'recipient', 'value': 'msg1asdfgh...'},
{'key': 'amount', 'value': '1000000umsg'},
],
}],
},
},
},
}
2.3 Cosmos SDK 模块事件
转账事件由 MsgSend 触发:
transfer_event = {
'type': 'transfer',
'attributes': [
{'key': 'recipient', 'value': 'msg1qwerty...'},
{'key': 'sender', 'value': 'msg1asdfgh...'},
{'key': 'amount', 'value': '5000000umsg'},
],
}
WASM 合约事件属性由合约逻辑自定义,_contract_address 为保留属性:
wasm_event = {
'type': 'wasm',
'attributes': [
{'key': '_contract_address', 'value': 'msg1contract...'},
{'key': 'action', 'value': 'increment'},
{'key': 'count', 'value': '42'},
],
}
质押事件包括委托、解委托和奖励提取:
delegate_event = {
'type': 'delegate',
'attributes': [
{'key': 'validator', 'value': 'msgvaloper1validator...'},
{'key': 'delegator', 'value': 'msg1delegator...'},
{'key': 'amount', 'value': '1000000umsg'},
{'key': 'new_shares', 'value': '1000000'},
],
}
治理事件包含提案提交、投票和结果:
submit_proposal_event = {
'type': 'submit_proposal',
'attributes': [
{'key': 'proposal_id', 'value': '42'},
{'key': 'proposer', 'value': 'msg1proposer...'},
{'key': 'voting_start_time', 'value': '2026-07-10T00:00:00Z'},
],
}
2.4 Agent 专属事件
agent_register_event = {
'type': 'agent',
'attributes': [
{'key': 'action', 'value': 'register'},
{'key': 'agent_id', 'value': 'agent:msg1agentaddress...'},
{'key': 'owner', 'value': 'msg1owner...'},
{'key': 'agent_type', 'value': 'autonomous'},
],
}
agent_payment_event = {
'type': 'agent',
'attributes': [
{'key': 'action', 'value': 'payment'},
{'key': 'agent_id', 'value': 'agent:msg1agentaddress...'},
{'key': 'from', 'value': 'msg1payer...'},
{'key': 'amount', 'value': '1000000umsg'},
{'key': 'payment_type', 'value': 'subscription'},
],
}
agent_a2a_event = {
'type': 'agent',
'attributes': [
{'key': 'action', 'value': 'a2a'},
{'key': 'from_agent', 'value': 'agent:msg1agentA...'},
{'key': 'to_agent', 'value': 'agent:msg1agentB...'},
{'key': 'message_type', 'value': 'request'},
{'key': 'message_id', 'value': 'msg_a2a_通用维护记录_001'},
],
}
2.5 事件解析封装
import json
from typing import Optional
class MsgChainEvent:
"""MSG Chain 事件统一封装。"""
def __init__(self, raw_event: dict):
self.raw = raw_event
self._parsed = self._parse(raw_event)
def _parse(self, raw: dict) -> dict:
events_map = {}
for event in raw.get('events', []):
etype = event.get('type', '')
attrs = {}
for attr in event.get('attributes', []):
key = attr.get('key', '')
value = attr.get('value', '')
if key in attrs:
if not isinstance(attrs[key], list):
attrs[key] = [attrs[key]]
attrs[key].append(value)
else:
attrs[key] = value
events_map[etype] = attrs
return events_map
@property
def tx_hash(self) -> Optional[str]:
return self.raw.get('tx_hash')
@property
def height(self) -> Optional[int]:
h = self.raw.get('height')
return int(h) if h else None
def has_event(self, event_type: str) -> bool:
return event_type in self._parsed
def get_attributes(self, event_type: str) -> Optional[dict]:
return self._parsed.get(event_type)
2.6 事件过滤表达式
QUERY_EXAMPLES = {
'all_tx': "tm.event='Tx'",
'contract_events': "tm.event='Tx' AND wasm._contract_address='msg1contract...'",
'agent_registrations': "tm.event='Tx' AND agent.action='register'",
'large_transfers': "tm.event='Tx' AND transfer.amount>1000000000umsg",
'a2a_messages': "tm.event='Tx' AND agent.action='a2a' AND agent.to_agent='agent:msg1myagent...'",
}
COMPLEX_QUERY = (
"tm.event='Tx' AND "
"(agent.action='payment' OR agent.action='a2a') AND "
"agent.agent_id='agent:msg1myagent...'"
)
3. WebSocket 订阅
3.1 连接端点
import asyncio
import websockets
import json
from typing import AsyncGenerator, Optional
TENDERMINT_WS = 'wss://rpc.msgchain.org/websocket'
AGENT_WS = 'wss://api.msgchain.org/api/v1/agent/events/subscribe'
3.2 Tendermint RPC WebSocket 订阅
async def subscribe_tendermint(
query: str,
ws_url: str = TENDERMINT_WS,
) -> AsyncGenerator[dict, None]:
async with websockets.connect(ws_url) as ws:
await ws.send(json.dumps({
'jsonrpc': '2.0',
'method': 'subscribe',
'params': {'query': query},
'id': 1,
}))
response = json.loads(await ws.recv())
if 'error' in response:
raise ConnectionError(f"Subscription failed: {response['error']}")
async for message in ws:
event = json.loads(message)
if event.get('id') is not None:
continue
if event.get('result', {}).get('type') == 'heartbeat':
continue
yield event
3.3 Agent API WebSocket 订阅(推荐)
async def subscribe_agent_events(
topics: list[str],
agent_id: Optional[str] = None,
ws_url: str = AGENT_WS,
) -> AsyncGenerator[dict, None]:
async with websockets.connect(ws_url) as ws:
request = {
'jsonrpc': '2.0',
'method': 'subscribe',
'params': {
'topics': topics,
'chain_id': 'msg-chain-1',
},
'id': 1,
}
if agent_id:
request['params']['agent_id'] = agent_id
await ws.send(json.dumps(request))
ack = json.loads(await ws.recv())
if 'error' in ack:
raise ConnectionError(
f"Agent event subscription failed: {ack['error']}"
)
async for message in ws:
yield json.loads(message)
3.4 多主题订阅管理
class SubscriptionManager:
def __init__(self, ws_url: str = AGENT_WS):
self.ws_url = ws_url
self.subscriptions: dict[str, dict] = {}
self._ws = None
async def connect(self):
self._ws = await websockets.connect(self.ws_url)
async def disconnect(self):
if self._ws:
await self._ws.close()
async def add_subscription(
self, sub_id: str, topics: list[str], agent_id: Optional[str] = None,
):
request = {
'jsonrpc': '2.0',
'method': 'subscribe',
'params': {
'topics': topics,
'chain_id': 'msg-chain-1',
'sub_id': sub_id,
},
'id': hash(sub_id) % (2**31),
}
if agent_id:
request['params']['agent_id'] = agent_id
await self._ws.send(json.dumps(request))
self.subscriptions[sub_id] = {'topics': topics, 'active': True}
async def remove_subscription(self, sub_id: str):
await self._ws.send(json.dumps({
'jsonrpc': '2.0',
'method': 'unsubscribe',
'params': {'sub_id': sub_id},
'id': hash(sub_id) % (2**31),
}))
self.subscriptions.pop(sub_id, None)
async def event_stream(self) -> AsyncGenerator[tuple[str, dict], None]:
async for message in self._ws:
data = json.loads(message)
sub_id = data.get('params', {}).get('sub_id', 'unknown')
yield sub_id, data
3.5 心跳与连接保活
import logging
logger = logging.getLogger(__name__)
class HeartbeatManager:
def __init__(self, ws, interval: float = 30.0, timeout: float = 10.0):
self.ws = ws
self.interval = interval
self.timeout = timeout
self._running = False
async def start(self):
self._running = True
while self._running:
await asyncio.sleep(self.interval)
try:
pong_waiter = await self.ws.ping()
await asyncio.wait_for(pong_waiter, timeout=self.timeout)
except (asyncio.TimeoutError, websockets.ConnectionClosed):
raise ConnectionError('WebSocket heartbeat failed')
def stop(self):
self._running = False
3.6 事件缓冲与背压
from collections import deque
from dataclasses import dataclass, field
from datetime import datetime
@dataclass
class BufferedEvent:
data: dict
received_at: datetime = field(default_factory=datetime.utcnow)
class EventBuffer:
def __init__(self, max_size: int = 10000, batch_timeout: float = 0.1):
self.queue: deque[BufferedEvent] = deque(maxlen=max_size)
self.max_size = max_size
self.batch_timeout = batch_timeout
self._new_event = asyncio.Event()
async def push(self, event: dict):
if len(self.queue) >= self.max_size:
logger.warning('Event buffer full, discarding oldest')
self.queue.popleft()
self.queue.append(BufferedEvent(data=event))
self._new_event.set()
async def pull_batch(self) -> list[dict]:
try:
await asyncio.wait_for(self._new_event.wait(), timeout=self.batch_timeout)
except asyncio.TimeoutError:
return []
self._new_event.clear()
batch = []
while self.queue:
batch.append(self.queue.popleft().data)
return batch
@property
def size(self) -> int:
return len(self.queue)
3.7 完整 WebSocket 客户端
class MsgChainWSClient:
def __init__(
self,
topics: list[str],
agent_id: Optional[str] = None,
ws_url: str = AGENT_WS,
heartbeat_interval: float = 30.0,
buffer_max_size: int = 10000,
reconnect_delay: float = 1.0,
max_reconnect_delay: float = 60.0,
):
self.topics = topics
self.agent_id = agent_id
self.ws_url = ws_url
self.heartbeat_interval = heartbeat_interval
self.reconnect_delay = reconnect_delay
self.max_reconnect_delay = max_reconnect_delay
self.buffer = EventBuffer(max_size=buffer_max_size)
self._running = False
self._event_handlers: list[callable] = []
self._ws = None
self._heartbeat = None
def on_event(self, handler: callable):
self._event_handlers.append(handler)
return handler
async def _connect_and_subscribe(self):
self._ws = await websockets.connect(
self.ws_url, ping_interval=None, ping_timeout=None, max_size=2**20,
)
self._heartbeat = HeartbeatManager(self._ws, self.heartbeat_interval)
asyncio.create_task(self._heartbeat.start())
await self._ws.send(json.dumps({
'jsonrpc': '2.0', 'method': 'subscribe',
'params': {
'topics': self.topics,
'chain_id': 'msg-chain-1',
'agent_id': self.agent_id,
} if self.agent_id else {
'topics': self.topics,
'chain_id': 'msg-chain-1',
},
'id': 1,
}))
ack = json.loads(await self._ws.recv())
if 'error' in ack:
raise ConnectionError(f"Subscription failed: {ack['error']}")
async def _message_loop(self):
while self._running:
try:
message = await self._ws.recv()
event = json.loads(message)
if event.get('id') is not None:
continue
await self.buffer.push(event)
for handler in self._event_handlers:
try:
if asyncio.iscoroutinefunction(handler):
await handler(event)
else:
handler(event)
except Exception:
logger.exception('Event handler failed')
except websockets.ConnectionClosed:
raise
async def _reconnect(self):
delay = self.reconnect_delay
while self._running:
try:
logger.info(f'Reconnecting in {delay:.1f}s...')
await asyncio.sleep(delay)
await self._connect_and_subscribe()
return
except (ConnectionError, websockets.WebSocketException) as e:
logger.error(f'Reconnect failed: {e}')
delay = min(delay * 2, self.max_reconnect_delay)
async def run(self):
self._running = True
while self._running:
try:
await self._connect_and_subscribe()
await self._message_loop()
except (ConnectionError, websockets.WebSocketException):
if self._running:
await self._reconnect()
except asyncio.CancelledError:
break
async def stop(self):
self._running = False
if self._heartbeat:
self._heartbeat.stop()
if self._ws:
await self._ws.close()
async def example_usage():
client = MsgChainWSClient(
topics=['tx', 'wasm', 'agent'],
agent_id='agent:msg1myagent...',
)
@client.on_event
async def handle_payment(event: dict):
print(f'Event received: {event.get("tx_hash")}')
try:
await client.run()
except KeyboardInterrupt:
await client.stop()
4. 事件历史查询
4.1 客户端基础
import httpx
from typing import Optional, AsyncGenerator
HISTORY_URL = 'https://api.msgchain.org/agent/v1/events/history'
class EventsHistoryClient:
def __init__(
self,
base_url: str = 'https://api.msgchain.org',
api_key: Optional[str] = None,
timeout: float = 30.0,
):
self.base_url = base_url.rstrip('/')
self._client = httpx.AsyncClient(
base_url=self.base_url,
timeout=timeout,
headers={'Content-Type': 'application/json',
'Accept': 'application/json',
**( {'X-API-Key': api_key} if api_key else {} )},
)
async def close(self):
await self._client.aclose()
4.2 基本查询与过滤
async def get_events(
self,
event_type: Optional[str] = None,
from_height: Optional[int] = None,
to_height: Optional[int] = None,
from_time: Optional[str] = None,
to_time: Optional[str] = None,
agent_id: Optional[str] = None,
tx_hash: Optional[str] = None,
limit: int = 100,
offset: int = 0,
) -> dict:
params = {'limit': min(limit, 1000), 'offset': offset}
if event_type: params['event_type'] = event_type
if from_height is not None: params['from_height'] = from_height
if to_height is not None: params['to_height'] = to_height
if from_time: params['from_time'] = from_time
if to_time: params['to_time'] = to_time
if agent_id: params['agent_id'] = agent_id
if tx_hash: params['tx_hash'] = tx_hash
response = await self._client.get('/agent/v1/events/history', params=params)
response.raise_for_status()
return response.json()
async def query_examples():
client = EventsHistoryClient()
agent_events = await client.get_events(event_type='agent', limit=20)
transfers = await client.get_events(
event_type='transfer', from_height=100000, to_height=101000, limit=100,
)
tx_events = await client.get_events(tx_hash='0xABCDEF1234567890...')
time_filtered = await client.get_events(
from_time='2026-07-08T00:00:00Z', to_time='2026-07-08T23:59:59Z',
)
4.3 分页遍历
async def get_events_paginated(self, **kwargs) -> AsyncGenerator[dict, None]:
page_size = min(kwargs.pop('limit', 100), 100)
offset = 0
while True:
result = await self.get_events(**kwargs, limit=page_size, offset=offset)
items = result.get('events', [])
if not items:
break
for event in items:
yield event
total = result.get('pagination', {}).get('total', 0)
offset += page_size
if offset >= total:
break
async def get_events_by_height_range(
self, from_height: int, to_height: int, event_type: Optional[str] = None,
) -> list[dict]:
all_events = []
current = from_height
while current <= to_height:
batch_end = min(current + 999, to_height)
async for e in self.get_events_paginated(
event_type=event_type, from_height=current, to_height=batch_end, limit=100,
):
all_events.append(e)
current = batch_end + 1
return all_events
4.4 事件重放
class EventReplayManager:
def __init__(self, history_client: EventsHistoryClient, state_store: Optional[dict] = None):
self.client = history_client
self.state = state_store or {}
self._processed: set[str] = set()
def _key(self, event: dict) -> str:
return f"{event.get('tx_hash', '')}:{event.get('event_index', 0)}"
def mark_processed(self, event: dict):
self._processed.add(self._key(event))
def is_processed(self, event: dict) -> bool:
return self._key(event) in self._processed
async def replay_missed(
self, last_height: int, current_height: int, event_type: Optional[str] = None,
) -> list[dict]:
if current_height <= last_height:
return []
missed = []
async for e in self.client.get_events_paginated(
event_type=event_type, from_height=last_height + 1, to_height=current_height, limit=100,
):
if not self.is_processed(e):
missed.append(e)
return missed
def save_checkpoint(self, height: int):
self.state['last_processed_height'] = height
def load_checkpoint(self) -> Optional[int]:
return self.state.get('last_processed_height')
async def replay_example():
client = EventsHistoryClient()
state = {'last_processed_height': 99999}
replay = EventReplayManager(client, state)
current_height = 100500
last = replay.load_checkpoint()
if last and current_height > last:
missed = await replay.replay_missed(last, current_height, 'agent')
for event in missed:
await process_event(event)
replay.mark_processed(event)
replay.save_checkpoint(current_height)
4.5 查询结果聚合
from collections import Counter, defaultdict
class EventAnalytics:
def __init__(self, events: list[dict]):
self.events = events
def count_by_type(self) -> Counter:
return Counter(e.get('event_type', 'unknown') for e in self.events)
def count_by_action(self) -> Counter:
actions = Counter()
for e in self.events:
actions[e.get('attributes', {}).get('action', 'unknown')] += 1
return actions
def volume_by_hour(self) -> dict[str, int]:
vol = defaultdict(int)
for e in self.events:
ts = e.get('timestamp', '')
if ts:
vol[ts[:13] + ':00:00'] += 1
return dict(vol)
def summary(self) -> dict:
heights = [e.get('height', 0) for e in self.events]
return {
'total': len(self.events),
'by_type': self.count_by_type(),
'by_action': self.count_by_action(),
'unique_tx': len(set(e.get('tx_hash', '') for e in self.events)),
'height_range': (min(heights), max(heights)) if heights else (0, 0),
}
async def analytics_example():
client = EventsHistoryClient()
data = await client.get_events(
from_time='2026-07-08T00:00:00Z', to_time='2026-07-08T23:59:59Z', limit=5000,
)
analytics = EventAnalytics(data.get('events', []))
print(analytics.summary())
5. 事件驱动 Agent 模式
5.1 Agent 基类
from abc import ABC, abstractmethod
from typing import Optional
class EventDrivenAgent(ABC):
def __init__(
self,
agent_id: str,
topics: list[str],
ws_client: Optional[MsgChainWSClient] = None,
replay_manager: Optional[EventReplayManager] = None,
):
self.agent_id = agent_id
self.topics = topics
self.ws_client = ws_client or MsgChainWSClient(topics=topics, agent_id=agent_id)
self.replay_manager = replay_manager
self._running = False
@abstractmethod
async def should_respond(self, event: dict) -> bool:
...
@abstractmethod
async def decide_action(self, event: dict) -> Optional[dict]:
...
@abstractmethod
async def execute_action(self, action: dict, context: Optional[dict] = None):
...
async def handle_event(self, event: dict):
try:
if not await self.should_respond(event):
return
action = await self.decide_action(event)
if action is None:
return
await self.execute_action(action, {
'event': event,
'agent_id': self.agent_id,
'timestamp': datetime.utcnow().isoformat(),
})
if self.replay_manager:
self.replay_manager.mark_processed(event)
except Exception:
logger.exception(f'Handle event failed: {event.get("tx_hash")}')
async def run(self):
self._running = True
if self.replay_manager:
last = self.replay_manager.load_checkpoint()
if last:
logger.info(f'Resuming from height {last}')
@self.ws_client.on_event
async def on_event(event: dict):
await self.handle_event(event)
await self.ws_client.run()
async def stop(self):
self._running = False
await self.ws_client.stop()
5.2 支付监控 Agent
class PaymentMonitorAgent(EventDrivenAgent):
def __init__(self, agent_id: str, expected_price: int = 1000000, price_denom: str = 'umsg'):
super().__init__(agent_id=agent_id, topics=['tx', 'agent'])
self.expected_price = expected_price
self.price_denom = price_denom
self.subscribers: dict[str, dict] = {}
self.payment_history: list[dict] = []
async def should_respond(self, event: dict) -> bool:
return MsgChainEvent(event).has_event('agent')
async def decide_action(self, event: dict) -> Optional[dict]:
parsed = MsgChainEvent(event)
attrs = parsed.get_attributes('agent')
if not attrs:
return None
action = attrs.get('action')
if action == 'payment' and attrs.get('agent_id') == self.agent_id:
amount = int(attrs.get('amount', '0').replace(self.price_denom, ''))
return {
'type': 'process_payment',
'from': attrs.get('from'),
'amount': amount,
'tx_hash': parsed.tx_hash,
}
if action == 'subscribe' and attrs.get('agent_id') == self.agent_id:
return {'type': 'add_subscriber', 'subscriber': attrs.get('from')}
if action == 'unsubscribe' and attrs.get('agent_id') == self.agent_id:
return {'type': 'remove_subscriber', 'subscriber': attrs.get('from')}
return None
async def execute_action(self, action: dict, context: Optional[dict] = None):
if action['type'] == 'process_payment':
self.payment_history.append({
'from': action['from'], 'amount': action['amount'],
'tx_hash': action['tx_hash'],
'at': datetime.utcnow().isoformat(),
})
if action['amount'] >= self.expected_price:
self.subscribers.setdefault(action['from'], {
'status': 'active',
'activated_at': datetime.utcnow().isoformat(),
'total_paid': 0,
})['total_paid'] += action['amount']
logger.info(f'Payment: {action["from"]} paid {action["amount"]}')
elif action['type'] == 'add_subscriber':
self.subscribers[action['subscriber']] = {
'status': 'active', 'activated_at': datetime.utcnow().isoformat(), 'total_paid': 0,
}
elif action['type'] == 'remove_subscriber':
s = self.subscribers.get(action['subscriber'])
if s: s['status'] = 'cancelled'
5.3 合约监控 Agent
class ContractMonitorAgent(EventDrivenAgent):
def __init__(self, agent_id: str, contract_address: str, watch_actions: Optional[list[str]] = None):
super().__init__(agent_id=agent_id, topics=['tx', 'wasm'])
self.contract_address = contract_address
self.watch_actions = watch_actions or ['*']
self.event_log: list[dict] = []
self.thresholds: dict[str, int] = {}
def set_threshold(self, action: str, value: int):
self.thresholds[action] = value
async def should_respond(self, event: dict) -> bool:
parsed = MsgChainEvent(event)
if not parsed.has_event('wasm'):
return False
attrs = parsed.get_attributes('wasm')
return bool(attrs and attrs.get('_contract_address') == self.contract_address and
('*' in self.watch_actions or attrs.get('action') in self.watch_actions))
async def decide_action(self, event: dict) -> Optional[dict]:
parsed = MsgChainEvent(event)
attrs = parsed.get_attributes('wasm')
action = attrs.get('action', '')
if action in self.thresholds:
curr = int(attrs.get('count', '0'))
if curr >= self.thresholds[action]:
return {'type': 'threshold_alert', 'action': action,
'current_value': curr, 'threshold': self.thresholds[action]}
return {'type': 'log_event', 'action': action, 'attributes': attrs, 'tx_hash': parsed.tx_hash}
async def execute_action(self, action: dict, context: Optional[dict] = None):
if action['type'] == 'log_event':
self.event_log.append(action)
elif action['type'] == 'threshold_alert':
logger.warning(f"THRESHOLD: {action['action']} at {action['current_value']}")
5.4 治理投票 Agent
class GovernanceAgent(EventDrivenAgent):
def __init__(self, agent_id: str, voter: str, strategy: Optional[dict] = None):
super().__init__(agent_id=agent_id, topics=['tx', 'governance'])
self.voter = voter
self.strategy = strategy or {
'default': 'VOTE_OPTION_ABSTAIN', 'Text': 'VOTE_OPTION_YES',
}
self.proposals: dict[str, dict] = {}
self.vote_history: list[dict] = []
async def should_respond(self, event: dict) -> bool:
p = MsgChainEvent(event)
return any(p.has_event(t) for t in ('submit_proposal', 'proposal_result', 'vote'))
async def decide_action(self, event: dict) -> Optional[dict]:
p = MsgChainEvent(event)
if p.has_event('submit_proposal'):
a = p.get_attributes('submit_proposal')
self.proposals[a['proposal_id']] = {'id': a['proposal_id'], 'type': a.get('proposal_type', 'Text'), 'status': 'pending'}
vote = self.strategy.get(a.get('proposal_type', ''), self.strategy['default'])
if vote != 'VOTE_OPTION_ABSTAIN':
return {'type': 'vote', 'proposal_id': a['proposal_id'], 'vote': vote}
if p.has_event('proposal_result'):
a = p.get_attributes('proposal_result')
pid = a['proposal_id']
if pid in self.proposals:
self.proposals[pid]['result'] = a['result']
self.proposals[pid]['status'] = 'completed'
return {'type': 'record_result', 'proposal_id': pid, 'result': a['result']}
return None
async def execute_action(self, action: dict, context: Optional[dict] = None):
if action['type'] == 'vote':
self.vote_history.append(action)
logger.info(f"Voted {action['vote']} on #{action['proposal_id']}")
elif action['type'] == 'record_result':
logger.info(f"Proposal #{action['proposal_id']}: {action['result']}")
5.5 A2A 消息 Agent
class A2AMessagingAgent(EventDrivenAgent):
def __init__(self, agent_id: str):
super().__init__(agent_id=agent_id, topics=[f"agent.action='a2a' AND agent.to_agent='{agent_id}'"])
self.inbox: list[dict] = []
self.sent: list[dict] = []
self.handlers: dict[str, callable] = {}
def register_handler(self, msg_type: str, handler: callable):
self.handlers[msg_type] = handler
async def should_respond(self, event: dict) -> bool:
p = MsgChainEvent(event)
a = p.get_attributes('agent')
return bool(a and a.get('action') == 'a2a' and a.get('to_agent') == self.agent_id)
async def decide_action(self, event: dict) -> Optional[dict]:
p = MsgChainEvent(event)
a = p.get_attributes('agent')
msg = {'message_id': a.get('message_id'), 'from': a.get('from_agent'),
'type': a.get('message_type'), 'tx_hash': p.tx_hash}
self.inbox.append(msg)
return {'type': 'process_message', 'message': msg}
async def execute_action(self, action: dict, context: Optional[dict] = None):
msg = action['message']
handler = self.handlers.get(msg['type'])
if handler:
try: await handler(msg)
except Exception: logger.exception(f'Handler failed: {msg["message_id"]}')
5.6 多 Agent 协调器
class AgentCoordinator:
def __init__(self):
self.agents: dict[str, EventDrivenAgent] = {}
self._tasks: dict[str, asyncio.Task] = {}
def register(self, agent_id: str, agent: EventDrivenAgent):
self.agents[agent_id] = agent
async def start_all(self):
for aid, agent in self.agents.items():
self._tasks[aid] = asyncio.create_task(agent.run(), name=aid)
async def stop_all(self):
for agent in self.agents.values():
await agent.stop()
for t in self._tasks.values():
t.cancel()
await asyncio.gather(*self._tasks.values(), return_exceptions=True)
async def multi_agent_example():
coord = AgentCoordinator()
coord.register('payment', PaymentMonitorAgent('agent:msg1payment...', expected_price=5000000))
coord.register('contract', ContractMonitorAgent('agent:msg1contract...', 'msg1monitored...'))
try:
await coord.start_all()
await asyncio.Future()
except KeyboardInterrupt:
await coord.stop_all()
6. 触发器与自动化工作流
6.1 事件-条件-动作流水线
class Condition:
async def evaluate(self, event: dict, context: dict) -> bool:
raise NotImplementedError
class Action:
async def execute(self, event: dict, context: dict):
raise NotImplementedError
class EventTrigger:
def __init__(self, name: str, condition: Condition, action: Action, desc: str = ''):
self.name = name
self.condition = condition
self.action = action
self.description = desc
self.execution_count = 0
self.last_execution: Optional[datetime] = None
async def evaluate_and_execute(self, event: dict, context: dict) -> bool:
if await self.condition.evaluate(event, context):
await self.action.execute(event, context)
self.execution_count += 1
self.last_execution = datetime.utcnow()
return True
return False
class TriggerPipeline:
def __init__(self):
self.triggers: list[EventTrigger] = []
def add(self, trigger: EventTrigger):
self.triggers.append(trigger)
async def process(self, event: dict, context: dict) -> list[str]:
executed = []
for t in self.triggers:
if await t.evaluate_and_execute(event, context):
executed.append(t.name)
return executed
6.2 常用条件实现
class EventTypeCondition(Condition):
def __init__(self, types: list[str]):
self.types = set(types)
async def evaluate(self, event: dict, context: dict) -> bool:
p = MsgChainEvent(event)
return bool(self.types & set(p.event_types))
class AttributeCondition(Condition):
def __init__(self, event_type: str, key: str, value: str):
self.event_type = event_type
self.key = key
self.value = value
async def evaluate(self, event: dict, context: dict) -> bool:
p = MsgChainEvent(event)
a = p.get_attributes(self.event_type)
return bool(a and a.get(self.key) == self.value)
class ThresholdCondition(Condition):
def __init__(self, event_type: str, key: str, threshold: int, op: str = '>='):
self.event_type = event_type
self.key = key
self.threshold = threshold
self.op = op
async def evaluate(self, event: dict, context: dict) -> bool:
p = MsgChainEvent(event)
a = p.get_attributes(self.event_type)
if not a: return False
val_str = a.get(self.key, '0')
for prefix in ['umsg', 'msg', 'u']:
if val_str.endswith(prefix):
val_str = val_str.replace(prefix, '')
value = int(val_str) if val_str.isdigit() else 0
ops = {'>=': lambda v, t: v >= t, '>': lambda v, t: v > t,
'<=': lambda v, t: v <= t, '<': lambda v, t: v < t,
'==': lambda v, t: v == t}
return ops.get(self.op, lambda v, t: False)(value, self.threshold)
class AndCondition(Condition):
def __init__(self, *conditions: Condition):
self.conditions = conditions
async def evaluate(self, event: dict, context: dict) -> bool:
for c in self.conditions:
if not await c.evaluate(event, context):
return False
return True
class OrCondition(Condition):
def __init__(self, *conditions: Condition):
self.conditions = conditions
async def evaluate(self, event: dict, context: dict) -> bool:
for c in self.conditions:
if await c.evaluate(event, context):
return True
return False
6.3 常用动作实现
class LogAction(Action):
def __init__(self, level: str = 'info'):
self.level = level
async def execute(self, event: dict, context: dict):
msg = f'Trigger: {event.get("tx_hash", "?")}'
if self.level == 'info': logger.info(msg)
elif self.level == 'warning': logger.warning(msg)
class WebhookAction(Action):
def __init__(self, url: str, headers: Optional[dict] = None):
self.url = url
self.headers = headers or {}
self._client = httpx.AsyncClient()
async def execute(self, event: dict, context: dict):
try:
r = await self._client.post(self.url, json={'event': event, 'context': context}, headers=self.headers, timeout=10.0)
r.raise_for_status()
except httpx.HTTPError as e:
logger.error(f'Webhook failed: {e}')
async def close(self):
await self._client.aclose()
class NotificationAction(Action):
def __init__(self, title_template: str, body_template: str, channels: list[str] | None = None):
self.title_tmpl = title_template
self.body_tmpl = body_template
self.channels = channels or ['log']
async def execute(self, event: dict, context: dict):
title = self.title_tmpl.format(event=event, context=context)
body = self.body_tmpl.format(event=event, context=context)
if 'log' in self.channels:
logger.info(f'[{title}] {body}')
6.4 Cron 时间触发器
import croniter
class CronTrigger:
def __init__(self, name: str, cron: str, action: Action, desc: str = ''):
self.name = name
self.cron = cron
self.action = action
self.description = desc
self._last_run: Optional[datetime] = None
async def should_run(self) -> bool:
base = self._last_run or datetime.utcnow()
return datetime.utcnow() >= croniter.croniter(self.cron, base).get_next(datetime)
async def execute(self, context: Optional[dict] = None):
await self.action.execute({'type': 'cron', 'trigger': self.name}, context or {})
self._last_run = datetime.utcnow()
class CronScheduler:
def __init__(self):
self.triggers: list[CronTrigger] = []
self._running = False
def add(self, trigger: CronTrigger):
self.triggers.append(trigger)
async def start(self):
self._running = True
asyncio.create_task(self._loop())
async def _loop(self):
while self._running:
for t in self.triggers:
if await t.should_run():
await t.execute()
await asyncio.sleep(1)
async def stop(self):
self._running = False
def setup_cron() -> CronScheduler:
s = CronScheduler()
s.add(CronTrigger('hourly_summary', '0 * * * *', LogAction()))
s.add(CronTrigger('daily_cleanup', '0 0 * * *', LogAction()))
return s
6.5 事件聚合
import time
from collections import defaultdict
class TimeWindowAggregator:
def __init__(self, window: float = 60.0, max_events: int = 100):
self.window = window
self.max_events = max_events
self.buckets: dict[str, list[dict]] = defaultdict(list)
self._callbacks: list[callable] = []
def on_flush(self, cb: callable):
self._callbacks.append(cb)
async def add(self, key: str, event: dict):
self.buckets[key].append({'event': event, 'ts': time.time()})
if len(self.buckets[key]) >= self.max_events:
await self._flush(key)
async def _flush(self, key: str):
batch = self.buckets.pop(key, [])
for cb in self._callbacks:
try:
if asyncio.iscoroutinefunction(cb):
await cb(key, batch)
else:
cb(key, batch)
except Exception:
logger.exception(f'Flush callback failed: {key}')
async def flush_expired(self):
now = time.time()
for key in list(self.buckets.keys()):
expired = [e for e in self.buckets[key] if e['ts'] < now - self.window]
if len(expired) == len(self.buckets[key]):
await self._flush(key)
async def start_auto_flush(self, interval: float = 10.0):
while True:
await asyncio.sleep(interval)
await self.flush_expired()
6.6 事件关联
class PatternStep:
def __init__(self, event_type: str, conditions: Optional[list[Condition]] = None):
self.event_type = event_type
self.conditions = conditions or []
async def matches(self, event: dict) -> bool:
p = MsgChainEvent(event)
if not p.has_event(self.event_type):
return False
for c in self.conditions:
if not await c.evaluate(event, {}):
return False
return True
class EventPattern:
def __init__(self, name: str, steps: list[PatternStep], timeout: float = 60.0):
self.name = name
self.steps = steps
self.timeout = timeout
async def on_match(self, events: list[dict]):
logger.info(f'Pattern matched: {self.name}')
class CorrelationEngine:
def __init__(self):
self.patterns: list[EventPattern] = {}
self._sessions: dict[str, '_Session'] = {}
def add(self, pattern: EventPattern):
self.patterns[pattern.name] = pattern
async def feed(self, event: dict):
for p in self.patterns.values():
if p.name not in self._sessions:
self._sessions[p.name] = _Session(p)
if await self._sessions[p.name].feed(event):
if self._sessions[p.name].is_complete:
await p.on_match(self._sessions[p.name].events)
del self._sessions[p.name]
class _Session:
def __init__(self, pattern: EventPattern):
self.pattern = pattern
self.events: list[dict] = []
self._step = 0
self._created = time.time()
@property
def is_complete(self) -> bool:
return self._step >= len(self.pattern.steps)
@property
def expired(self) -> bool:
return (time.time() - self._created) > self.pattern.timeout
async def feed(self, event: dict) -> bool:
if self.is_complete or self.expired:
return False
if await self.pattern.steps[self._step].matches(event):
self.events.append(event)
self._step += 1
return True
return False
def setup_correlation():
engine = CorrelationEngine()
pattern = EventPattern('large_tx_and_stake', [
PatternStep('transfer', [ThresholdCondition('transfer', 'amount', 1000000000)]),
PatternStep('delegate'),
], timeout=300.0)
engine.add(pattern)
return engine
6.7 完整工作流编排
class AutomatedWorkflow:
def __init__(self):
self.pipeline = TriggerPipeline()
self.cron = CronScheduler()
self.aggregator = TimeWindowAggregator()
self.correlation = CorrelationEngine()
self._running = False
def add_trigger(self, trigger: EventTrigger):
self.pipeline.add(trigger)
def add_cron(self, trigger: CronTrigger):
self.cron.add(trigger)
async def handle_event(self, event: dict):
await self.pipeline.process(event, {'ts': datetime.utcnow().isoformat()})
async def run(self):
self._running = True
asyncio.create_task(self.cron.start())
client = MsgChainWSClient(topics=['tx', 'wasm', 'transfer', 'agent'])
@client.on_event
async def on_event(event: dict):
await self.handle_event(event)
try:
await client.run()
finally:
self._running = False
await self.cron.stop()
async def workflow_example():
wf = AutomatedWorkflow()
wf.add_trigger(EventTrigger(
'large_transfer',
AndCondition(
EventTypeCondition(['transfer']),
ThresholdCondition('transfer', 'amount', 10000000000),
),
NotificationAction('Large Transfer', '{event.get("tx_hash")}', ['log']),
))
wf.add_cron(CronTrigger('hourly', '0 * * * *', LogAction()))
await wf.run()
7. 边界与最佳实践
7.1 WebSocket 重连策略
class ReconnectStrategy:
def __init__(self, initial: float = 1.0, max_delay: float = 60.0, factor: float = 2.0, jitter: float = 0.1, max_retries: Optional[int] = None):
self.initial = initial
self.max_delay = max_delay
self.factor = factor
self.jitter = jitter
self.max_retries = max_retries
self._attempts = 0
def get_delay(self) -> float:
self._attempts += 1
d = min(self.initial * (self.factor ** (self._attempts - 1)), self.max_delay)
import random
return max(0.1, d + random.uniform(-d * self.jitter, d * self.jitter))
def should_retry(self) -> bool:
return self.max_retries is None or self._attempts < self.max_retries
def reset(self):
self._attempts = 0
class ResilientClient(MsgChainWSClient):
def __init__(self, *args, strategy: Optional[ReconnectStrategy] = None, **kwargs):
super().__init__(*args, **kwargs)
self._strategy = strategy or ReconnectStrategy()
async def _connect_with_retry(self) -> bool:
while self._running and self._strategy.should_retry():
try:
await self._connect_and_subscribe()
self._strategy.reset()
return True
except (ConnectionError, websockets.WebSocketException) as e:
d = self._strategy.get_delay()
logger.warning(f'Retry {self._strategy._attempts} in {d:.1f}s: {e}')
await asyncio.sleep(d)
logger.error('Max retries exceeded')
return False
7.2 事件去重
class EventDeduplicator:
def __init__(self, capacity: int = 100000):
self.capacity = capacity
self._seen: set[str] = set()
def key(self, event: dict) -> str:
return f"{event.get('tx_hash', '')}:{event.get('event_index', 0)}"
def is_duplicate(self, event: dict) -> bool:
return self.key(event) in self._seen
def mark(self, event: dict):
self._seen.add(self.key(event))
if len(self._seen) > self.capacity * 2:
self._seen = set(list(self._seen)[-self.capacity:])
def check(self, event: dict) -> bool:
if self.is_duplicate(event):
return True
self.mark(event)
return False
class DedupedAgent(EventDrivenAgent):
def __init__(self, *args, dedup_capacity: int = 100000, **kwargs):
super().__init__(*args, **kwargs)
self.dedup = EventDeduplicator(capacity=dedup_capacity)
async def handle_event(self, event: dict):
if self.dedup.check(event):
return
await super().handle_event(event)
7.3 丢失事件恢复
class EventRecovery:
def __init__(self, history: EventsHistoryClient, checkpoint_file: str = '/tmp/msgchain.json'):
self.history = history
self.checkpoint_file = checkpoint_file
self.last_height: Optional[int] = None
async def recover(self):
last = await self._load()
if last:
logger.info(f'Recovering from height {last + 1}')
current = await self._current_height()
if current > last:
async for e in self.history.get_events_paginated(from_height=last + 1, to_height=current, limit=100):
yield e
async def save_checkpoint(self, height: int):
self.last_height = height
import json
with open(self.checkpoint_file, 'w') as f:
json.dump({'height': height}, f)
async def _load(self) -> Optional[int]:
import json, os
if os.path.exists(self.checkpoint_file):
with open(self.checkpoint_file) as f:
return json.load(f).get('height')
return None
async def _current_height(self) -> int:
r = await self.history.get_events(limit=1)
events = r.get('events', [])
return int(events[0]['height']) if events else 0
7.4 速率限制
from collections import deque
import time
class RateLimiter:
def __init__(self, max_calls: int, period: float = 1.0):
self.max_calls = max_calls
self.period = period
self.calls: deque[float] = deque()
async def acquire(self):
now = time.time()
while self.calls and self.calls[0] < now - self.period:
self.calls.popleft()
if len(self.calls) >= self.max_calls:
wait = self.calls[0] + self.period - now
if wait > 0:
await asyncio.sleep(wait)
now = time.time()
while self.calls and self.calls[0] < now - self.period:
self.calls.popleft()
self.calls.append(now)
class Throttle:
def __init__(self, max_per_second: float = 10.0):
self.min_interval = 1.0 / max_per_second
self._last = 0.0
async def wait(self):
now = time.time()
elapsed = now - self._last
if elapsed < self.min_interval:
await asyncio.sleep(self.min_interval - elapsed)
self._last = now
class EventProcessor:
def __init__(self, buffer: EventBuffer, max_concurrent: int = 10, rate: float = 50.0):
self.buffer = buffer
self.semaphore = asyncio.Semaphore(max_concurrent)
self.throttle = Throttle(rate)
self.processed = 0
self.errors = 0
async def run(self):
while True:
batch = await self.buffer.pull_batch()
if not batch:
continue
async with self.semaphore:
await self.throttle.wait()
for event in batch:
try:
await self.process(event)
self.processed += 1
except Exception:
self.errors += 1
async def process(self, event: dict):
raise NotImplementedError
7.5 回退与断路器
class RetryHandler:
def __init__(self, max_retries: int = 3, base: float = 0.5, max_delay: float = 30.0,
retryable: Optional[tuple] = None):
self.max_retries = max_retries
self.base = base
self.max_delay = max_delay
self.retryable = retryable or (ConnectionError, TimeoutError, httpx.HTTPError)
async def execute(self, func: callable, *args, **kwargs):
last = None
for attempt in range(self.max_retries + 1):
try:
return await func(*args, **kwargs)
except self.retryable as e:
last = e
if attempt < self.max_retries:
d = min(self.base * (2 ** attempt), self.max_delay)
await asyncio.sleep(d)
raise last
class CircuitBreaker:
def __init__(self, threshold: int = 5, timeout: float = 30.0, half_open_max: int = 1):
self.threshold = threshold
self.timeout = timeout
self.half_open_max = half_open_max
self.state = 'closed'
self.failures = 0
self.last_failure: Optional[float] = None
self.half_open_requests = 0
async def call(self, func: callable, *args, **kwargs):
if self.state == 'open':
if time.time() - self.last_failure >= self.timeout:
self.state = 'half-open'
self.half_open_requests = 0
else:
raise Exception('Circuit breaker open')
if self.state == 'half-open' and self.half_open_requests >= self.half_open_max:
raise Exception('Circuit breaker half-open')
self.half_open_requests += 1
try:
result = await func(*args, **kwargs)
if self.state == 'half-open':
self.state = 'closed'
self.failures = 0
return result
except Exception as e:
self.failures += 1
self.last_failure = time.time()
if self.failures >= self.threshold:
self.state = 'open'
raise e
7.6 安全最佳实践
import re
import hashlib
import hmac
def validate_msg_address(addr: str) -> bool:
return bool(re.match(r'^msg[a-z0-9]{10,64}$', addr))
def sanitize_event(event: dict) -> dict:
safe = event.copy()
sensitive = {'private_key', 'seed', 'mnemonic', 'password'}
for ev in safe.get('events', []):
for attr in ev.get('attributes', []):
if attr.get('key', '').lower() in sensitive:
attr['value'] = '***REDACTED***'
return safe
def verify_signature(event: dict, signature: str, secret: str) -> bool:
payload = json.dumps(event, sort_keys=True).encode()
expected = hmac.new(secret.encode(), payload, hashlib.sha256).hexdigest()
return hmac.compare_digest(expected, signature)
7.7 生产就绪性差距
| 项 | 当前实现 | 生产要求 |
|---|---|---|
| 持久化 | 内存存储 | PostgreSQL/Redis |
| 认证 | API Key | JWT/OAuth2 + 密钥管理 |
| 监控 | 仅日志 | Prometheus + Grafana + 告警 |
| 消息可靠性 | 至多一次 | 至少一次 + 死信队列 |
| 配置 | 硬编码 | 环境变量 + 配置中心 |
| 部署 | 单进程 | 容器化 + 编排 + 健康检查 |
STUB 端点说明:
STUB_ENDPOINTS = {
'/agent/v1/events/subscribe': 'Beta 阶段',
'/agent/v1/events/history': 'Beta 阶段,数据可能有延迟',
}
STUB_WARNING = """
Agent 专属端点处于 Beta 阶段。建议:
1. 使用 wss://rpc.msgchain.org/websocket 为主事件源
2. 将 /agent/v1/* 作为备用数据源
3. 测试网验证后再切换到主网
4. 始终保留 fallback 机制
"""
7.8 性能参考
PERF = {
'ws_latency': {'p50': '50ms', 'p95': '200ms', 'p99': '500ms'},
'throughput': {'normal': '50 evt/s', 'peak': '200 evt/s'},
'history_api': {'p50': '100ms', 'p95': '500ms', 'limit': '10 req/s/IP'},
}
TIPS = """
1. 并发订阅不超过 5 个
2. 复杂操作放入后台任务队列
3. 历史查询每次控制在 10000 区块以内
4. 事件缓冲使用有界队列 + 背压
5. 批量处理比逐个处理吞吐量高 10-50 倍
"""
7.9 诊断工具
class Diagnostics:
def __init__(self):
self.received = 0
self.processed = 0
self.errors = 0
self.latencies: list[float] = []
def record(self, latency_ms: float):
self.latencies.append(latency_ms)
if len(self.latencies) > 10000:
self.latencies = self.latencies[-5000:]
def summary(self) -> dict:
s = sorted(self.latencies) if self.latencies else [0]
return {
'received': self.received, 'processed': self.processed,
'errors': self.errors,
'error_rate': round(self.errors / max(self.received, 1) * 100, 2),
'avg_latency_ms': round(sum(s) / len(s), 2),
'p95_latency_ms': round(s[int(len(s)*0.95)], 2),
}
7.10 部署检查清单
[ ] WebSocket: WSS + 指数退避重连 + 心跳保活
[ ] 事件: 去重 + 缓冲背压 + 检查点恢复
[ ] 错误: 重试 + 断路器 + 异常日志
[ ] 监控: 连接状态 + 延迟(P50/P95/P99) + 吞吐量 + 告警
[ ] 安全: 密钥加密存储 + Webhook 签名验证 + 速率限制
[ ] 测试: 单元测试 + 集成测试(重连) + 压力测试
[ ] 部署: 容器化 + 健康检查 + 优雅关闭 + 日志聚合
7.11 常见问题
Q: WebSocket 频繁断连?
- 检查网络和防火墙,确保启用心跳(30s 间隔),检查连接数限制
Q: 事件丢失?
- 通过
/agent/v1/events/history验证连续性,启用检查点恢复机制
Q: 处理延迟过高?
- 检查事件处理器中的阻塞操作,使用异步任务队列,增加并发数
Q: 订阅不生效?
- 验证查询语法和
chain_id(应为msg-chain-1),地址前缀应为msg
附录
A. 事件类型速查
| 类型 | 模块 | 说明 | 关键属性 |
|---|---|---|---|
transfer |
Bank | 转账 | sender, recipient, amount |
delegate |
Staking | 委托 | validator, delegator, amount |
undelegate |
Staking | 解委托 | + completion_time |
submit_proposal |
Governance | 提交提案 | proposal_id, proposer |
vote |
Governance | 投票 | proposal_id, voter, option |
wasm |
WASM | 合约事件 | _contract_address, action |
agent |
Agent | Agent 事件 | action, agent_id |
B. 查询表达式
| 表达式 | 说明 |
|---|---|
tm.event='Tx' |
所有交易 |
tm.event='Tx' AND wasm._contract_address='msg1...' |
特定合约 |
tm.event='Tx' AND agent.action='register' |
Agent 注册 |
tm.event='Tx' AND (agent.action='payment' OR agent.action='a2a') |
Agent 支付或通信 |
C. 参考资源
- MSG Chain 文档:
https://docs.msgchain.org/ - Cosmos SDK 事件:
https://docs.cosmos.network/v0.50/core/events.html - Tendermint WebSocket:
https://docs.cometbft.com/v0.38/guides/websocket/
本指南由 AI 生成,代码示例仅供参考。生产环境部署前请充分测试。端点和 API 参数以 MSG Chain 官方文档为准。
