dApp Docs/Webhook 事件推送与通知服务指南
Development reference. Not independently verified for production.

Webhook 事件推送与通知服务指南

数据来源:MSG Chain 代码库核实

主网状态: No-Go — 当前 MSGChain 主网裁决为 No-Go,以下内容反映代码实际状态,不代表生产可用。

1. Webhook 在区块链通知中的角色

1.1 为什么需要 Webhook

Webhook 是一种用户定义的 HTTP 回调机制,当链上特定事件发生时,服务端主动向用户注册的端点推送事件数据。与轮询(Polling)模式相比,Webhook 在区块链通知场景中具有根本性优势:

轮询的问题:

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 适用场景

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


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