Webhook 事件推送与通知服务指南
数据来源:MSG Chain 代码库核实
主网状态: No-Go — 当前 MSGChain 主网裁决为 No-Go,以下内容反映代码实际状态,不代表生产可用。
1. Webhook 在区块链通知中的角色
1.1 为什么需要 Webhook
Webhook 是一种用户定义的 HTTP 回调机制,当链上特定事件发生时,服务端主动向用户注册的端点推送事件数据。与轮询(Polling)模式相比,Webhook 在区块链通知场景中具有根本性优势:
轮询的问题:
- 客户端每隔固定间隔(如 15 秒)拉取一次最新事件
- 空查询浪费大量带宽和计算资源(90%+ 的请求无新数据)
- 实时性受轮询间隔限制,无法做到秒级响应
- 高频率轮询会被视为 DoS 攻击,触发 API 限流
- 每个用户的轮询请求叠加,服务端压力随用户数线性增长
Webhook 的优势:
| 维度 | 轮询 | Webhook |
|---|---|---|
| 实时性 | 受间隔限制(秒级~分钟级) | 事件产生即推送(毫秒级) |
| 资源效率 | 大量空请求浪费带宽 | 仅推送有效事件 |
| 服务端负载 | 与用户数成正比 | 与事件量成正比 |
| 客户端复杂度 | 需维护定时器和状态 | 仅需 HTTP 端点接收 |
| 扩展性 | 每增 100 用户需 100x 轮询 | 事件合并推送,边缘扩展 |
1.2 Webhook 在 MSG Chain 中的位置
在 MSG Chain 生态中,Webhook 服务位于链上事件生产者和下游消费者之间,承担事件路由、转换和可靠投递的职责:
区块链节点 (Tendermint WSS / RPC)
|
事件索引器 (Indexer)
|
Webhook 路由器 (Router)
| +------------------+
+---> 用户服务 A (HTTP) |
+---> 用户服务 B (HTTP) |
+---> 用户 Bot C (HTTP) |
+---> 死信队列 (DLQ) |
+------------------+
1.3 适用场景
- 支付通知:监听
transfer或agent.payment事件,收到付款后推送通知到商家后端 - 合约事件监控:当监控的 CosmWasm 合约触发特定 action 时,回调用户 API
- 治理提醒:新提案、投票开始/结束、结果公布时推送通知
- 跨链桥事件:IBC 转账完成、中继失败等事件的通知
- AI Agent 回调:Agent A2A 消息到达、任务状态变更时通知 Agent 后端
1.4 端点速查
WEBHOOK_API_BASE = 'https://webhook.msgchain.org'
ENDPOINTS = {
'register_webhook': '/api/v1/webhooks',
'update_webhook': '/api/v1/webhooks/{id}',
'delete_webhook': '/api/v1/webhooks/{id}',
'list_webhooks': '/api/v1/webhooks',
'get_webhook': '/api/v1/webhooks/{id}',
'webhook_logs': '/api/v1/webhooks/{id}/logs',
'test_webhook': '/api/v1/webhooks/{id}/test',
'rotate_secret': '/api/v1/webhooks/{id}/secret',
'delivery_stats': '/api/v1/webhooks/{id}/stats',
}
链参数:
Chain ID: msg-chain-1
地址前缀: msg
原生代币: MSG (最小单位 umsg, 1 MSG = 10^6 umsg)
RPC WebSocket: wss://rpc.msgchain.org/websocket
REST API: https://api.msgchain.org
Webhook API: https://webhook.msgchain.org
2. 架构概览
2.1 整体架构
Webhook 通知服务由以下子系统构成:
+----------------------------------------------------------+
| 事件监听器 (Event Listener) |
| . 连接 Tendermint WSS 订阅 tm.event='Tx' |
| . 缓冲、去重、解析原始事件 |
| . 故障恢复、检查点持久化 |
+----------------------------------------------------------+
| 规范化事件
v
+----------------------------------------------------------+
| 过滤器引擎 (Filter Engine) |
| . 地址过滤: _contract_address, sender, recipient |
| . 事件类型过滤: transfer, wasm, agent, delegate |
| . 参数匹配: action='register', amount>1000 |
| . 条件表达式: AND/OR/NOT 组合 |
+----------------------------------------------------------+
| 匹配的事件
v
+----------------------------------------------------------+
| 回调规则引擎 (Rule Engine) |
| . 静态规则: 预定义匹配模式 |
| . 动态规则: 运行时加载, 热更新 |
| . 用户自定义 DSL: 条件表达式 -> 动作 |
+----------------------------------------------------------+
| 路由决策
v
+----------------------------------------------------------+
| Webhook 路由器 (Webhook Router) |
| . 并发投递 (HTTP POST) |
| . 负载均衡、用户配额管理 |
| . 重试策略、指数退避 |
| . 签名计算 (HMAC-SHA256) |
+----------+---------------------+----------------+----------+
| | |
成功 失败(可重试) 失败(不可重试)
v v v
+----------------+ +------------------+ +----------------+
| 成功记录 | | 重试队列(Retry Q)| | 死信队列(DLQ) |
| (PostgreSQL) | | Redis 延迟队列 | | S3/持久化 |
+----------------+ +------------------+ +----------------+
| 重试
v
+------------------+
| 告警 & 监控 |
| Prometheus + |
| Grafana + Alert |
+------------------+
2.2 核心组件职责
| 组件 | 职责 | 技术选型 |
|---|---|---|
| Event Listener | 连接链节点、订阅事件、解析标准化 | Python asyncio / Rust |
| Filter Engine | 多维度事件匹配 | 表达式树 + LRU 缓存 |
| Rule Engine | 用户规则管理与评估 | DSL 解释器 / CEL |
| Webhook Router | HTTP 回调投递、重试、签名 | aiohttp / httpx |
| Retry Queue | 失败事件延迟重试 | Redis Streams |
| DLQ | 最终失败的持久化存储 | S3 / PostgreSQL |
| Metrics | 延迟、成功率、吞吐监控 | Prometheus + Grafana |
2.3 数据流时序
链节点 Event Listener Filter+Rule Router 用户服务
| | | | |
|-- 新区块/Tx ------| | | |
| |-- 解析事件 -------| | |
| | |-- 匹配规则 ------| |
| | | |-- HTTP POST ----|
| | | |<- 200 OK -------|
| | | |-- 记录日志 -----|
| | | | |
| | | | (若失败) |
| | | |-- 延迟重试 ------|
| | | |-- 最终失败 > DLQ |
3. Webhook 服务部署
3.1 事件监听器
事件监听器是 Webhook 服务的入口,负责从 MSG Chain 节点持续获取链上事件。
import asyncio
import json
import logging
from datetime import datetime
from typing import AsyncGenerator, Optional
import httpx
import websockets
logger = logging.getLogger(__name__)
CHAIN_ID = 'msg-chain-1'
RPC_WS = 'wss://rpc.msgchain.org/websocket'
REST_API = 'https://api.msgchain.org'
class ChainEventListener:
"""链事件监听器:从 Tendermint WebSocket 订阅实时事件。"""
def __init__(
self,
ws_url: str = RPC_WS,
rest_url: str = REST_API,
query: str = "tm.event='Tx'",
buffer_size: int = 5000,
checkpoint_interval: int = 100,
):
self.ws_url = ws_url
self.rest_url = rest_url
self.query = query
self.buffer: asyncio.Queue[dict] = asyncio.Queue(maxsize=buffer_size)
self.checkpoint_interval = checkpoint_interval
self.last_height: Optional[int] = None
self._running = False
self._http = httpx.AsyncClient(base_url=rest_url, timeout=30.0)
async def _ws_listener(self):
"""WebSocket 事件订阅主循环。"""
async for event in self._subscribe():
try:
result = event.get('result', {})
data = result.get('data', {})
value = data.get('value', {})
tx_result = value.get('TxResult', {})
height_str = tx_result.get('height', '0')
height = int(height_str)
parsed = {
'height': height,
'tx_hash': tx_result.get('txhash', ''),
'timestamp': datetime.utcnow().isoformat() + 'Z',
'events': self._normalize_events(
tx_result.get('events', [])
),
'chain_id': CHAIN_ID,
}
await self.buffer.put(parsed)
self.last_height = height
except Exception as e:
logger.error(f'Parse event failed: {e}')
async def _subscribe(self) -> AsyncGenerator[dict, None]:
"""连接到 Tendermint WebSocket 并订阅事件。"""
while self._running:
try:
async with websockets.connect(self.ws_url) as ws:
await ws.send(json.dumps({
'jsonrpc': '2.0',
'method': 'subscribe',
'params': {'query': self.query},
'id': 1,
}))
response = json.loads(await ws.recv())
if 'error' in response:
raise ConnectionError(
'Subscribe failed: {response["error"]}'
)
async for message in ws:
if not self._running:
return
event = json.loads(message)
if event.get('id') is not None:
continue
if event.get('result', {}).get('type') == 'heartbeat':
continue
yield event
except (websockets.ConnectionClosed, ConnectionError) as e:
logger.warning(f'WS disconnected: {e}, reconnecting in 5s...')
await asyncio.sleep(5)
def _normalize_events(self, raw_events: list[dict]) -> list[dict]:
"""将 Tendermint 原始事件规范化为统一格式。"""
normalized = []
for event in raw_events:
attrs = {}
for attr in event.get('attributes', []):
key = attr.get('key', '')
value = attr.get('value', '')
if isinstance(key, bytes):
key = key.decode('utf-8')
if isinstance(value, bytes):
value = value.decode('utf-8')
if key in attrs:
if not isinstance(attrs[key], list):
attrs[key] = [attrs[key]]
attrs[key].append(value)
else:
attrs[key] = value
normalized.append({
'type': event.get('type', ''),
'attributes': attrs,
})
return normalized
async def get_event(self) -> dict:
"""从缓冲队列获取一个事件(带背压)。"""
return await self.buffer.get()
async def get_batch(
self, max_size: int = 100, timeout: float = 0.5
) -> list[dict]:
"""批量获取事件,用于批处理优化。"""
batch = []
try:
batch.append(await asyncio.wait_for(
self.buffer.get(), timeout=timeout
))
while len(batch) < max_size and not self.buffer.empty():
batch.append(self.buffer.get_nowait())
except asyncio.TimeoutError:
pass
return batch
async def run(self):
self._running = True
await self._ws_listener()
async def stop(self):
self._running = False
await self._http.aclose()
3.2 过滤器引擎
从原始事件流中筛选出匹配用户规则的事件。
from abc import ABC, abstractmethod
from enum import Enum
from typing import Any, Optional
class FilterMatchType(Enum):
EXACT = 'exact'
PREFIX = 'prefix'
SUFFIX = 'suffix'
CONTAINS = 'contains'
REGEX = 'regex'
GT = 'gt'
GTE = 'gte'
LT = 'lt'
LTE = 'lte'
EXISTS = 'exists'
NOT_EXISTS = 'not_exists'
IN = 'in'
class FilterCondition:
"""单条过滤条件。"""
def __init__(
self,
event_type: str,
attribute_key: str,
match_type: FilterMatchType,
value: Any = None,
):
self.event_type = event_type
self.attribute_key = attribute_key
self.match_type = match_type
self.value = value
def evaluate(self, event_type: str, attrs: dict) -> bool:
if event_type != self.event_type:
return False
attr_value = attrs.get(self.attribute_key)
if self.match_type == FilterMatchType.EXISTS:
return attr_value is not None
if self.match_type == FilterMatchType.NOT_EXISTS:
return attr_value is None
if attr_value is None:
return False
if self.match_type == FilterMatchType.EXACT:
return str(attr_value) == str(self.value)
elif self.match_type == FilterMatchType.PREFIX:
return str(attr_value).startswith(str(self.value))
elif self.match_type == FilterMatchType.SUFFIX:
return str(attr_value).endswith(str(self.value))
elif self.match_type == FilterMatchType.CONTAINS:
return str(self.value) in str(attr_value)
elif self.match_type == FilterMatchType.REGEX:
import re
return bool(re.search(str(self.value), str(attr_value)))
elif self.match_type in (FilterMatchType.GT, FilterMatchType.GTE,
FilterMatchType.LT, FilterMatchType.LTE):
return self._compare_numeric(attr_value, self.match_type)
elif self.match_type == FilterMatchType.IN:
return attr_value in (self.value or [])
return False
def _compare_numeric(self, actual: Any, op: FilterMatchType) -> bool:
try:
a = self._to_numeric(actual)
b = self._to_numeric(self.value)
if op == FilterMatchType.GT: return a > b
elif op == FilterMatchType.GTE: return a >= b
elif op == FilterMatchType.LT: return a < b
elif op == FilterMatchType.LTE: return a <= b
except (ValueError, TypeError):
return False
return False
def _to_numeric(self, val: Any) -> int:
s = str(val)
for suffix in ['umsg', 'msg', 'u']:
if s.endswith(suffix):
s = s[:-len(suffix)]
break
return int(s)
class FilterGroup:
class Logic(Enum):
AND = 'and'
OR = 'or'
def __init__(self, logic: Logic = Logic.AND):
self.logic = logic
self.conditions: list[FilterCondition] = []
self.groups: list['FilterGroup'] = []
def add_condition(self, condition: FilterCondition):
self.conditions.append(condition)
def add_group(self, group: 'FilterGroup'):
self.groups.append(group)
def evaluate(self, normalized_events: list[dict]) -> bool:
if not self.conditions and not self.groups:
return True
results = []
for cond in self.conditions:
matched = False
for ev in normalized_events:
if cond.evaluate(ev['type'], ev.get('attributes', {})):
matched = True
break
results.append(matched)
for g in self.groups:
results.append(g.evaluate(normalized_events))
if not results:
return True
if self.logic == FilterGroup.Logic.AND:
return all(results)
else:
return any(results)
class FilterEngine:
"""过滤器引擎:对规范化事件执行多层次条件匹配。"""
def __init__(self):
self.filters: dict[str, FilterGroup] = {}
def register_filter(self, filter_id: str, group: FilterGroup):
self.filters[filter_id] = group
def unregister_filter(self, filter_id: str):
self.filters.pop(filter_id, None)
def match(self, normalized_events: list[dict]) -> list[str]:
"""返回所有匹配的 filter_id 列表。"""
matched = []
for fid, group in self.filters.items():
if group.evaluate(normalized_events):
matched.append(fid)
return matched
# 预设过滤模板
def filter_transfer_from(sender_addr: str) -> FilterGroup:
g = FilterGroup(FilterGroup.Logic.AND)
g.add_condition(FilterCondition(
'transfer', 'sender', FilterMatchType.EXACT, sender_addr
))
return g
def filter_transfer_above(min_amount: int, denom: str = 'umsg') -> FilterGroup:
g = FilterGroup(FilterGroup.Logic.AND)
g.add_condition(FilterCondition(
'transfer', 'amount', FilterMatchType.GTE, f'{min_amount}{denom}'
))
return g
def filter_contract_action(contract_addr: str, action: Optional[str] = None) -> FilterGroup:
g = FilterGroup(FilterGroup.Logic.AND)
g.add_condition(FilterCondition(
'wasm', '_contract_address', FilterMatchType.EXACT, contract_addr
))
if action:
g.add_condition(FilterCondition(
'wasm', 'action', FilterMatchType.EXACT, action
))
return g
def filter_agent_event(agent_id: str, *actions: str) -> FilterGroup:
g = FilterGroup(FilterGroup.Logic.AND)
g.add_condition(FilterCondition(
'agent', 'agent_id', FilterMatchType.EXACT, agent_id
))
if actions:
ag = FilterGroup(FilterGroup.Logic.OR)
for act in actions:
ag.add_condition(FilterCondition(
'agent', 'action', FilterMatchType.EXACT, act
))
g.add_group(ag)
return g
def filter_or(*filters: FilterGroup) -> FilterGroup:
combined = FilterGroup(FilterGroup.Logic.OR)
for f in filters: combined.add_group(f)
return combined
def filter_and(*filters: FilterGroup) -> FilterGroup:
combined = FilterGroup(FilterGroup.Logic.AND)
for f in filters: combined.add_group(f)
return combined
3.3 HTTP 回调与重试策略
Webhook 路由器向用户端点投递事件的完整实现。
import hashlib
import hmac
import json
import time
from dataclasses import dataclass, field
from datetime import datetime
from typing import Optional, Callable
import httpx
class SignatureScheme(Enum):
HMAC_SHA256 = 'hmac-sha256'
HMAC_SHA512 = 'hmac-sha512'
NONE = 'none'
@dataclass
class WebhookEndpoint:
"""用户注册的 Webhook 端点配置。"""
id: str
url: str
secret: str
signature_scheme: SignatureScheme = SignatureScheme.HMAC_SHA256
retry_max: int = 5
retry_base_delay: float = 1.0
retry_max_delay: float = 60.0
timeout: float = 10.0
headers: dict = field(default_factory=lambda: {
'Content-Type': 'application/json',
'User-Agent': 'MSGChain-Webhook/1.0',
})
enabled: bool = True
created_at: str = ''
updated_at: str = ''
def compute_signature(
payload: bytes, secret: str, scheme: SignatureScheme
) -> str:
if scheme == SignatureScheme.NONE:
return ''
algo = hashlib.sha256 if scheme == SignatureScheme.HMAC_SHA256 else hashlib.sha512
return hmac.new(secret.encode(), payload, algo).hexdigest()
class WebhookDelivery:
"""单次 Webhook 投递。"""
def __init__(
self,
endpoint: WebhookEndpoint,
event: dict,
attempt: int = 1,
):
self.endpoint = endpoint
self.event = event
self.attempt = attempt
self.status_code: Optional[int] = None
self.response_body: Optional[str] = None
self.error: Optional[str] = None
self.latency_ms: Optional[float] = None
self.timestamp: str = datetime.utcnow().isoformat() + 'Z'
async def send(self, client: httpx.AsyncClient) -> bool:
"""执行一次 HTTP POST 回调。"""
payload = json.dumps({
'event': self.event,
'webhook_id': self.endpoint.id,
'delivery_attempt': self.attempt,
'timestamp': self.timestamp,
}).encode()
signature = compute_signature(
payload, self.endpoint.secret, self.endpoint.signature_scheme
)
headers = dict(self.endpoint.headers)
if signature:
headers['X-MSGChain-Signature'] = signature
headers['X-MSGChain-Signature-Algorithm'] = (
self.endpoint.signature_scheme.value
)
headers['X-MSGChain-Delivery-Id'] = (
f'{self.endpoint.id}:{self.event.get("tx_hash", "")}'
)
headers['X-MSGChain-Attempt'] = str(self.attempt)
headers['X-MSGChain-Timestamp'] = self.timestamp
start = time.monotonic()
try:
response = await client.post(
self.endpoint.url,
content=payload,
headers=headers,
timeout=self.endpoint.timeout,
)
self.latency_ms = (time.monotonic() - start) * 1000
self.status_code = response.status_code
self.response_body = response.text[:1024]
return 200 <= response.status_code < 300
except httpx.TimeoutException as e:
self.latency_ms = (time.monotonic() - start) * 1000
self.error = f'Timeout: {e}'
return False
except httpx.RequestError as e:
self.latency_ms = (time.monotonic() - start) * 1000
self.error = f'RequestError: {e}'
return False
class RetryPolicy:
"""重试策略:指数退避 + 抖动。"""
def __init__(
self,
max_retries: int = 5,
base_delay: float = 1.0,
max_delay: float = 60.0,
jitter: float = 0.1,
retryable_statuses: Optional[set[int]] = None,
):
self.max_retries = max_retries
self.base_delay = base_delay
self.max_delay = max_delay
self.jitter = jitter
self.retryable_statuses = retryable_statuses or {
408, 429, 500, 502, 503, 504,
}
def should_retry(self, attempt: int, delivery: WebhookDelivery) -> bool:
"""判断是否应该重试。"""
if attempt >= self.max_retries:
return False
if delivery.status_code and delivery.status_code < 500:
if delivery.status_code not in self.retryable_statuses:
return False
if delivery.error and 'Timeout' in delivery.error:
return True
return True
def get_delay(self, attempt: int) -> float:
"""计算第 N 次重试的延迟。"""
import random
delay = min(
self.base_delay * (2 ** (attempt - 1)),
self.max_delay,
)
jitter_amount = delay * self.jitter
return max(0.1, delay + random.uniform(-jitter_amount, jitter_amount))
class WebhookRouter:
"""Webhook 路由器:管理端点、执行投递、处理重试。"""
def __init__(
self,
retry_policy: Optional[RetryPolicy] = None,
):
self.endpoints: dict[str, WebhookEndpoint] = {}
self.retry_policy = retry_policy or RetryPolicy()
self._client = httpx.AsyncClient(
limits=httpx.Limits(max_keepalive_connections=200, max_connections=1000),
timeout=30.0,
)
self._on_delivery_complete: Optional[Callable] = None
def register_endpoint(self, endpoint: WebhookEndpoint):
self.endpoints[endpoint.id] = endpoint
def unregister_endpoint(self, endpoint_id: str):
self.endpoints.pop(endpoint_id, None)
def set_delivery_callback(self, cb: Callable):
self._on_delivery_complete = cb
async def deliver(
self, endpoint_id: str, event: dict
) -> WebhookDelivery:
"""向指定端点投递事件,含重试逻辑。"""
endpoint = self.endpoints.get(endpoint_id)
if not endpoint or not endpoint.enabled:
raise ValueError(f'Endpoint {endpoint_id} not found or disabled')
attempt = 1
while attempt <= endpoint.retry_max:
delivery = WebhookDelivery(endpoint, event, attempt=attempt)
success = await delivery.send(self._client)
if self._on_delivery_complete:
try:
self._on_delivery_complete(delivery)
except Exception:
logger.exception('Delivery callback failed')
if success:
return delivery
if not self.retry_policy.should_retry(attempt, delivery):
return delivery
delay = self.retry_policy.get_delay(attempt)
logger.warning(
f'Delivery to {endpoint.url} failed (attempt {attempt}), '
f'retrying in {delay:.1f}s: {delivery.error or delivery.status_code}'
)
await asyncio.sleep(delay)
attempt += 1
return delivery
async def deliver_to_matched(
self, event: dict, matched_filter_ids: list[str]
):
"""向匹配的所有端点投递。"""
tasks = []
for fid in matched_filter_ids:
if fid in self.endpoints:
tasks.append(self.deliver(fid, event))
if tasks:
results = await asyncio.gather(*tasks, return_exceptions=True)
for result in results:
if isinstance(result, Exception):
logger.error(f'Bulk delivery error: {result}')
async def close(self):
await self._client.aclose()
3.4 完整服务编排
class WebhookService:
"""Webhook 通知服务的主编排器。"""
def __init__(
self,
listener: ChainEventListener,
filter_engine: FilterEngine,
router: WebhookRouter,
):
self.listener = listener
self.filter_engine = filter_engine
self.router = router
self._running = False
self._stats = {
'events_received': 0,
'events_matched': 0,
'deliveries_sent': 0,
'deliveries_success': 0,
'deliveries_failed': 0,
}
async def process_batch(self, events: list[dict]):
"""处理一批事件的完整流程。"""
for event in events:
self._stats['events_received'] += 1
normalized = event.get('events', [])
matched = self.filter_engine.match(normalized)
if not matched:
continue
self._stats['events_matched'] += 1
await self.router.deliver_to_matched(event, matched)
async def run(self):
"""运行事件处理循环。"""
self._running = True
listener_task = asyncio.create_task(self.listener.run())
try:
while self._running:
batch = await self.listener.get_batch(max_size=50, timeout=1.0)
if batch:
await self.process_batch(batch)
finally:
listener_task.cancel()
await self.router.close()
async def stop(self):
self._running = False
await self.listener.stop()
def stats(self) -> dict:
return dict(self._stats)
4. 事件过滤
4.1 地址过滤
地址过滤是最常用的过滤维度。在 MSG Chain 中,所有地址以 msg 前缀开头。
ADDRESS_FILTERS = {
'sender': '转账发送方地址',
'recipient': '转账接收方地址',
'_contract_address': '合约地址',
'agent_id': 'Agent 标识',
'delegator': '委托方',
'validator': '验证人',
'proposer': '提案人',
'voter': '投票人',
}
ADDRESS_PATTERN = r'^msg[a-z0-9]{10,64}$'
def build_address_filter(
event_type: str,
address_key: str,
address_value: str,
) -> FilterGroup:
"""构建地址过滤条件。"""
import re
if not re.match(ADDRESS_PATTERN, address_value):
raise ValueError(f'Invalid MSG address: {address_value}')
g = FilterGroup(FilterGroup.Logic.AND)
g.add_condition(FilterCondition(
event_type, address_key, FilterMatchType.EXACT, address_value
))
return g
def build_address_list_filter(
event_type: str,
address_key: str,
addresses: list[str],
) -> FilterGroup:
"""匹配地址列表中的任意一个。"""
g = FilterGroup(FilterGroup.Logic.OR)
for addr in addresses:
g.add_condition(FilterCondition(
event_type, address_key, FilterMatchType.EXACT, addr
))
return g
4.2 事件类型过滤
MSG Chain 的核心事件类型:
| 事件类型 | 模块 | 触发条件 | 关键属性 |
|---|---|---|---|
transfer |
bank | MsgSend 执行 | sender, recipient, amount |
delegate |
staking | MsgDelegate | validator, delegator, amount, new_shares |
undelegate |
staking | MsgUndelegate | + completion_time |
redelegate |
staking | MsgBeginRedelegate | + src_validator, dst_validator |
submit_proposal |
gov | MsgSubmitProposal | proposal_id, proposal_type, proposer |
deposit |
gov | MsgDeposit | depositor, amount |
vote |
gov | MsgVote | voter, option |
wasm |
wasm | 合约执行 | _contract_address, action, ... |
agent |
agent | Agent 操作 | agent_id, action, owner |
message |
wasm | MsgExecuteContract | sender, contract, msg |
coin_received |
bank | 代币接收 | receiver, amount |
coin_spent |
bank | 代币花费 | spender, amount |
slash |
slashing | 验证人惩罚 | validator, amount, reason |
class EventTypeFilter:
"""按事件类型过滤。"""
def __init__(self, allowed_types: set[str]):
self.allowed = allowed_types
def matches(self, normalized_events: list[dict]) -> bool:
return any(ev['type'] in self.allowed for ev in normalized_events)
class CompositeFilter:
"""多层组合过滤:先过类型,再过滤属性。"""
def __init__(self):
self.type_filter: Optional[EventTypeFilter] = None
self.attribute_filters: list[FilterGroup] = []
def set_types(self, types: set[str]):
self.type_filter = EventTypeFilter(types)
def add_attribute_filter(self, fg: FilterGroup):
self.attribute_filters.append(fg)
def match(self, normalized_events: list[dict]) -> bool:
if self.type_filter and not self.type_filter.matches(normalized_events):
return False
for af in self.attribute_filters:
if not af.evaluate(normalized_events):
return False
return True
4.3 参数匹配
from typing import Union, Callable
class ParamMatcher:
"""支持多种匹配模式的参数匹配器。"""
def __init__(self):
self.rules: list[tuple[str, str, Callable]] = []
def exact(self, event_type: str, key: str, value: str):
self.rules.append((
event_type, key,
lambda actual: str(actual) == value
))
def prefix(self, event_type: str, key: str, prefix_val: str):
self.rules.append((
event_type, key,
lambda actual: str(actual).startswith(prefix_val)
))
def suffix(self, event_type: str, key: str, suffix_val: str):
self.rules.append((
event_type, key,
lambda actual: str(actual).endswith(suffix_val)
))
def contains(self, event_type: str, key: str, substr: str):
self.rules.append((
event_type, key,
lambda actual: str(substr) in str(actual)
))
def numeric_gt(self, event_type: str, key: str, threshold: Union[int, float]):
self.rules.append((
event_type, key,
lambda actual: self._parse_amount(actual) > threshold
))
def numeric_gte(self, event_type: str, key: str, threshold: Union[int, float]):
self.rules.append((
event_type, key,
lambda actual: self._parse_amount(actual) >= threshold
))
def numeric_lt(self, event_type: str, key: str, threshold: Union[int, float]):
self.rules.append((
event_type, key,
lambda actual: self._parse_amount(actual) < threshold
))
def numeric_lte(self, event_type: str, key: str, threshold: Union[int, float]):
self.rules.append((
event_type, key,
lambda actual: self._parse_amount(actual) <= threshold
))
def in_set(self, event_type: str, key: str, values: set):
self.rules.append((
event_type, key,
lambda actual: actual in values
))
def _parse_amount(self, val) -> float:
s = str(val)
for suffix in ['umsg', 'msg']:
if s.endswith(suffix):
s = s[:-len(suffix)]
break
try:
return float(s)
except ValueError:
return 0.0
def match(self, normalized_events: list[dict]) -> bool:
for event_type, key, predicate in self.rules:
found = False
for ev in normalized_events:
if ev['type'] != event_type:
continue
actual = ev.get('attributes', {}).get(key)
if actual is not None and predicate(actual):
found = True
break
if not found:
return False
return True
# 使用示例
matcher = ParamMatcher()
matcher.exact('wasm', 'action', 'increment')
matcher.numeric_gt('transfer', 'amount', 1000000000)
matcher.in_set('agent', 'action', {'register', 'payment'})
4.4 条件表达式 DSL
除了代码 API,Webhook 服务还提供字符串形式的条件表达式,方便用户通过 API 配置过滤规则。
"""
条件表达式语法 (Condition Expression Language):
表达式格式:
<event_type>.<attribute> <operator> <value>
运算符:
= 精确匹配
!= 不等于
~ 正则匹配
>, >= 大于 (数值比较)
<, <= 小于 (数值比较)
IN 集合包含 a IN (x, y, z)
EXISTS 存在判断
NOT 取反
组合:
AND 逻辑与
OR 逻辑或
( ) 分组
示例:
transfer.amount >= 1000000umsg
wasm._contract_address = 'msg1contract...'
agent.action IN ('register', 'payment')
transfer.amount >= 1000000umsg AND agent.action = 'payment'
(wasm.action = 'increment' OR wasm.action = 'reset')
"""
class CondExprParser:
"""条件表达式解析器。"""
def tokenize(self, expr: str) -> list[dict]:
tokens = []
i = 0
while i < len(expr):
c = expr[i]
if c.isspace():
i += 1
continue
if c == '(':
tokens.append({'type': 'LPAREN'})
i += 1
elif c == ')':
tokens.append({'type': 'RPAREN'})
i += 1
elif c in "'\"":
j = i + 1
while j < len(expr) and expr[j] != c:
j += 1
tokens.append({'type': 'STRING', 'value': expr[i+1:j]})
i = j + 1
elif expr[i:i+3] == 'AND':
tokens.append({'type': 'AND'})
i += 3
elif expr[i:i+2] == 'OR':
tokens.append({'type': 'OR'})
i += 2
elif expr[i:i+6] == 'EXISTS':
tokens.append({'type': 'EXISTS'})
i += 6
elif expr[i:i+3] == 'NOT':
tokens.append({'type': 'NOT'})
i += 3
elif expr.startswith('IN', i):
tokens.append({'type': 'IN'})
i += 2
elif expr[i:i+2] == '>=':
tokens.append({'type': 'OP', 'value': '>='})
i += 2
elif expr[i:i+2] == '<=':
tokens.append({'type': 'OP', 'value': '<='})
i += 2
elif expr[i:i+2] == '!=':
tokens.append({'type': 'OP', 'value': '!='})
i += 2
elif c == '>':
tokens.append({'type': 'OP', 'value': '>'})
i += 1
elif c == '<':
tokens.append({'type': 'OP', 'value': '<'})
i += 1
elif c == '=':
tokens.append({'type': 'OP', 'value': '='})
i += 1
elif c == '~':
tokens.append({'type': 'OP', 'value': '~'})
i += 1
elif c == ',':
tokens.append({'type': 'COMMA'})
i += 1
elif c == '.':
tokens.append({'type': 'DOT'})
i += 1
else:
j = i
while j < len(expr) and expr[j] not in ' \t\n\r()=><!,~':
j += 1
tok_val = expr[i:j]
if tok_val in ('AND', 'OR', 'IN', 'NOT', 'EXISTS'):
tokens.append({'type': tok_val})
else:
tokens.append({'type': 'IDENT', 'value': tok_val})
i = j
return tokens
def parse_to_filter(self, expr: str) -> FilterGroup:
"""将条件表达式解析为 FilterGroup。"""
tokens = self.tokenize(expr)
return self._parse_or(tokens, 0)[0]
def _parse_or(self, tokens: list[dict], pos: int):
left, pos = self._parse_and(tokens, pos)
while pos < len(tokens) and tokens[pos]['type'] == 'OR':
right, pos = self._parse_and(tokens, pos + 1)
g = FilterGroup(FilterGroup.Logic.OR)
g.add_group(left)
g.add_group(right)
left = g
return left, pos
def _parse_and(self, tokens: list[dict], pos: int):
left, pos = self._parse_unary(tokens, pos)
while pos < len(tokens) and tokens[pos]['type'] == 'AND':
right, pos = self._parse_unary(tokens, pos + 1)
g = FilterGroup(FilterGroup.Logic.AND)
g.add_group(left)
g.add_group(right)
left = g
return left, pos
def _parse_unary(self, tokens: list[dict], pos: int):
if pos >= len(tokens):
raise ValueError('Unexpected end of expression')
if tokens[pos]['type'] == 'LPAREN':
inner, pos = self._parse_or(tokens, pos + 1)
if pos >= len(tokens) or tokens[pos]['type'] != 'RPAREN':
raise ValueError('Missing closing parenthesis')
return inner, pos + 1
return self._parse_condition(tokens, pos)
def _parse_condition(self, tokens: list[dict], pos: int):
if pos + 2 >= len(tokens):
raise ValueError('Incomplete condition')
if tokens[pos]['type'] != 'IDENT':
raise ValueError(f'Expected event type, got {tokens[pos]}')
event_type = tokens[pos]['value']
pos += 1
if tokens[pos]['type'] != 'DOT':
raise ValueError('Expected "."')
pos += 1
if tokens[pos]['type'] != 'IDENT':
raise ValueError('Expected attribute key')
key = tokens[pos]['value']
pos += 1
if tokens[pos]['type'] == 'EXISTS':
g = FilterGroup()
g.add_condition(FilterCondition(
event_type, key, FilterMatchType.EXISTS, None
))
return g, pos + 1
if tokens[pos]['type'] != 'OP':
raise ValueError(f'Expected operator, got {tokens[pos]}')
op = tokens[pos]['value']
pos += 1
if tokens[pos]['type'] == 'IN':
pos += 1
if pos >= len(tokens) or tokens[pos]['type'] != 'LPAREN':
raise ValueError('Expected "(" after IN')
pos += 1
values = []
while pos < len(tokens) and tokens[pos]['type'] != 'RPAREN':
if tokens[pos]['type'] in ('STRING', 'IDENT'):
values.append(tokens[pos]['value'])
pos += 1
if pos < len(tokens) and tokens[pos]['type'] == 'COMMA':
pos += 1
if pos >= len(tokens) or tokens[pos]['type'] != 'RPAREN':
raise ValueError('Expected ")" after IN list')
g = FilterGroup(FilterGroup.Logic.OR)
for v in values:
g.add_condition(FilterCondition(
event_type, key, FilterMatchType.EXACT, v
))
return g, pos + 1
value = tokens[pos]['value'] if tokens[pos]['type'] in ('STRING', 'IDENT') else ''
pos += 1
match_map = {
'=': FilterMatchType.EXACT,
'~': FilterMatchType.REGEX,
'>': FilterMatchType.GT,
'>=': FilterMatchType.GTE,
'<': FilterMatchType.LT,
'<=': FilterMatchType.LTE,
}
match_type = match_map.get(op)
if match_type is None:
raise ValueError(f'Unknown operator: {op}')
condition = FilterCondition(event_type, key, match_type, value)
g = FilterGroup()
g.add_condition(condition)
return g, pos
class DSLFilterEngine:
"""基于 DSL 的过滤引擎,支持从字符串表达式创建过滤器。"""
def __init__(self):
self.parser = CondExprParser()
def create_filter(self, expression: str) -> FilterGroup:
return self.parser.parse_to_filter(expression)
def match_expression(
self, expression: str, normalized_events: list[dict]
) -> bool:
return self.create_filter(expression).evaluate(normalized_events)
# DSL 使用示例
dsl = DSLFilterEngine()
f1 = dsl.create_filter("wasm._contract_address = 'msg1contract...' AND wasm.action = 'increment'")
f2 = dsl.create_filter("transfer.amount >= 1000000000umsg")
f3 = dsl.create_filter("agent.action IN ('register', 'payment', 'a2a')")
f4 = dsl.create_filter("(transfer.amount >= 500000000umsg OR delegate.amount >= 100000000umsg)")
5. 回调规则引擎
5.1 静态规则
静态规则在 Webhook 注册时定义,规则内容在生命周期内不变。
from dataclasses import dataclass, field
from datetime import datetime
from typing import Optional
@dataclass
class StaticRule:
"""静态规则:注册后不可修改。"""
id: str
user_id: str
name: str
filter_group: FilterGroup
endpoint: WebhookEndpoint
description: str = ''
enabled: bool = True
created_at: str = field(default_factory=lambda: datetime.utcnow().isoformat() + 'Z')
class StaticRuleManager:
"""静态规则管理器。"""
def __init__(self):
self.rules: dict[str, StaticRule] = {}
def add_rule(self, rule: StaticRule):
self.rules[rule.id] = rule
def remove_rule(self, rule_id: str):
self.rules.pop(rule_id, None)
def get_rules_for_user(self, user_id: str) -> list[StaticRule]:
return [r for r in self.rules.values() if r.user_id == user_id and r.enabled]
def match_rules(self, normalized_events: list[dict]) -> list[StaticRule]:
matched = []
for rule in self.rules.values():
if not rule.enabled:
continue
if rule.filter_group.evaluate(normalized_events):
matched.append(rule)
return matched
5.2 动态规则
动态规则支持运行时加载、热更新,无需重启服务。
import hashlib
from typing import Optional
class DynamicRule:
"""动态规则:支持运行时热更新。"""
def __init__(
self,
rule_id: str,
user_id: str,
name: str,
filter_expr: str,
endpoint_url: str,
endpoint_secret: str,
priority: int = 0,
enabled: bool = True,
version: int = 1,
):
self.rule_id = rule_id
self.user_id = user_id
self.name = name
self.filter_expr = filter_expr
self.endpoint_url = endpoint_url
self.endpoint_secret = endpoint_secret
self.priority = priority
self.enabled = enabled
self.version = version
self._filter_group: Optional[FilterGroup] = None
self._compiled = False
def compile(self, parser: CondExprParser):
if self.filter_expr:
self._filter_group = parser.parse_to_filter(self.filter_expr)
self._compiled = True
def evaluate(self, normalized_events: list[dict]) -> bool:
if not self._compiled:
raise RuntimeError(f'Rule {self.rule_id} not compiled')
if not self._filter_group:
return True
return self._filter_group.evaluate(normalized_events)
@property
def checksum(self) -> str:
raw = f'{self.filter_expr}:{self.endpoint_url}:{self.version}'
return hashlib.sha256(raw.encode()).hexdigest()[:16]
class DynamicRuleStore:
"""动态规则存储:支持从外部数据源加载和热更新。"""
def __init__(self, parser: Optional[CondExprParser] = None):
self.parser = parser or CondExprParser()
self.rules: dict[str, DynamicRule] = {}
self._version_map: dict[str, int] = {}
def load_rules(self, rules: list[DynamicRule]):
for rule in rules:
rule.compile(self.parser)
self.rules[rule.rule_id] = rule
self._version_map[rule.rule_id] = rule.version
def reload_rules(self, rules: list[DynamicRule]) -> int:
"""热更新规则。仅更新版本号变化的规则。"""
updated = 0
for rule in rules:
existing = self.rules.get(rule.rule_id)
if existing and existing.version >= rule.version:
continue
rule.compile(self.parser)
self.rules[rule.rule_id] = rule
self._version_map[rule.rule_id] = rule.version
updated += 1
return updated
def get_enabled_rules(self) -> list[DynamicRule]:
return sorted(
[r for r in self.rules.values() if r.enabled],
key=lambda r: r.priority,
reverse=True,
)
def match_rules(self, normalized_events: list[dict]) -> list[DynamicRule]:
matched = []
for rule in self.get_enabled_rules():
if rule.evaluate(normalized_events):
matched.append(rule)
return matched
def remove_rule(self, rule_id: str):
self.rules.pop(rule_id, None)
self._version_map.pop(rule_id, None)
def size(self) -> int:
return len(self.rules)
# 热更新示例
async def hot_reload_loop(store: DynamicRuleStore, interval: float = 30.0):
"""定期从数据源拉取规则进行热更新。"""
while True:
await asyncio.sleep(interval)
fresh_rules = await fetch_rules_from_db()
updated = store.reload_rules(fresh_rules)
if updated:
logger.info(f'Hot reloaded {updated} rules')
async def fetch_rules_from_db() -> list[DynamicRule]:
"""模拟从数据库获取规则。"""
return [
DynamicRule(
rule_id='rule_001',
user_id='user_msg1abc...',
name='Large Transfer Alert',
filter_expr="transfer.amount >= 10000000000umsg",
endpoint_url='https://myapp.com/webhooks/msgchain',
endpoint_secret='whsec_abc123',
priority=10,
),
]
5.3 用户自定义规则 DSL
提供完整的 DSL 供用户通过 API 自定义规则。
class UserRule:
"""用户自定义规则。"""
def __init__(
self,
rule_id: str,
user_id: str,
name: str,
condition_expr: str,
actions: list[dict],
options: Optional[dict] = None,
):
self.rule_id = rule_id
self.user_id = user_id
self.name = name
self.condition_expr = condition_expr
self.actions = actions
self.options = options or {}
self._condition: Optional[FilterGroup] = None
def compile(self, parser: CondExprParser):
self._condition = parser.parse_to_filter(self.condition_expr)
def evaluate(self, normalized_events: list[dict]) -> bool:
if not self._condition:
return True
return self._condition.evaluate(normalized_events)
def render_body(self, event: dict, webhook_id: str) -> dict:
"""渲染包含模板变量的请求体。"""
body = {
'webhook_id': webhook_id,
'rule_id': self.rule_id,
'rule_name': self.name,
'event': event,
'timestamp': datetime.utcnow().isoformat() + 'Z',
}
return body
class UserRuleManager:
"""用户规则管理器:提供完整的 CRUD 和评估功能。"""
def __init__(self):
self.parser = CondExprParser()
self.rules: dict[str, UserRule] = {}
self._user_index: dict[str, set[str]] = {}
async def create_rule(self, rule: UserRule) -> str:
rule.compile(self.parser)
rid = rule.rule_id
self.rules[rid] = rule
self._user_index.setdefault(rule.user_id, set()).add(rid)
return rid
async def update_rule(self, rule_id: str, rule: UserRule):
rule.compile(self.parser)
old = self.rules.get(rule_id)
if old and old.user_id != rule.user_id:
self._user_index.get(old.user_id, set()).discard(rule_id)
self.rules[rule_id] = rule
self._user_index.setdefault(rule.user_id, set()).add(rule_id)
async def delete_rule(self, rule_id: str):
rule = self.rules.pop(rule_id, None)
if rule:
self._user_index.get(rule.user_id, set()).discard(rule_id)
async def get_user_rules(self, user_id: str) -> list[UserRule]:
rids = self._user_index.get(user_id, set())
return [self.rules[rid] for rid in rids if rid in self.rules]
async def match_rules(self, event: dict, user_id: str) -> list[UserRule]:
normalized = event.get('events', [])
matched = []
for rule in await self.get_user_rules(user_id):
if rule.evaluate(normalized):
matched.append(rule)
return matched
# REST API 路由模拟
class WebhookAPI:
"""Webhook 管理和规则配置 API。"""
def __init__(self, rule_mgr: UserRuleManager):
self.rule_mgr = rule_mgr
async def register_webhook(self, request: dict) -> dict:
"""POST /api/v1/webhooks"""
user_id = request['user_id']
rule = UserRule(
rule_id=self._gen_id(),
user_id=user_id,
name=request['name'],
condition_expr=request.get('condition', ''),
actions=request.get('actions', []),
options=request.get('options', {}),
)
rid = await self.rule_mgr.create_rule(rule)
return {'webhook_id': rid, 'status': 'active'}
async def update_webhook(self, webhook_id: str, request: dict) -> dict:
"""PUT /api/v1/webhooks/{id}"""
old = self.rule_mgr.rules.get(webhook_id)
if not old:
raise ValueError('Webhook not found')
updated = UserRule(
rule_id=webhook_id,
user_id=old.user_id,
name=request.get('name', old.name),
condition_expr=request.get('condition', old.condition_expr),
actions=request.get('actions', old.actions),
options=request.get('options', old.options),
)
await self.rule_mgr.update_rule(webhook_id, updated)
return {'webhook_id': webhook_id, 'status': 'updated'}
async def delete_webhook(self, webhook_id: str) -> dict:
"""DELETE /api/v1/webhooks/{id}"""
await self.rule_mgr.delete_rule(webhook_id)
return {'webhook_id': webhook_id, 'status': 'deleted'}
async def list_webhooks(self, user_id: str) -> list[dict]:
"""GET /api/v1/webhooks?user_id=..."""
rules = await self.rule_mgr.get_user_rules(user_id)
return [{'id': r.rule_id, 'name': r.name, 'condition': r.condition_expr} for r in rules]
async def test_webhook(self, webhook_id: str) -> dict:
"""POST /api/v1/webhooks/{id}/test"""
rule = self.rule_mgr.rules.get(webhook_id)
if not rule:
raise ValueError('Webhook not found')
test_event = {
'height': 123456,
'tx_hash': '0xTEST00000000000000000000000000000000000000000000000000000',
'events': [{'type': 'transfer', 'attributes': {'sender': 'msg1s...', 'recipient': 'msg1r...', 'amount': '1000000umsg'}}],
}
matched = rule.evaluate(test_event.get('events', []))
return {'webhook_id': webhook_id, 'test_event': test_event, 'matched': matched}
def _gen_id(self) -> str:
import uuid
return f'wh_{uuid.uuid4().hex[:24]}'
6. 可靠投递
6.1 At-Least-Once 语义
Webhook 服务保证每条匹配事件至少投递一次。
class AtLeastOnceDelivery:
"""At-Least-Once 投递保证:重试直到成功或达到上限。"""
def __init__(self, router: WebhookRouter, retry_queue: 'RetryQueue', max_retries: int = 10):
self.router = router
self.retry_queue = retry_queue
self.max_retries = max_retries
async def deliver(self, endpoint_id: str, event: dict) -> dict:
attempt = 0
while attempt < self.max_retries:
attempt += 1
delivery = WebhookDelivery(self.router.endpoints[endpoint_id], event, attempt=attempt)
success = await delivery.send(self.router._client)
if success:
return {'status': 'delivered', 'attempts': attempt, 'endpoint_id': endpoint_id}
if attempt < self.max_retries:
delay = self.router.retry_policy.get_delay(attempt)
await self.retry_queue.schedule(endpoint_id, event, attempt + 1, delay)
await asyncio.sleep(delay)
await self.retry_queue.send_to_dlq(endpoint_id, event, attempt)
return {'status': 'failed', 'attempts': attempt, 'endpoint_id': endpoint_id, 'reason': 'max_retries_exceeded'}
6.2 Exactly-Once 语义
class DedupStore:
"""去重存储:记录投递状态。"""
def __init__(self, redis_client=None, ttl: int = 86400):
self._redis = redis_client
self._ttl = ttl
self._memory_store: dict[str, str] = {}
async def is_delivered(self, delivery_id: str) -> bool:
if self._redis:
return await self._redis.get(f'dedup:{delivery_id}') == 'delivered'
return self._memory_store.get(delivery_id) == 'delivered'
async def mark_delivered(self, delivery_id: str):
if self._redis:
await self._redis.set(f'dedup:{delivery_id}', 'delivered', ex=self._ttl)
else:
self._memory_store[delivery_id] = 'delivered'
async def mark_in_progress(self, delivery_id: str):
if self._redis:
await self._redis.set(f'dedup:{delivery_id}', 'in_progress', ex=self._ttl)
else:
self._memory_store[delivery_id] = 'in_progress'
async def mark_failed(self, delivery_id: str):
if self._redis:
await self._redis.set(f'dedup:{delivery_id}', 'failed', ex=self._ttl)
else:
self._memory_store[delivery_id] = 'failed'
class ExactlyOnceDelivery:
"""Exactly-Once 投递保证:幂等性 + 去重。"""
def __init__(self, dedup_store: DedupStore, router: WebhookRouter):
self.dedup_store = dedup_store
self.router = router
def _delivery_id(self, endpoint_id: str, tx_hash: str, event_index: int = 0) -> str:
return f'{endpoint_id}:{tx_hash}:{event_index}'
async def deliver(self, endpoint_id: str, event: dict) -> dict:
tx_hash = event.get('tx_hash', '')
did = self._delivery_id(endpoint_id, tx_hash)
if await self.dedup_store.is_delivered(did):
return {'status': 'already_delivered', 'delivery_id': did}
await self.dedup_store.mark_in_progress(did)
try:
delivery = WebhookDelivery(self.router.endpoints.get(endpoint_id), event)
success = await delivery.send(self.router._client)
if success:
await self.dedup_store.mark_delivered(did)
return {'status': 'delivered', 'delivery_id': did}
else:
await self.dedup_store.mark_failed(did)
return {'status': 'failed', 'delivery_id': did, 'error': delivery.error or str(delivery.status_code)}
except Exception as e:
await self.dedup_store.mark_failed(did)
raise
6.3 幂等性保证
用户端点应利用 X-MSGChain-Delivery-Id 头实现幂等性。
"""
from fastapi import FastAPI, Request, HTTPException
app = FastAPI()
_processed: set[str] = set()
@app.post('/webhooks/msgchain')
async def handle_webhook(request: Request):
delivery_id = request.headers.get('X-MSGChain-Delivery-Id')
if not delivery_id:
raise HTTPException(400, 'Missing delivery ID')
if delivery_id in _processed:
return {'status': 'already_processed'}
body = await request.json()
signature = request.headers.get('X-MSGChain-Signature', '')
if not verify_signature(body, signature, 'whsec_...'):
raise HTTPException(401, 'Invalid signature')
_processed.add(delivery_id)
return {'status': 'ok', 'tx_hash': body.get('event', {}).get('tx_hash', '')}
def verify_signature(payload: dict, signature: str, secret: str) -> bool:
import hmac, hashlib, json
raw = json.dumps(payload, sort_keys=True).encode()
expected = hmac.new(secret.encode(), raw, hashlib.sha256).hexdigest()
return hmac.compare_digest(expected, signature)
"""
6.4 死信队列
当事件经过所有重试后仍无法投递成功,会被送入死信队列(DLQ)。
from enum import Enum
class DLQEntryStatus(Enum):
PENDING = 'pending'
PROCESSING = 'processing'
RESOLVED = 'resolved'
DISCARDED = 'discarded'
@dataclass
class DLQEntry:
"""死信队列条目。"""
id: str
endpoint_id: str
user_id: str
event: dict
error: str
attempts: int
last_attempt_at: str
status: DLQEntryStatus = DLQEntryStatus.PENDING
created_at: str = field(default_factory=lambda: datetime.utcnow().isoformat() + 'Z')
class DeadLetterQueue:
"""死信队列:持久化存储无法投递的事件。"""
def __init__(self):
self._entries: list[DLQEntry] = []
self._on_new_entry: Optional[Callable] = None
def on_new_entry(self, cb: Callable):
self._on_new_entry = cb
async def push(self, endpoint_id: str, user_id: str, event: dict, error: str, attempts: int) -> str:
import uuid
entry_id = f'dlq_{uuid.uuid4().hex[:16]}'
entry = DLQEntry(id=entry_id, endpoint_id=endpoint_id, user_id=user_id, event=event, error=error, attempts=attempts, last_attempt_at=datetime.utcnow().isoformat() + 'Z')
self._entries.append(entry)
return entry_id
async def retry(self, entry_id: str, router: WebhookRouter) -> bool:
"""重新投递死信队列中的事件。"""
entry = next((e for e in self._entries if e.id == entry_id), None)
if not entry:
raise ValueError(f'DLQ entry {entry_id} not found')
entry.status = DLQEntryStatus.PROCESSING
try:
delivery = WebhookDelivery(router.endpoints[entry.endpoint_id], entry.event)
success = await delivery.send(router._client)
entry.status = DLQEntryStatus.RESOLVED if success else DLQEntryStatus.PENDING
return success
except Exception:
entry.status = DLQEntryStatus.PENDING
raise
async def retry_all(self, router: WebhookRouter) -> dict:
"""重新投递所有待处理死信。"""
results = {'success': 0, 'failed': 0}
for entry in self._entries:
if entry.status != DLQEntryStatus.PENDING:
continue
try:
if await self.retry(entry.id, router):
results['success'] += 1
else:
results['failed'] += 1
except Exception:
results['failed'] += 1
return results
async def discard(self, entry_id: str):
entry = next((e for e in self._entries if e.id == entry_id), None)
if entry:
entry.status = DLQEntryStatus.DISCARDED
async def list_entries(self, status: Optional[DLQEntryStatus] = None, limit: int = 100) -> list[DLQEntry]:
if status:
return [e for e in self._entries if e.status == status][:limit]
return self._entries[:limit]
async def count(self) -> dict:
counts = {}
for s in DLQEntryStatus:
counts[s.value] = sum(1 for e in self._entries if e.status == s)
return counts
# 自动重试死信机制
async def dlq_retry_loop(dlq: DeadLetterQueue, router: WebhookRouter, interval: float = 300.0):
"""定期尝试重新投递死信队列中的事件。"""
while True:
await asyncio.sleep(interval)
entries = await dlq.list_entries(status=DLQEntryStatus.PENDING, limit=50)
for entry in entries:
try:
await dlq.retry(entry.id, router)
except Exception as e:
logger.error(f'DLQ retry error {entry.id}: {e}')
6.5 Redis 重试队列
使用 Redis Streams 实现可靠的重试队列。
class RedisRetryQueue:
"""基于 Redis Streams 的重试队列。"""
def __init__(self, redis_client, stream_key: str = 'webhook:retry'):
self.redis = redis_client
self.stream = stream_key
async def schedule(self, endpoint_id: str, event: dict, next_attempt: int, delay: float):
"""安排延迟重试。"""
import json
payload = {'endpoint_id': endpoint_id, 'event': json.dumps(event), 'attempt': str(next_attempt), 'scheduled_at': str(time.time() + delay)}
await self.redis.xadd(self.stream, payload)
async def send_to_dlq(self, endpoint_id: str, event: dict, attempts: int):
"""超过重试次数,送入死信队列。"""
import json
await self.redis.xadd('webhook:dlq', {'endpoint_id': endpoint_id, 'event': json.dumps(event), 'attempts': str(attempts), 'failed_at': str(time.time())})
7. 签名与认证
7.1 HMAC 签名
所有 Webhook 回调携带 X-MSGChain-Signature 头,用户端点应验证签名以确保事件确实来自 MSG Chain Webhook 服务。
import hashlib
import hmac
import json
from typing import Optional
class WebhookSigner:
"""Webhook 回调签名器。"""
def __init__(self, secret: str, algorithm: str = 'sha256'):
self.secret = secret
self.algo = algorithm
def sign(self, payload: bytes) -> str:
hash_fn = hashlib.sha256 if self.algo == 'sha256' else hashlib.sha512
return hmac.new(self.secret.encode('utf-8'), payload, hash_fn).hexdigest()
def sign_event(self, event: dict, timestamp: str) -> str:
"""对事件载荷签名。"""
raw = json.dumps({'event': event, 'timestamp': timestamp}, sort_keys=True, separators=(',', ':')).encode('utf-8')
return self.sign(raw)
class WebhookSignatureVerifier:
"""回调签名验证器(端点侧)。"""
def __init__(self, secret: str):
self.secret = secret
def verify(self, payload: bytes, signature: str) -> bool:
expected = hmac.new(self.secret.encode('utf-8'), payload, hashlib.sha256).hexdigest()
return hmac.compare_digest(expected, signature)
def verify_request(self, body: bytes, signature_header: str, timestamp_header: Optional[str] = None, tolerance: int = 300) -> bool:
"""验证完整的 Webhook 请求。"""
if not signature_header:
return False
sig_value = None
for part in signature_header.split(','):
part = part.strip()
if part.startswith('v1='):
sig_value = part[3:]
break
if not sig_value:
return False
if timestamp_header:
try:
ts = int(timestamp_header)
if abs(time.time() - ts) > tolerance:
return False
except ValueError:
return False
payload = f'{timestamp_header}.'.encode() + body
else:
payload = body
return self.verify(payload, sig_value)
7.2 JWT 认证
import jwt as pyjwt
from datetime import datetime, timedelta
from typing import Optional
JWT_SECRET = 'your-256-bit-secret'
JWT_ALGORITHM = 'HS256'
class JWTManager:
"""JWT 令牌管理器。"""
def __init__(self, secret: str = JWT_SECRET, algorithm: str = JWT_ALGORITHM):
self.secret = secret
self.algorithm = algorithm
def generate_token(self, user_id: str, permissions: Optional[list[str]] = None, expiry_hours: int = 24) -> str:
payload = {'sub': user_id, 'iat': datetime.utcnow(), 'exp': datetime.utcnow() + timedelta(hours=expiry_hours), 'permissions': permissions or ['webhook:read', 'webhook:write'], 'iss': 'webhook.msgchain.org'}
return pyjwt.encode(payload, self.secret, algorithm=self.algorithm)
def verify_token(self, token: str) -> Optional[dict]:
try:
return pyjwt.decode(token, self.secret, algorithms=[self.algorithm], issuer='webhook.msgchain.org')
except (pyjwt.ExpiredSignatureError, pyjwt.InvalidTokenError) as e:
logger.warning(f'JWT verification failed: {e}')
return None
def refresh_token(self, token: str) -> Optional[str]:
payload = self.verify_token(token)
return self.generate_token(payload['sub'], payload.get('permissions')) if payload else None
class ApiKeyManager:
"""API Key 管理器。"""
def __init__(self):
self.keys: dict[str, dict] = {}
def generate_key(self, user_id: str, label: str = 'default') -> str:
import uuid, hashlib
raw = f'mg_{uuid.uuid4().hex}'
hashed = hashlib.sha256(raw.encode()).hexdigest()
self.keys[hashed] = {'user_id': user_id, 'label': label, 'created_at': datetime.utcnow().isoformat() + 'Z', 'enabled': True}
return raw
def validate_key(self, api_key: str) -> Optional[dict]:
import hashlib
hashed = hashlib.sha256(api_key.encode()).hexdigest()
entry = self.keys.get(hashed)
return entry if entry and entry['enabled'] else None
def revoke_key(self, api_key: str):
import hashlib
hashed = hashlib.sha256(api_key.encode()).hexdigest()
if hashed in self.keys:
self.keys[hashed]['enabled'] = False
def list_keys(self, user_id: str) -> list[dict]:
return [{'label': v['label'], 'created_at': v['created_at'], 'enabled': v['enabled']} for v in self.keys.values() if v['user_id'] == user_id]
7.3 IP 白名单
import ipaddress
from typing import Optional
class IPWhitelist:
"""IP 白名单管理器。"""
MSGCHAIN_WEBHOOK_CIDRS = [
'35.188.0.0/16',
'34.64.0.0/16',
'10.0.0.0/8',
]
def __init__(self):
self.allowed = [ipaddress.ip_network(cidr) for cidr in self.MSGCHAIN_WEBHOOK_CIDRS]
self.user_allowed: dict[str, list] = {}
def add_user_cidr(self, user_id: str, cidr: str):
self.user_allowed.setdefault(user_id, []).append(ipaddress.ip_network(cidr))
def remove_user_cidr(self, user_id: str, cidr: str):
if user_id in self.user_allowed:
net = ipaddress.ip_network(cidr)
self.user_allowed[user_id] = [n for n in self.user_allowed[user_id] if n != net]
def is_allowed(self, ip_str: str, user_id: Optional[str] = None) -> bool:
try:
ip = ipaddress.ip_address(ip_str)
except ValueError:
return False
for net in self.allowed:
if ip in net:
return True
if user_id and user_id in self.user_allowed:
for net in self.user_allowed[user_id]:
if ip in net:
return True
return False
7.4 完整认证中间件
class WebhookAuthMiddleware:
"""API 认证中间件。"""
def __init__(self, jwt_mgr: JWTManager, api_key_mgr: ApiKeyManager, ip_whitelist: IPWhitelist):
self.jwt = jwt_mgr
self.api_key = api_key_mgr
self.ip_whitelist = ip_whitelist
async def authenticate(self, authorization: Optional[str], api_key: Optional[str], remote_ip: str) -> dict:
"""认证并返回用户信息。"""
if not self.ip_whitelist.is_allowed(remote_ip):
raise PermissionError(f'IP {remote_ip} not allowed')
if authorization and authorization.startswith('Bearer '):
token = authorization[7:]
payload = self.jwt.verify_token(token)
if payload:
return {'user_id': payload['sub'], 'permissions': payload.get('permissions', []), 'auth_method': 'jwt'}
if api_key:
entry = self.api_key.validate_key(api_key)
if entry:
return {'user_id': entry['user_id'], 'permissions': ['webhook:read', 'webhook:write'], 'auth_method': 'api_key'}
raise PermissionError('Authentication required')
8. 速率限制与节流
8.1 用户配额
每个用户注册的 Webhook 数量、每分钟投递次数受配额限制。
from dataclasses import dataclass, field
QUOTA_TIERS = {
'standard': {'max_webhooks': 10, 'max_deliveries_per_minute': 100, 'max_concurrent_deliveries': 10, 'max_retries': 5},
'pro': {'max_webhooks': 50, 'max_deliveries_per_minute': 1000, 'max_concurrent_deliveries': 50, 'max_retries': 10},
'enterprise': {'max_webhooks': 500, 'max_deliveries_per_minute': 10000, 'max_concurrent_deliveries': 200, 'max_retries': 20},
}
@dataclass
class UserQuota:
"""用户配额配置。"""
user_id: str
max_webhooks: int = 50
max_deliveries_per_minute: int = 1000
max_concurrent_deliveries: int = 100
max_body_size_bytes: int = 1024 * 512
max_retries: int = 10
tier: str = 'standard'
class QuotaManager:
"""用户配额管理器。"""
def __init__(self):
self.quotas: dict[str, UserQuota] = {}
def get_quota(self, user_id: str) -> UserQuota:
if user_id not in self.quotas:
self.quotas[user_id] = UserQuota(user_id=user_id)
return self.quotas[user_id]
def update_tier(self, user_id: str, tier: str):
if tier not in QUOTA_TIERS:
raise ValueError(f'Unknown tier: {tier}')
config = QUOTA_TIERS[tier]
quota = self.get_quota(user_id)
for k, v in config.items():
setattr(quota, k, v)
quota.tier = tier
def check_webhook_limit(self, user_id: str, current_count: int) -> bool:
return current_count < self.get_quota(user_id).max_webhooks
def check_rate_limit(self, user_id: str, current_minute_count: int) -> bool:
return current_minute_count < self.get_quota(user_id).max_deliveries_per_minute
8.2 速率限制实现
基于滑动窗口 + 令牌桶的速率限制器。
import time
from collections import defaultdict, deque
class SlidingWindowRateLimiter:
"""滑动窗口速率限制器。"""
def __init__(self, window_size: float = 60.0):
self.window_size = window_size
self.windows: dict[str, deque] = defaultdict(deque)
def _clean(self, key: str):
now = time.time()
w = self.windows[key]
while w and w[0] < now - self.window_size:
w.popleft()
def allow(self, key: str, max_calls: int) -> bool:
self._clean(key)
if len(self.windows[key]) >= max_calls:
return False
self.windows[key].append(time.time())
return True
def remaining(self, key: str, max_calls: int) -> int:
self._clean(key)
return max(0, max_calls - len(self.windows[key]))
class TokenBucketLimiter:
"""令牌桶算法:支持突发流量。"""
def __init__(self, rate: float, burst: int):
self.rate = rate
self.burst = burst
self.tokens = float(burst)
self.last_refill = time.monotonic()
def _refill(self):
now = time.monotonic()
elapsed = now - self.last_refill
self.tokens = min(float(self.burst), self.tokens + elapsed * self.rate)
self.last_refill = now
def allow(self, tokens: int = 1) -> bool:
self._refill()
if self.tokens >= tokens:
self.tokens -= tokens
return True
return False
def wait_time(self) -> float:
self._refill()
if self.tokens >= 1:
return 0.0
return (1.0 - self.tokens) / self.rate
class MultiLevelRateLimiter:
"""多层级速率限制:用户级 + 端点级 + 全局。"""
def __init__(self):
self.global_limiter = TokenBucketLimiter(rate=10000, burst=20000)
self.user_limiters: dict[str, TokenBucketLimiter] = {}
self.endpoint_limiters: dict[str, TokenBucketLimiter] = {}
def check(self, user_id: str, endpoint_id: str, user_rate: float, user_burst: int) -> bool:
if not self.global_limiter.allow():
return False
if user_id not in self.user_limiters:
self.user_limiters[user_id] = TokenBucketLimiter(user_rate, user_burst)
if not self.user_limiters[user_id].allow():
return False
if endpoint_id not in self.endpoint_limiters:
self.endpoint_limiters[endpoint_id] = TokenBucketLimiter(100, 200)
if not self.endpoint_limiters[endpoint_id].allow():
return False
return True
8.3 背压机制
当用户端点响应缓慢或 Webhook 服务负载过高时,背压机制自动降级。
from enum import Enum
from dataclasses import dataclass, field
class BackpressureMode(Enum):
NORMAL = 'normal'
THROTTLED = 'throttled'
DROPPING = 'dropping'
BLOCKED = 'blocked'
@dataclass
class BackpressureState:
"""背压状态。"""
mode: BackpressureMode = BackpressureMode.NORMAL
queue_depth: int = 0
latency_p99_ms: float = 0.0
error_rate: float = 0.0
last_changed: float = field(default_factory=time.time)
class BackpressureController:
"""背压控制器:根据系统负载自动调整投递行为。"""
def __init__(self, max_queue_depth: int = 10000, max_latency_p99: float = 5000.0, max_error_rate: float = 0.1):
self.max_queue_depth = max_queue_depth
self.max_latency_p99 = max_latency_p99
self.max_error_rate = max_error_rate
self.state = BackpressureState()
def update(self, queue_depth: int, latency_p99: float, error_rate: float):
self.state.queue_depth = queue_depth
self.state.latency_p99_ms = latency_p99
self.state.error_rate = error_rate
new_mode = self._compute_mode()
if new_mode != self.state.mode:
self.state.mode = new_mode
self.state.last_changed = time.time()
logger.warning(f'Backpressure mode changed to {new_mode.value}')
def _compute_mode(self) -> BackpressureMode:
if self.state.queue_depth > self.max_queue_depth:
return BackpressureMode.BLOCKED
if self.state.latency_p99_ms > self.max_latency_p99 or self.state.error_rate > self.max_error_rate * 2:
return BackpressureMode.DROPPING
if self.state.queue_depth > self.max_queue_depth * 0.7 or self.state.error_rate > self.max_error_rate:
return BackpressureMode.THROTTLED
return BackpressureMode.NORMAL
def should_deliver(self, priority: str = 'normal') -> bool:
mode = self.state.mode
if mode == BackpressureMode.BLOCKED:
return False
if mode == BackpressureMode.DROPPING and priority == 'low':
return False
if mode == BackpressureMode.THROTTLED:
import random
return random.random() < 0.5
return True
def get_sampling_rate(self) -> float:
return {BackpressureMode.NORMAL: 1.0, BackpressureMode.THROTTLED: 0.5, BackpressureMode.DROPPING: 0.2, BackpressureMode.BLOCKED: 0.0}[self.state.mode]
def get_delay_multiplier(self) -> float:
return {BackpressureMode.NORMAL: 1.0, BackpressureMode.THROTTLED: 2.0, BackpressureMode.DROPPING: 5.0, BackpressureMode.BLOCKED: float('inf')}[self.state.mode]
9. 监控与日志
9.1 投递延迟与成功率
from dataclasses import dataclass, field
from collections import defaultdict
from datetime import datetime, timedelta
from typing import Optional
@dataclass
class DeliveryMetric:
"""单次投递指标。"""
endpoint_id: str
user_id: str
status: str
latency_ms: float
attempt: int
timestamp: str
status_code: Optional[int] = None
error: Optional[str] = None
class MetricsCollector:
"""投递指标收集器。"""
def __init__(self, window_size: int = 10000):
self.metrics: list[DeliveryMetric] = []
self.window_size = window_size
def record(self, metric: DeliveryMetric):
self.metrics.append(metric)
if len(self.metrics) > self.window_size:
self.metrics = self.metrics[-self.window_size:]
def success_rate(self, minutes: int = 5) -> float:
recent = self._recent(minutes)
if not recent:
return 1.0
return sum(1 for m in recent if m.status == 'success') / len(recent)
def latency_percentile(self, percentile: float, minutes: int = 5) -> float:
recent = self._recent(minutes)
latencies = sorted([m.latency_ms for m in recent if m.status == 'success'])
if not latencies:
return 0.0
idx = int(len(latencies) * percentile / 100)
return latencies[min(idx, len(latencies) - 1)]
def error_breakdown(self, minutes: int = 60) -> dict:
recent = self._recent(minutes)
breakdown: dict[str, int] = defaultdict(int)
for m in recent:
if m.status == 'failed':
breakdown[m.error or f'http_{m.status_code or 0}'] += 1
return dict(breakdown)
def user_stats(self, user_id: str, minutes: int = 60) -> dict:
user_metrics = [m for m in self._recent(minutes) if m.user_id == user_id]
total = len(user_metrics)
successes = sum(1 for m in user_metrics if m.status == 'success')
latencies = [m.latency_ms for m in user_metrics if m.status == 'success']
return {'total': total, 'success': successes, 'failed': total - successes, 'success_rate': successes / max(total, 1), 'avg_latency_ms': sum(latencies) / max(len(latencies), 1)}
def _recent(self, minutes: int) -> list[DeliveryMetric]:
cutoff = (datetime.utcnow() - timedelta(minutes=minutes)).isoformat()
return [m for m in self.metrics if m.timestamp >= cutoff]
class PrometheusMetricsExporter:
"""Prometheus 指标导出。"""
def __init__(self, namespace: str = 'msgchain_webhook'):
self.namespace = namespace
self._enabled = False
try:
from prometheus_client import Counter, Histogram, Gauge
self.deliveries_total = Counter(f'{namespace}_deliveries_total', 'Total webhook deliveries', ['endpoint_id', 'status'])
self.delivery_latency = Histogram(f'{namespace}_delivery_latency_ms', 'Delivery latency', ['endpoint_id'], buckets=[5, 10, 25, 50, 100, 250, 500, 1000, 2500, 5000])
self.queue_depth = Gauge(f'{namespace}_queue_depth', 'Event queue depth')
self.active_endpoints = Gauge(f'{namespace}_active_endpoints', 'Active endpoints')
self.dlq_size = Gauge(f'{namespace}_dlq_size', 'Dead letter queue size')
self._enabled = True
except ImportError:
pass
def observe_delivery(self, endpoint_id: str, status: str, latency_ms: float):
if self._enabled:
self.deliveries_total.labels(endpoint_id=endpoint_id, status=status).inc()
self.delivery_latency.labels(endpoint_id=endpoint_id).observe(latency_ms)
9.2 结构化日志
import json
import logging
import sys
from datetime import datetime
class StructuredFormatter(logging.Formatter):
"""结构化 JSON 日志格式化器。"""
def format(self, record: logging.LogRecord) -> str:
log_entry = {'timestamp': datetime.utcnow().isoformat() + 'Z', 'level': record.levelname, 'logger': record.name, 'message': record.getMessage()}
if hasattr(record, 'props'):
log_entry['props'] = record.props
if record.exc_info and record.exc_info[0]:
log_entry['exception'] = {'type': record.exc_info[0].__name__, 'message': str(record.exc_info[1])}
return json.dumps(log_entry, ensure_ascii=False)
def setup_logging(level: str = 'INFO', json_format: bool = True):
handler = logging.StreamHandler(sys.stdout)
if json_format:
handler.setFormatter(StructuredFormatter())
root = logging.getLogger()
root.setLevel(getattr(logging, level.upper(), logging.INFO))
root.handlers.clear()
root.addHandler(handler)
class DeliveryLogger:
"""投递事件日志器。"""
def __init__(self, logger_name: str = 'webhook.delivery'):
self.logger = logging.getLogger(logger_name)
def log_delivery(self, endpoint_id: str, user_id: str, event: dict, attempt: int, status: str, latency_ms: float, status_code=None, error=None):
extra = {'props': {'endpoint_id': endpoint_id, 'user_id': user_id, 'tx_hash': event.get('tx_hash', ''), 'height': event.get('height', 0), 'attempt': attempt, 'latency_ms': round(latency_ms, 2), 'status_code': status_code, 'error': error}}
if status == 'success':
self.logger.info(f'Delivered to {endpoint_id}', extra=extra)
else:
self.logger.warning(f'Failed delivery to {endpoint_id} (attempt {attempt}): {error or status_code}', extra=extra)
9.3 告警规则
ALERT_RULES = {
'high_failure_rate': {'condition': 'delivery_success_rate < 0.95 for 5m', 'severity': 'warning', 'summary': 'Webhook 投递成功率低于 95%'},
'critical_failure_rate': {'condition': 'delivery_success_rate < 0.80 for 5m', 'severity': 'critical', 'summary': 'Webhook 投递成功率低于 80%'},
'high_latency': {'condition': 'delivery_latency_p99 > 5000ms for 5m', 'severity': 'warning', 'summary': '投递延迟异常'},
'dlq_growth': {'condition': 'dlq_size > 100', 'severity': 'warning', 'summary': '死信队列堆积'},
'queue_backlog': {'condition': 'queue_depth > 5000 for 1m', 'severity': 'critical', 'summary': '事件队列堆积'},
'endpoint_down': {'condition': 'up{job="webhook"} == 0', 'severity': 'critical', 'summary': 'Webhook 服务实例宕机'},
}
10. MSG Chain 集成
10.1 事件索引
Webhook 服务需要从 MSG Chain 节点索引事件。除了直接订阅 WebSocket,还可以通过 REST API 进行历史事件补录。
class MsgChainIndexer:
"""MSG Chain 事件索引器。"""
def __init__(self, base_url: str = 'https://api.msgchain.org'):
self.base_url = base_url
self._client = httpx.AsyncClient(base_url=base_url, timeout=30.0)
async def get_block_events(self, height: int) -> list[dict]:
"""获取指定区块的所有事件。"""
response = await self._client.get(f'/cosmos/tx/v1beta1/txs/block/{height}')
response.raise_for_status()
data = response.json()
events = []
for tx in data.get('txs', []):
tx_hash = tx.get('txhash', '')
tx_response = tx.get('tx_response', {})
for event in tx_response.get('events', []):
events.append({'tx_hash': tx_hash, 'height': height, 'type': event.get('type', ''), 'attributes': self._normalize_attrs(event.get('attributes', []))})
return events
async def index_range(self, from_height: int, to_height: int, concurrency: int = 10) -> list[dict]:
"""批量索引区块范围的事件。"""
sem = asyncio.Semaphore(concurrency)
async def fetch(h):
async with sem:
try:
return await self.get_block_events(h)
except Exception as e:
logger.error(f'Failed to index block {h}: {e}')
return []
tasks = [fetch(h) for h in range(from_height, to_height + 1)]
results = await asyncio.gather(*tasks)
return [ev for batch in results for ev in batch]
def _normalize_attrs(self, raw_attrs: list[dict]) -> dict:
attrs = {}
for attr in raw_attrs:
key = attr.get('key', '')
value = attr.get('value', '')
if isinstance(key, bytes): key = key.decode()
if isinstance(value, bytes): value = value.decode()
attrs[key] = value
return attrs
async def close(self):
await self._client.aclose()
10.2 GraphQL Subscription
Webhook 服务可以将事件发布到 GraphQL 端点,支持客户端通过 Subscription 实时消费。
class GraphQLSubscriptionPublisher:
"""将事件发布到 GraphQL Subscription。"""
def __init__(self):
self.subscribers: dict[str, list[asyncio.Queue]] = {}
def subscribe(self, sub_id: str) -> asyncio.Queue:
queue: asyncio.Queue = asyncio.Queue(maxsize=100)
self.subscribers.setdefault(sub_id, []).append(queue)
return queue
async def publish(self, event: dict):
msg = {'data': {'txHash': event.get('tx_hash', ''), 'height': event.get('height', 0), 'chainId': 'msg-chain-1', 'timestamp': datetime.utcnow().isoformat() + 'Z'}}
for queues in self.subscribers.values():
for q in list(queues):
try:
q.put_nowait(msg)
except asyncio.QueueFull:
pass
GraphQL Schema 定义:
type Event {
txHash: String!
height: Int!
type: String!
attributes: JSON!
chainId: String!
timestamp: String!
}
type Subscription {
events(types: [String!], contractAddress: String, agentId: String): Event!
deliveryStatus(endpointIds: [String!]): WebhookDelivery!
}
type Query {
deliveryLogs(endpointId: String!, limit: Int = 100, offset: Int = 0): [WebhookDelivery!]!
deliveryStats(endpointId: String!, since: String!, until: String!): DeliveryStats!
}
type DeliveryStats {
total: Int! success: Int! failed: Int! avgLatencyMs: Float! p95LatencyMs: Float!
}
10.3 WebSocket 桥接
除了标准 HTTP 回调,服务支持 WebSocket 桥接模式。
import websockets
class WebSocketBridge:
"""WebSocket 桥接:通过 WSS 推送事件。"""
def __init__(self, host: str = '0.0.0.0', port: int = 9500):
self.host = host
self.port = port
self.connections: dict[str, set] = {}
self._server = None
async def start(self):
self._server = await websockets.serve(self._handler, self.host, self.port, ping_interval=30, ping_timeout=10)
logger.info(f'WS Bridge listening on {self.host}:{self.port}')
async def stop(self):
if self._server:
self._server.close()
await self._server.wait_closed()
async def _handler(self, ws):
user_id = None
try:
auth_msg = await asyncio.wait_for(ws.recv(), timeout=10)
auth = json.loads(auth_msg)
user_id = auth.get('user_id', '')
if not user_id: return
self.connections.setdefault(user_id, set()).add(ws)
await ws.send(json.dumps({'type': 'connected', 'user_id': user_id}))
async for msg in ws:
data = json.loads(msg)
if data.get('type') == 'ping':
await ws.send(json.dumps({'type': 'pong'}))
except Exception:
pass
finally:
if user_id:
self.connections.get(user_id, set()).discard(ws)
async def broadcast(self, event: dict, user_id: str = None):
"""向特定用户或所有连接的客户端广播事件。"""
message = json.dumps({'type': 'event', 'data': event, 'chain_id': 'msg-chain-1'})
targets = self.connections.get(user_id, set()) if user_id else set().union(*self.connections.values())
for ws in list(targets):
try:
await ws.send(message)
except websockets.ConnectionClosed:
pass
10.4 GraphQL Mutation API
Webhook 服务提供完整的 GraphQL Mutation 用于管理配置。
type Mutation {
createWebhook(input: WebhookInput!): WebhookPayload!
updateWebhook(id: String!, input: WebhookInput!): WebhookPayload!
deleteWebhook(id: String!): DeletePayload!
testWebhook(id: String!): TestResult!
rotateSecret(id: String!): SecretPayload!
retryDLQ(entryId: String!): RetryPayload!
}
input WebhookInput {
name: String! url: String! condition: String secret: String
retryMax: Int timeout: Float enabled: Boolean headers: JSON
}
type WebhookPayload { id: String! name: String! url: String! status: String! createdAt: String! }
type TestResult { webhookId: String! matched: Boolean! event: JSON! message: String! }
type SecretPayload { webhookId: String! secret: String! message: String! }
11. 部署实践
11.1 Webhook 网关架构
生产环境推荐的多层网关架构:
+-----------------+
| 负载均衡器 |
| (Nginx / LB) |
+--------+--------+
|
+----------------+----------------+
v v v
+--------------+ +--------------+ +--------------+
| Webhook GW 1 | | Webhook GW 2 | | Webhook GW 3 |
| (stateless) | | (stateless) | | (stateless) |
+------+-------+ +------+-------+ +------+-------+
| | |
+----------------+----------------+
|
+---------------+----------------+
| | |
v v v
+--------------+ +--------------+ +----------------+
| Redis | | PostgreSQL | | 事件索引器 |
| (重试队列+ | | (规则存储+ | | (Indexer) |
| 去重+缓存) | | 配置+日志) | | |
+--------------+ +--------------+ +----------------+
|
v
+----------------+
| MSG Chain |
| RPC/WSS 节点 |
+----------------+
DOCKER_COMPOSE_CONFIG = '''
version: '3.8'
services:
webhook-gateway:
image: msgchain/webhook-gateway:latest
ports:
- "8080:8080"
environment:
- REDIS_URL=redis://redis:6379
- DATABASE_URL=postgresql://user:pass@postgres:5432/webhooks
- CHAIN_RPC_WS=wss://rpc.msgchain.org/websocket
- CHAIN_ID=msg-chain-1
- LOG_LEVEL=INFO
depends_on:
- redis
- postgres
deploy:
replicas: 3
resources:
limits:
cpus: '2'
memory: 4G
redis:
image: redis:7-alpine
volumes:
- redis_data:/data
postgres:
image: postgres:15-alpine
environment:
- POSTGRES_DB=webhooks
- POSTGRES_USER=webhook
- POSTGRES_PASSWORD=changeme
volumes:
- postgres_data:/var/lib/postgresql/data
volumes:
redis_data:
postgres_data:
'''
11.2 水平扩展策略
class HorizontalScaler:
"""水平扩展策略。"""
SCALING_POLICIES = {
'cpu_based': {'metric': 'cpu_utilization', 'target': 0.7, 'min_replicas': 2, 'max_replicas': 20},
'queue_based': {'metric': 'queue_depth', 'target': 1000, 'min_replicas': 2, 'max_replicas': 20},
'delivery_rate': {'metric': 'deliveries_per_second', 'target': 500, 'min_replicas': 2, 'max_replicas': 30},
}
@staticmethod
def consistent_hash(user_id: str, num_buckets: int) -> int:
import hashlib
return int(hashlib.md5(user_id.encode()).hexdigest()[:8], 16) % num_buckets
11.3 K8s 部署
apiVersion: apps/v1
kind: Deployment
metadata:
name: msgchain-webhook
namespace: webhook-system
spec:
replicas: 3
selector:
matchLabels:
app: msgchain-webhook
template:
metadata:
labels:
app: msgchain-webhook
annotations:
prometheus.io/scrape: "true"
prometheus.io/port: "8080"
spec:
affinity:
podAntiAffinity:
preferredDuringSchedulingIgnoredDuringExecution:
- weight: 100
podAffinityTerm:
labelSelector:
matchExpressions:
- key: app
operator: In
values:
- msgchain-webhook
topologyKey: kubernetes.io/hostname
containers:
- name: webhook
image: msgchain/webhook-gateway:latest
ports:
- containerPort: 8080
name: http
- containerPort: 9500
name: ws
env:
- name: CHAIN_ID
value: "msg-chain-1"
- name: CHAIN_RPC_WS
value: "wss://rpc.msgchain.org/websocket"
- name: REDIS_URL
valueFrom:
secretKeyRef:
name: webhook-secrets
key: redis-url
- name: DATABASE_URL
valueFrom:
secretKeyRef:
name: webhook-secrets
key: database-url
resources:
requests:
cpu: "500m"
memory: "512Mi"
limits:
cpu: "2000m"
memory: "2Gi"
livenessProbe:
httpGet:
path: /health
port: 8080
initialDelaySeconds: 10
periodSeconds: 15
readinessProbe:
httpGet:
path: /ready
port: 8080
initialDelaySeconds: 5
periodSeconds: 10
---
apiVersion: v1
kind: Service
metadata:
name: msgchain-webhook
namespace: webhook-system
spec:
selector:
app: msgchain-webhook
ports:
- name: http
port: 8080
targetPort: 8080
- name: ws
port: 9500
targetPort: 9500
type: ClusterIP
---
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
name: msgchain-webhook-hpa
namespace: webhook-system
spec:
scaleTargetRef:
apiVersion: apps/v1
kind: Deployment
name: msgchain-webhook
minReplicas: 2
maxReplicas: 20
metrics:
- type: Resource
resource:
name: cpu
target:
type: Utilization
averageUtilization: 70
- type: Pods
pods:
metric:
name: webhook_queue_depth
target:
type: AverageValue
averageValue: 1000
---
apiVersion: networking.k8s.io/v1
kind: Ingress
metadata:
name: msgchain-webhook-ingress
namespace: webhook-system
annotations:
kubernetes.io/ingress.class: nginx
nginx.ingress.kubernetes.io/ssl-redirect: "true"
cert-manager.io/cluster-issuer: letsencrypt-prod
spec:
tls:
- hosts:
- webhook.msgchain.org
secretName: webhook-tls
rules:
- host: webhook.msgchain.org
http:
paths:
- path: /
pathType: Prefix
backend:
service:
name: msgchain-webhook
port:
number: 8080
11.4 配置管理
import os
from typing import Optional
class WebhookServiceConfig:
"""Webhook 服务配置。"""
def __init__(self):
self.chain_id = os.getenv('CHAIN_ID', 'msg-chain-1')
self.rpc_ws = os.getenv('CHAIN_RPC_WS', 'wss://rpc.msgchain.org/websocket')
self.rest_api = os.getenv('CHAIN_REST_API', 'https://api.msgchain.org')
self.host = os.getenv('HOST', '0.0.0.0')
self.port = int(os.getenv('PORT', '8080'))
self.ws_bridge_port = int(os.getenv('WS_BRIDGE_PORT', '9500'))
self.log_level = os.getenv('LOG_LEVEL', 'INFO')
self.database_url = os.getenv('DATABASE_URL', '')
self.redis_url = os.getenv('REDIS_URL', '')
self.jwt_secret = os.getenv('JWT_SECRET', '')
self.jwt_algorithm = os.getenv('JWT_ALGORITHM', 'HS256')
self.global_rate = float(os.getenv('GLOBAL_RATE', '10000'))
self.global_burst = int(os.getenv('GLOBAL_BURST', '20000'))
self.default_user_rate = float(os.getenv('DEFAULT_USER_RATE', '100'))
self.default_user_burst = int(os.getenv('DEFAULT_USER_BURST', '200'))
self.default_retry_max = int(os.getenv('DEFAULT_RETRY_MAX', '5'))
self.default_retry_base_delay = float(os.getenv('DEFAULT_RETRY_BASE_DELAY', '1.0'))
self.default_retry_max_delay = float(os.getenv('DEFAULT_RETRY_MAX_DELAY', '60.0'))
self.metrics_enabled = os.getenv('METRICS_ENABLED', 'true').lower() == 'true'
def validate(self):
errors = []
if not self.jwt_secret: errors.append('JWT_SECRET required')
if not self.database_url: errors.append('DATABASE_URL required')
if errors:
raise ValueError(f'Config errors: {", ".join(errors)}')
11.5 健康检查与优雅关闭
import signal
import time
class HealthChecker:
"""健康检查端点。"""
def __init__(self):
self.start_time = time.time()
async def check_liveness(self) -> dict:
return {'status': 'ok', 'uptime_seconds': time.time() - self.start_time, 'timestamp': datetime.utcnow().isoformat() + 'Z'}
async def check_readiness(self) -> dict:
return {'status': 'ready'}
class GracefulShutdown:
"""优雅关闭管理器。"""
def __init__(self, services: list):
self.services = services
self._shutdown_requested = False
def install_signal_handlers(self):
loop = asyncio.get_event_loop()
for sig in (signal.SIGTERM, signal.SIGINT):
loop.add_signal_handler(sig, lambda: asyncio.create_task(self.shutdown()))
async def shutdown(self):
if self._shutdown_requested:
return
self._shutdown_requested = True
logger.info('Initiating graceful shutdown...')
tasks = [s.stop() for s in self.services if hasattr(s, 'stop')]
tasks += [s.close() for s in self.services if hasattr(s, 'close')]
if tasks:
await asyncio.gather(*tasks, return_exceptions=True)
logger.info('Graceful shutdown complete')
12. 总结
12.1 架构要点回顾
Webhook 事件推送与通知服务是 MSG Chain 生态中连接链上事件与下游消费者的核心基础设施。
| 模块 | 核心设计 | 生产建议 |
|---|---|---|
| 事件监听 | Tendermint WSS 订阅 + 缓冲队列 + 检查点恢复 | 多节点冗余订阅,避免单点故障 |
| 过滤引擎 | 地址/类型/参数多维匹配 + 条件表达式 DSL | 编译缓存,LRU 淘汰策略 |
| 规则引擎 | 静态+动态规则,DSL 支持热更新 | 规则存储使用 PostgreSQL,Redis 缓存活跃规则 |
| 可靠投递 | At-Least-Once + 幂等性 + 死信队列 | 端点需实现幂等性接口 |
| 签名认证 | HMAC-SHA256 + JWT + API Key + IP 白名单 | 密钥定期轮换,审计日志 |
| 速率限制 | 令牌桶 + 滑动窗口 + 多层级配额 | 用户级+全局双重限制 |
| 监控告警 | Prometheus + Grafana + 结构化日志 | 成功率<95%告警,延迟 P99<5s |
| 部署 | 无状态 + K8s HPA + 多副本 | 跨 AZ 部署,避免单可用区故障 |
12.2 与 AI Agent 事件驱动文档的区别
本文档聚焦于 Webhook 即服务基础设施,与 AI Agent 事件驱动指南定位不同:
| 维度 | AI Agent 事件驱动指南 | Webhook 服务指南(本文) |
|---|---|---|
| 视角 | Agent 开发者(事件消费者) | 平台开发者/运维(事件分发者) |
| 核心 | Agent 如何监听和响应事件 | 如何构建可靠事件投递服务 |
| 关注点 | 业务逻辑、Agent 决策 | 基础设施、可靠性、扩展性 |
| 关键能力 | WebSocket 客户端、事件处理 | 过滤引擎、重试策略、DLQ、限流 |
| 目标用户 | Agent 开发者 | 平台团队、DevOps |
12.3 生产就绪检查清单
[ ] 事件接收
[ ] Tendermint WSS 生产连接 (wss://rpc.msgchain.org/websocket)
[ ] 多节点冗余订阅
[ ] 检查点持久化(每 100 区块)
[ ] 事件缓冲队列(有界,背压感知)
[ ] 过滤匹配
[ ] 条件表达式编译缓存
[ ] 过滤器性能基准(< 1ms/规则)
[ ] 通配符和正则匹配限制
[ ] 回调投递
[ ] HTTP 连接池优化(keepalive, 最大连接数)
[ ] 幂等性投递 ID (X-MSGChain-Delivery-Id)
[ ] 签名验证 (X-MSGChain-Signature)
[ ] 重试策略(指数退避 + 抖动)
[ ] 死信队列持久化
[ ] 认证安全
[ ] API Key / JWT 认证
[ ] IP 白名单
[ ] 密钥定期轮换策略
[ ] TLS 加密(全链路)
[ ] 限流保护
[ ] 用户级配额(Webhook 数量、投递频率)
[ ] 全局速率限制
[ ] 背压降级机制
[ ] 监控告警
[ ] 投递成功率 (SLI: >99%)
[ ] 投递延迟 P50/P95/P99 (SLI: P99 < 5s)
[ ] 死信队列大小告警
[ ] 服务实例健康检查
[ ] 结构化 JSON 日志
[ ] 部署运维
[ ] 无状态部署(水平可扩展)
[ ] K8s HPA 配置(基于 CPU + 队列深度)
[ ] 跨可用区部署
[ ] 蓝绿部署 / 滚动更新
[ ] 灾备方案(多区域)
12.4 参考资源
- 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/ - Webhook 最佳实践 (Stripe):
https://stripe.com/docs/webhooks - Webhook 安全指南 (Svix):
https://docs.svix.com/ - CosmWasm 合约开发:
https://docs.cosmwasm.com - Prometheus 监控:
https://prometheus.io/docs/introduction/overview/
本指南由 AI 生成,代码示例仅供参考。生产环境部署前请充分测试。端点和 API 参数以 MSG Chain 官方文档为准。
