dApp Docs/AI Agent 链上事件驱动架构与实时响应指南
Development reference. Not independently verified for production.

AI Agent 链上事件驱动架构与实时响应指南

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

1. 概述

1.1 什么是事件驱动架构

事件驱动架构(EDA)是一种软件范式,其中系统组件通过生产和消费事件通信。对于链上 AI Agent,这意味着 Agent 不再被动等待用户指令,而是主动监听区块链上的状态变更并实时响应。

在 MSG Chain 生态中,AI Agent 需处理四类事件:

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 链参数

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 频繁断连?

Q: 事件丢失?

Q: 处理延迟过高?

Q: 订阅不生效?


附录

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. 参考资源


本指南由 AI 生成,代码示例仅供参考。生产环境部署前请充分测试。端点和 API 参数以 MSG Chain 官方文档为准。