ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

fhEVM Listener Notifier Library 深度解析:面向 zama 组件的 RabbitMQ 事件消费与日志通知库

2026/9/13 15:33:30 拓冰建站 浏览量
fhEVM Listener Notifier Library 深度解析:面向 zama 组件的 RabbitMQ 事件消费与日志通知库 fhEVM Listener Notifier Library 深度解析面向 zama 组件的 RabbitMQ 事件消费与日志通知库【免费下载链接】fhevmFHEVM, a full-stack framework for integrating Fully Homomorphic Encryption (FHE) with blockchain applications项目地址: https://gitcode.com/GitHub_Trending/fh/fhevm导读listener/docs/library_notifier.md描述的 Notifier Library通知器库是 fhEVM Listener 组件体系中的关键一环它以一个独立 Rust 库的形式消费 Listener Core 广播到 RabbitMQ 的区块、交易与回执消息按数据库中的过滤器filter与 ABI 进行匹配再把命中的事件log转发给 zama 各内部组件如 Gateway、Relayer 等从而以订阅-消费的方式模拟链上 poller 或 websocket 流。本文以该设计文档为骨架结合 listener 目录下 broker 抽象、事件契约、Postgres 迁移与过滤器管理等实际源码逐层拆解其目标、运行逻辑、功能清单、表结构设计与实战接入方式帮助读者理解从链事件到组件内部逻辑的完整管线。一、设计目标让组件像订阅事件流一样消费链上日志原文档开门见山地给出了 Notifier Library 的核心目标GoalIntegrate event consuming for specific logs and filters for zama components, and mimic the behaviour of a poller or a websocket stream, but from the listener standalone component.翻译成工程语言就是为 zama 各组件提供按特定日志logs与过滤器filters消费事件的能力让它们体验到类似轮询器poller或 websocket 流的实时性但底层并不直接连链而是从 Listener 这一独立组件中间接获取数据。这个定位带来两个关键好处解耦组件与链节点Gateway、Relayer 等组件不再各自维护 RPC 长连接或轮询逻辑统一由 Listener 负责区块抓取、重组reorg检测与广播见 listener_core.md组件侧只负责消费自己关心的事件。弹性扩展HPA友好由于事件来自 RabbitMQ 消息队列而非有状态的长连接水平伸缩副本Horizontal Pod Autoscaler不会破坏消费语义——多个副本竞争同一队列即可这正是原文档在通用特性中强调的Resilient to HPA: (Ok by design: consuming from rmq)的含义。二、核心运行逻辑一个以回执为核心的解析器Receipt Parser文档用一段话概括了组件的数据流LogicThis component consume blocks, transactions, receipts from the different queues declared on rabbitmq, checks into its table to see if there is relevant filters or abi to watch, or even from or to sources if we need to watch for some transactions, register logs in a table, and forward to the internal logic of the components that need logs.可归纳为如下处理链消费从 RabbitMQ 上声明的多条队列中分别消费区块blocks、交易transactions、回执receipts消息匹配查询本地数据表判断是否存在需要关注的过滤器filters或 ABI以及是否需要按交易的from/to来源做定向监控登记将命中的日志logs写入数据表转发把日志交给内部需要它的组件逻辑继续处理。文档随后点明本质Basically, the library is receipt parser.——整个库本质上是一个回执解析器。链上事件全部以 log 形式存在于交易回执中对应listener/docs/listener_core.md中 Transaction receipts contains all the logs 这条常识所以解析回执并按 watcher 规则匹配就能不漏掉任何一个事件。原文档还注明库的匹配逻辑存在一个已有的参考实现可供内部咨询For inspiration regarding this library, There is an existing implementation for logic, ask for access if needed.新实现可在其基础上演进。三、功能特性全景3.1 通用特性GeneralRust 库以 crate 形式发布可被各组件以依赖方式接入。HPA 弹性通过从 RabbitMQ 消费天然获得无需额外状态同步。清晰的订阅/消费 API对外提供subscribe与consume事件流的清晰接口。钩子hooks / handlers消费函数支持传入回调处理器触发 zama 内部逻辑要求易于集成到现有组件。策略模式预留Solana为未来接入 Solana 预留策略模式strategy pattern扩展点——这一点与 Listener Core 中strategy patternhandling chains that doesnt supporteth_getBlockReceiptsmethod, and solana later的设计一脉相承。3.2 队列消费Queue consumer在 Listener 与 Notifier 库之间共享同一个 RabbitMQ 库对应仓库中的 shared/broker crate并自带重试队列等可靠性设施。从仓库源码看shared/broker/README.md 提供了一套同时支持Redis Streams 与 RabbitMQ的统一 broker 抽象AMQP 模式下自动声明main/retry/dlx三个共享交换机消息先进入主队列处理失败后经重试交换机带 TTL 延迟返回主队列重试超过重试次数上限后进入死信队列dead-letter。这正对应文档要求的retry queues。3.3 数据存储Data storage采用RDS 数据库可被不同组件复用使用Postgres 表存储日志logs与各类 watcher。仓库中 listener_core/migrations/20260224175428_init.sql 与 20260827120000_add_filter_type_and_final_blocks.sql 提供了实际落地证据详见第五节。3.4 区块链相关最小功能集Minimal features持久化区块高度并对失败模式具备弹性不因进程重启或队列抖动丢失进度多链设计为每条链分别消费多类队列blocks、transactions、receipts → logs全量消费即便某事件当前没有使用者也必须消费掉Should consume all events even if they are not used (or rmq memory will grow)否则 RabbitMQ 内存会持续增长——这是防止队列积压的硬性要求动态 ABI 通知器声明将 log watcher 类型存入 Postgres将 watched logs 存入 Postgres支持运行时动态声明新的 watcher。3.5 进阶Should have功能多 watcher 与确认块数可为不同 watcher 声明不同的区块确认数block confirmations例如按finality、safe或N confirmations来区分交付语义若需要向 RPC 查询确认状态则可能要求配置 RPC URL。仓库中 event.rs 的BlockFlow枚举Live/Reorged/Catchup/Final/FinalCatchup与 routing.rs 中的final-event、final-catchup-event路由键正是最终性交付流的实现印证感知新链能够感知 RabbitMQ 上出现的新链命名空间多种 watcher 类型logs 与 tx 两类 watcher可选取消被重组的事件Cancel reorged events可选回放历史区块Replay past blocks——文档同时指出由于消息队列本身具有缓冲能力通常不需要回放去重检查重复日志问题必要时可对日志施加**语义哈希semantic hash**保证唯一性由 zama 内部按需处理Metrics 与 Alerting提供指标与告警能力。仓库中 listener_core/src/metrics.rs 与 broker 的 metrics.rs 可视为该能力的落点。3.6 示例Example要求实现一个最小可运行示例演示如何使用该库。仓库中 listener/crates/example crate 提供了main.rs、live_events.rs、final_events.rs、full_block.rs、transfer.rs等多个示例源文件分别覆盖实时事件、最终性事件、全区块与转账场景可直接作为接入模板。四、Postgres 表结构提案watcher 与 logs原文档在 Postgres table 一节给出了两张核心表的设计提案表 1watcher| 字段 | 说明 | | --- | --- | |uuid| 主键 | |chainId| 所属链 ID | | 确认块数conf block | 该 watcher 需要的区块确认数 | |ABI| 关注的合约 ABI | | watcher 类型 |tx或contract|表 2logs| 字段 | 说明 | | --- | --- | |uuid| 主键 | |watcher_uuid| 关联的 watcher | |block_number| 日志所在区块高度 | |released| 是否已释放TRUE / FALSE用于标记该日志是否已转发给内部逻辑 | |log| 日志内容可反序列化或原始存储 | |UNCLE?| 是否为叔块日志非必选——若依赖区块确认数机制则可不记录 |这两个released布尔位与UNCLE标记的设计意图值得展开released用于区分已消费与待消费的日志是幂等交付与断点续传的基础UNCLE对应重组后被废弃的叔块事件文档明确表示如果不做重放/取消且依赖区块确认数则此列非必须——因为等待足够确认数后叔块事件根本不会被释放。五、仓库中的实际落地filters、blocks 与 final_blocks将设计文档与当前仓库源码对照可以看到提案在 Listener 中已具体化为如下迁移与模型注意以下以仓库实际实现为准与文档提案存在命名上的演进5.1 filters 表对应 watcher 表20260224175428_init.sql 创建了filters表CREATE TABLE IF NOT EXISTS filters ( id UUID PRIMARY KEY, chain_id BIGINT NOT NULL, consumer_id VARCHAR(128) NOT NULL, from VARCHAR(42), to VARCHAR(42), log_address VARCHAR(42), created_at TIMESTAMPTZ NOT NULL DEFAULT NOW() ); CREATE UNIQUE INDEX IF NOT EXISTS idx_filters_unique_chain_consumer_from_to ON filters(chain_id, consumer_id, COALESCE(from, ), COALESCE(to, ), COALESCE(log_address, ));每个过滤器以(chain_id, consumer_id, from, to, log_address)唯一标识consumer_id即消费方身份如 gateway与文档中watcher 类型 chainId的提案一一对应COALESCE(..., )将 NULL 统一映射为空串使 NULL 参与唯一性约束兼容 PostgreSQL 15 之前没有NULLS NOT DISTINCT的版本后续迁移 20260827120000_add_filter_type_and_final_blocks.sql 进一步新增filter_type列枚举LIVE/FINAL使同一(chain, consumer, addresses)组合可按 watcher 类型各存在一条并将唯一索引升级为含filter_type的版本同时用idx_filters_chain_consumer_live/idx_filters_chain_consumer_final两个部分索引分别加速热路径查询。这与文档中不同 watcher 可声明不同确认数finality / safe / N confirmations的设想直接呼应FINAL类型对应最终性交付流。Rust 侧的模型见 filter_model.rsFilter结构体完整映射上述列FilterType枚举通过sqlx映射到数据库枚举filter_type。5.2 blocks 表对应 logs 表的区块维度同一次迁移还创建了blocks表用于记录区块最小元数据CREATE TYPE block_status AS ENUM (CANONICAL, FINALIZED, UNCLE); CREATE TABLE IF NOT EXISTS blocks ( id UUID PRIMARY KEY, chain_id BIGINT NOT NULL, block_number BIGINT NOT NULL, block_hash BYTEA NOT NULL, parent_hash BYTEA NOT NULL, status block_status NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), CONSTRAINT uq_blocks_chain_block_hash UNIQUE (chain_id, block_hash) ); CREATE UNIQUE INDEX IF NOT EXISTS idx_blocks_unique_canonical_per_number ON blocks(chain_id, block_number) WHERE status CANONICAL;status枚举CANONICAL / FINALIZED / UNCLE精确对应文档中关于叔块UNCLE与确认状态的讨论同一高度只允许一条CANONICAL记录而UNCLE/FINALIZED可多条并存——这是 Cursor 算法完整性get_latest_canonical_block()必须返回唯一链尖的关键约束final_blocks表则专用于最终性流程finalized 区块永不重组因此每链每高度唯一见 迁移 SQL。5.3 动态 watcher 生命周期管理crates/listener_core/src/core/filters.rs 提供了按链隔离的Filters管理器封装了add_filter/remove_filter两个异步方法参数为consumer_id、from、to、log_address、filter_type并在每次操作时输出结构化日志。该类注释明确建议 Wrap inArcfor sharing between handlers and evm_listener说明它是控制面消息 → 数据库持久化之间的桥梁正是文档Declare new watchers dynamically的实现载体。六、消息契约事件负载与路由键Notifier 库消费的每一条消息都遵循仓库中 shared/primitives/src/event.rs 定义的契约BlockPayload区块级负载含flow处理流LIVE/REORGED/CATCHUP/FINAL/FINAL_CATCHUP、chain_id、block_number、block_hash、parent_hash、timestamp及有序的transactions数组TransactionPayloadfrom/to合约创建为None/hash/transaction_index/value/data/logsIndexedLoglog_index、address、topics、data即回执中单条日志的完整表示FilterCommandconsumer_idfrom/to/log_address均可选filter_type用于control.watch/control.unwatch控制面其validate()会 trimconsumer_id并拒绝空值四个地址字段全空时表示全区块wildcard订阅——消费方接收该区块的全部交易与日志对应文档消费所有事件的设计CatchupPayloadconsumer_idblock_start/block_end请求历史区块回放对应文档Replay past blocks的可选能力ReorgBacktrackEvent重组回溯消息携带block_number、新规范区块哈希与父哈希。路由键集中在 shared/primitives/src/routing.rsnew-event、final-event、catchup-event、final-catchup-event事件面以及control.watch/control.unwatch/catchup/range-catchup/backtrack-reorg控制面。消费端路由由consumer_new_event_routing(consumer_id)等函数生成consumer_id.event形式保证各消费者事件互不串扰routing.rs 中配有路由不冲突的单元测试。这些契约与路由正是subscribe / consume 清晰 API与不同类型 watcherlogs、tx的协议基础。七、可靠性机制重试、死信与熔断器文档在队列消费特性中点名要求retry queues。仓库中 shared/broker 的实现超出了字面需求提供了三层可靠性保障错误分类AsyncHandlerPayloadClassified处理器允许回调显式返回HandlerError::Transient基础设施故障无限重试并可触发熔断或HandlerError::Execution消息本身非法计入max_retries后进入死信队列而AsyncHandlerPayloadOnly将一切错误视为永久错误适用于简单场景重试预算与死信AMQP 模式下消息经retry交换机按amqp_retry_delay延迟返回主队列超过max_retries后进入dlx交换机对应的错误队列熔断器连续瞬态错误超过阈值如 3 次后消费暂停Open冷却期后进入 Half-Open 探针状态探针成功则恢复Closed失败则重新打开——这能有效防止下游 RPC/数据库故障期间向死信队列灌入海量消息状态机详见 broker README。对消费所有事件以免 rmq 内存增长的要求broker 还提供prefetch缓冲消息数、redis_claim_min_idle/redis_claim_intervalRedis Streams 模式下认领卡死消息等消费参数配合队列深度depth监控与 metrics 输出形成完整的队列健康保障。八、如何接入从示例到组件文档要求提供最小可运行示例仓库的 example crate 已给出可直接对照的接入模板live_events.rs订阅实时LIVE事件流final_events.rs订阅最终性FINAL事件流full_block.rs全区块订阅接收区块内所有交易与日志对应 wildcard filtertransfer.rs按地址/转账场景定向过滤的示例main.rs各示例的入口装配。同时consumer crate 提供client、options、error等模块可作为消费端客户端 API 的参考入口。接入方的大致路径是通过控制面队列control.watch以FilterCommand注册/注销自己的 watcher → 在事件队列new-event/final-event上以自身consumer_id订阅BlockPayload→ 在回调handler中解析IndexedLog并触发组件内部逻辑。结合 Listener Core 的 config.yaml 等配置含链与 RPC 配置项即可在真实环境中端到端运行。九、总结Notifier Library 的设计本质上是把链上事件 → 组件内部逻辑这条链路中的匹配与交付环节独立成库上游由 Listener Core 负责抓块、防重与广播下游由该库负责回执解析、过滤匹配与可靠转发。原文档提出的 watcher/logs 双表存储、动态 ABI 注册、多链多队列消费、确认块数分级、语义哈希去重、重试队列与可观测性等要点在仓库的filters/blocks/final_blocks迁移、FilterCommand/BlockPayload消息契约、共享 broker 的重试-死信-熔断机制以及 example/consumer 示例 crate 中均已找到对应实现或演进形态。对希望为 zama 组件接入链上事件能力的开发者而言这份设计与仓库源码互为印证是一套完整、可落地的事件订阅基础设施参考。【免费下载链接】fhevmFHEVM, a full-stack framework for integrating Fully Homomorphic Encryption (FHE) with blockchain applications项目地址: https://gitcode.com/GitHub_Trending/fh/fhevm创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考