区块链 区块链技术 比特币公众号手机端

使用 QuickNode Streams 在 Solana 上构建 AML 监控系统

liumuhui 6小时前 阅读数 1 #区块链

概述

反洗钱(AML)监控项目最常见于持有客户资金或连接银行系统的加密货币业务中。中心化交易所、法币出入金通道、支付处理商、托管钱包和稳定币发行方通常作为金融机构受到监管,无论是美国在 FinCEN 注册的货币服务企业、FATF 框架下的虚拟资产服务提供商,还是根据欧盟 MiCA 获得许可的加密资产服务提供商。这些类别的项目通常涉及交易监控、记录保存和可疑活动报告。

DeFi 原生的案例不太明显,但越来越常见。代币化国债和其他真实世界资产平台引入投资者,然后监控投资者在链上的行为。将存款置于 KYC 门槛之后的许可借贷池也运行类似的监控。为原本无需许可的协议运营前端或 relayer 的团队通常会进行筛查和监控,即使协议本身不在监管范围内。

监控、记录保存和报告都建立在可被证明完整的交易记录之上。这三者都假设日志包含审查窗口内发生的所有事情,这使得 AML 项目与它背后的记录一样可靠。当审计人员询问监控系统在特定 slot 看到了什么时,描述数据收集方式并不能回答这个问题。记录本身必须证明没有遗漏任何内容。

Solana 使这种证明比乍看起来更难。轮询循环可能会错过一个 slot,或者记录一个集群后来放弃的 slot,两种情况都不会引发错误。跳过的 slot 意味着缺少 slot 编号并不能证明出了问题,因此通常的完整性检查不适用。结果是产生一类容易制造、难以发现、事后解释成本很高的缺口。

本指南详解了一个位于 AML 监控项目之下的捕获层示例:一个自纠正的、记录触及给定 Solana 程序的交易的日志,在交易到达时对照制裁观察名单进行筛查,使用 Streams 和 Key-Value Store 构建。这是数据管道的一个起点,而不是一个完整的监控系统。

非法律或合规建议

本指南涵盖数据工程而非监管分析,其中的任何内容都不是法律建议,也不是说明以这种方式构建的系统是合规的。哪些要求适用于你的业务,以及任何监控系统是否满足这些要求,需要你自己的合规团队和法律顾问来确定。请将以下内容视为一个示例,根据你自己的义务进行调整和验证。

不完整交易记录的真实代价

交易记录是交付物。审计人员不会抽象地评估检测逻辑,他们会抽样日志并询问系统在特定时刻看到了什么,而每个答案都依赖于记录是完整且可证明完整的。

缺失一个区块在这里比在普通工程中代价更高。通常这是一个 bug:你注意到它、重放它、继续前进,付出的代价是几个工程工时。而在监控项目中,同样的缺失区块会留下一个你无法证明自己在观测的活动窗口,审计查看的是记录本身,而不是产生记录的努力。因此,必须主动发现并填补缺口,而当一个缺口存在时,风险在于你的组织,而不是任何丢失数据的提供商。

轮询并不适合这项工作。一个循环调用 getBlock、写入行并推进游标的工作进程,每当进程在 slot 中途重启或调用超时而游标仍然推进时,就会丢失数据,并且当它记录了一个集群后来放弃的 slot 时,会写入永远不会成为规范的数据。这些情况都不会引发错误,所以日志无论如何看起来都是健康的。

在 Solana 上,slot 可能被跳过,这意味着缺失的 slot 编号在设计中就存在歧义。它可能完全正常,也可能是你的缺口,仅凭 slot 编号无法区分。一个简单的不间断递增计数器完整性检查在这里不起作用。另外,在区块最终确定之前,集群可能在一个竞争分叉上达成共识,因此你观察并记录的某笔交易可能最终不在其他人同意的链上。包含非规范交易的记录与缺少规范交易的记录一样是问题,而且更难解释。

Solana 交易一旦被确认就不可逆转,因此当一笔被标记的交易到达你的日志时,它已经是最终的了。不存在银行在电汇离开前拒绝的对应物。检测使得下游操作成为可能,如冻结账户、拒绝未来活动或提交报告。只有阻止来自已被标记地址的未来活动这一狭窄情况才是真正预防性的,而且永远不是触发标记的那笔交易。

捕获层必须做对的事情

这些是本示例所围绕的属性。它们是数据需求而非合规要求,并且驱动下一节中的产品选择。

有序交付并标记纠正

数据必须按链顺序到达,当交付是纠正而非新活动时,管道必须在负载中说明这一点。无法区分新数据和纠正数据的消费者要么丢弃纠正,要么重复计数,两者都会产生一种事后难以发现的错误记录。

链变化时自动纠正

需要两种机制,它们不能互相替代。第一种是将交付保持在链尖端后方固定距离,这减少了记录后来被替换的 slot 的频率。它是概率性的,所以减少暴露但无法消除。

第二种是在确实发生分叉时,用规范数据重新交付受影响的区间,并加上标记以便消费者知道要覆盖写入。仅靠距离尖端的位置会让你在漏网的情况下永久错误。

填补后来发现的缺口的方法

自动纠正确实处理了链在尖端附近对你做的事情。但它对后来发现的缺口无能为力:管道暂停了一整天、早于系统存在的监控窗口、或审计人员要求你提供的回溯期。这需要第二种机制,由操作员发起,通过相同的筛查逻辑重放一个明确的 slot 区间并写入同一个表,使恢复的记录在形态和筛查上与实时捕获的记录完全相同。

不依赖程序自身检查的筛查

假设你监控的程序已经在链上拒绝已知不良地址。外部检查仍然重要,原因有两个不同的方面:

  • 过时性: 链上黑名单的时效性取决于其最后一次部署,而合规团队自己的名单(来自 OFAC 或商业提供商)可以每天更新多次。每次名单修订都重新部署程序是不现实的。
  • 规避: 一个动机明确的参与者可能完全绕过了程序的检查,通过一个未被标记的钱包、一个中间跳转,或程序自身执行逻辑中的一个缺口。

两种情况都指向同一个方向:筛查实际交易数据,独立于程序逻辑所做的事情。如果你唯一的控制是攻击者正在绕过的那个,那你实际上没有控制。

每个筛查结果都指明使用的名单

没有产生该结果的名单版本,筛查结果就毫无意义。制裁名单会变化,所以"我们检查时这个地址是干净的"只有在你能显示当时名单内容的情况下才是证据。

修复记录而不会两次通知分析师

这些是看起来相似的不同问题。数据集必须在 slot 变化时被纠正。分析师不能被同一笔底层交易页两次。解决其中一个并不能解决另一个,混淆它们往往会产生一个要么向分析师刷屏、要么悄悄丢弃纠正的系统。

独立审计的基础设施

你的控件建立在供应商的控件之上,审计人员会询问两者。SOC 1、SOC 2 和 ISO 27001 认证是这个问题实际采取的形式。

Quicknode 产品如何满足这些需求

两个 Quicknode 产品覆盖了上述大部分属性,交付到你拥有的数据库中:

Streams

Streams 将区块链数据作为有序的批次序列交付,在交付前在服务端应用你的过滤器,并将结果路由到你选择的目的地:webhook、S3 兼容存储、PostgreSQL、Azure Blob Storage 或 Kafka。

它的六个属性对应上面的列表:

  • 确认的顺序交付。 Streams 按顺序处理区块,在当前批次被确认交付之前不会前进到下一个批次。这是与轮询循环的结构性差异,在轮询循环中游标推进是因为你自己的代码决定的。这里游标推进是因为目的地确认了接收。
  • 失败暂停 Stream。 如果交付失败,Streams 根据你的配置重试,在最大尝试次数后暂停 Stream 而不是继续前进。一个损坏的消费者会变成你可以看到并恢复的停止管道,而不是记录中一个沉默的洞。
  • 最新区块延迟。 交付可以被保持在链尖端后面的可配置 slot 数量处。这是纠正需求的概率性一半。
  • 重组时重新流式传输。 这是确定性的一半。当 Streams 检测到它交付的链不再与规范链匹配时,它用规范数据重新交付受影响的区间,并在批次元数据中将交付标记为纠正。
  • 历史区间重放。 Stream 可以被赋予明确的开始和结束 slot,而不是跟随尖端,这就是你关闭已知缺口或产生回溯期的方式。因为重放运行相同的过滤器到相同的目的地,恢复的记录在形态和筛查上与实时记录无法区分。
  • 服务端过滤器。 你的筛查逻辑在管道内部运行,用 JavaScript 或 Go,因此交易在到达你的任何系统之前就被评估。这也意味着筛查结果随记录一起传输,而不是事后应用,因此结果是可审计的而不是重构的。

付费计划要求

每个付费计划都包括针对尖端流式传输的 重组时重新流式传输 选项,这是纠正步骤所依赖的机制。没有它,被替换的 slot 会以未纠正的状态留在你的记录中。

你可以交付到哪个目的地也取决于你的计划。本指南使用的 webhook 目的地是所有计划都可用的:

目的地 可用范围
Webhook 所有计划,包括免费试用
PostgreSQL Build 及以上
S3 兼容存储 Build 及以上
Kafka Build 及以上
Azure Blob Storage Accelerate 及以上

单个 Stream 可以供给的目的地数量也以同样的方式限制:免费试用和 Build 上为 1 个,Accelerate、Scale 和 Business 上最多 4 个,Enterprise 上最多 6 个。有关当前矩阵,请参阅定价页面。

如果你还没有账户,创建一个 Quicknode 账户并选择付费计划,然后在 Stream 设置中启用该选项。

Key-Value Store

Key-Value Store 是用于列表和键值对的托管存储,可通过两种方式访问:从 Streams 过滤器内部,以及从独立于 Streams 工作的 REST API。这种双重访问就是它比自建数据库更适合这里的原因。

它扮演两个角色,两者都存在是因为过滤器在没有网络访问的情况下运行,这使得 Key-Value Store 成为它可以读取的唯一可变状态:

  • 制裁观察名单,作为一个列表。因为 REST API 可以在不接触 Stream 的情况下更新列表,你的合规团队的名单可以与信息源发布一样频繁地变化,而过滤器代码保持不变。这直接回应了使链上黑名单不够用的过时性问题。列表作用于你的 Quicknode 账户而不是单个 Stream,因此同一个观察名单服务于每个 Stream 和任何需要检查成员资格的外部服务。
  • 列表版本,作为一个键值对。过滤器在每个批次上读取它并将其附加到每条记录上,因此筛查结果引用了产生它的修订版本。只有当你也归档每个修订版本包含的内容时,这个引用才有价值,实现部分会涵盖这一点。

警报去重是第三个状态片段,但它属于数据库而不是这里。抑制重复页面必须由知道写入成功的组件决定,而警报上的唯一键在与记录相同的事务中做到这一点。把它放在过滤器中会引入一个失败模式:被重试的交付发现键已经设置,并完全丢弃警报。

你控制的数据库

审计日志落在一个你拥有的数据库中。警报走相同的路径,所以永远只有一条路由进入你的系统,所有经过它的内容都已经通过了重组处理和观察名单检查。

拥有目的地也是使日志可被查询以进行调查、并可按照你自己的保留策略要求保留的原因,对于合规记录来说,保留期通常以年计。

参考架构

两个输出,一次数据遍历。日志是被动的、全面的。警报是狭窄的、中断性的。它们保持分离,因为它们有不同的消费者、不同的保留要求和不同的失败模式。

image.png

  • 捕获: Stream 读取 Solana 区块数据集,该数据集以 getBlock 返回的形状交付完整区块。每个区块都通过过滤器,过滤器只保留调用被监控程序的交易,并在任何内容交付之前丢弃其余部分。Streams 将每个新区块的父哈希与它之前交付的区块的哈希进行比较,在不匹配时向后轮询链以找到两个历史最后一次一致的地方。

  • 筛查: 对于每笔保留的交易,过滤器将移动价值的账户规范化为按地址的记录,在一次批量查找中检查它们是否在观察名单中,并用列表版本标记每条记录。一个负载离开过滤器时携带四个数组:批次中的区块、匹配的交易、按地址的筛查结果和任何警报。

  • 纠正: 每个批次都会告诉你它覆盖了哪些 slot,在 batch_start_rangebatch_end_range 中。摄取服务不是根据该区间是否向前移动来分支,而是删除它已经为该区间持有的任何内容,以及交付命名为重组的任何 slot,并重新写入本次交付。Streams 检测到分叉并重新交付,但只有你的代码才能让表反映这一点。

  • 去重: 因为纠正会重新交付已经见过的 slot,警报抑制以交易和地址对为键,而不是以交付为键。纠正记录和抑制重复页面有意地发生在不同层。

这里没有需要运行的索引器,也没有需要维护的跟随链的状态机,因为排序和纠正保证来自管道而不是来自你拥有的代码。

实现

接下来将端到端地构建整个东西:Key-Value Store 中的观察名单、配置为有序和重组纠正交付的 Stream、只保留触及你程序的交易并筛查每个移动价值的地址的过滤器、在 slot 被替换时保持表诚实的摄取服务,以及一个证明记录没有漏洞的查询。

该示例监控单个程序。因为该地址是一个常量,它保持硬编码在过滤器中,而 Key-Value Store 仅用于过滤器无法硬编码的两件事:观察名单(按合规团队的时间表变化)和标识哪个修订版本产生了给定筛查结果的版本字符串。

你需要什么

  • 付费计划上的 Quicknode 账户(重组时重新流式传输和 Solana 历史重放都需要)
  • 你的 Quicknode Key-Value Store REST API 密钥
  • Node.js(v20.6+,支持 --env-file
  • 本地安装 PostgreSQL(v14+),如下所示,或使用 Supabase 等主机上的数据库
  • ngrok 用于在开发时暴露你的本地摄取服务
  • 你想监控的 Solana 程序的地址

设置项目

观察名单同步任务和摄取服务共享一个小型 Node 项目。先创建它,以便后续每一步读取的配置有一个家:

mkdir solana-aml && cd solana-aml
npm init -y
npm install express@4.21.2 pg@8.13.1

两个服务都从环境变量读取配置。三个值涵盖它们需要的所有内容,所以把它们放在项目根目录的 .env 文件中:

.env

DATABASE_URL="postgres://localhost:5432/solana_aml"
QN_API_KEY="your-quicknode-api-key"
QN_SECURITY_TOKEN=""

然后在放入任何真实内容之前把它排除在版本控制之外:

echo ".env" >> .gitignore

DATABASE_URL 指定了将在下一节中创建的数据库。如果你按照那里的本地安装操作,就保持原样。

QN_API_KEY 是你的 Quicknode API 密钥,它属于你的账户而不是任何单个 Stream。所有出站调用 Quicknode 的内容都将它作为 x-api-key 标头发送:发布列表到 Key-Value Store 的观察名单同步任务,以及稍后创建 Streams 的 curl 命令。

QN_SECURITY_TOKEN 方向相反,目前保持为空。 Quicknode 为每个 Stream 生成它,你永远不会把它发送到任何地方,摄取服务使用它来验证传入的 webhook 确实来自你的 Stream。签发它的 Stream 还不存在,所以你将按照"编写 Stream 过滤器"创建 Stream,并在运行摄取服务之前从其设置选项卡中填入这个值。

Node 直接使用 --env-file 读取该文件,所以不需要 dotenv 包:

node --env-file=.env sync-watchlist.js

设置数据库

本指南中的所有内容都读写一个 PostgreSQL 数据库,位于你刚刚放入 DATABASE_URL 的连接字符串处。如果你已经有数据库,把连接字符串放在那里,跳到下一节。

否则,安装 Postgres 并启动服务器。在安装了 Homebrew 的 macOS 上:

brew install postgresql@16
brew services start postgresql@16

在 Debian 或 Ubuntu 上:

sudo apt install postgresql
sudo systemctl start postgresql

然后创建 DATABASE_URL 所指定的数据库:

createdb solana_aml

在创建任何内容之前确认你可以连接:

psql solana_aml -c "SELECT version();"

创建审计表

五张表,每张一个职责。保存为 schema.sql

schema.sql

-- 每个交付的区块一行。这是完整性主干。
CREATE TABLE blocks (
  slot               BIGINT PRIMARY KEY,
  blockhash          TEXT NOT NULL,
  previous_blockhash TEXT NOT NULL,
  block_time         TIMESTAMPTZ,
  received_at        TIMESTAMPTZ NOT NULL DEFAULT now(),
  corrected_at       TIMESTAMPTZ
);

-- 每个匹配的交易一行,关联到它所在的区块。
CREATE TABLE transactions (
  signature    TEXT    NOT NULL,
  slot         BIGINT  NOT NULL REFERENCES blocks(slot) ON DELETE CASCADE,
  signer       TEXT    NOT NULL,
  fee_lamports BIGINT,
  succeeded    BOOLEAN NOT NULL,
  PRIMARY KEY (signature, slot)
);

-- 每个移动价值的地址一行,包含筛查结果和产生它的列表修订版本。
CREATE TABLE screened_accounts (
  signature    TEXT    NOT NULL,
  slot         BIGINT  NOT NULL REFERENCES blocks(slot) ON DELETE CASCADE,
  address      TEXT    NOT NULL,
  sol_delta    BIGINT,
  sanctioned   BOOLEAN NOT NULL,
  list_version TEXT    NOT NULL,
  PRIMARY KEY (signature, slot, address)
);

-- 调查直接按地址查找。上面的主键以 signature 开头,所以无法单独服务于该查询。
CREATE INDEX screened_accounts_address_idx ON screened_accounts (address);

-- 已发出的警报。有意不关联区块:见摄取服务部分。
CREATE TABLE alerts (
  alert_key     TEXT PRIMARY KEY,
  signature     TEXT   NOT NULL,
  slot          BIGINT NOT NULL,
  address       TEXT   NOT NULL,
  list_version  TEXT   NOT NULL,
  raised_at     TIMESTAMPTZ NOT NULL DEFAULT now(),
  dispatched_at TIMESTAMPTZ
);

-- 已归档的观察名单修订,使标记的 list_version 可以解析为产生它的确切内容。
CREATE TABLE watchlist_versions (
  list_version TEXT PRIMARY KEY,
  content_hash TEXT   NOT NULL,
  addresses    TEXT[] NOT NULL,
  published_at TIMESTAMPTZ NOT NULL DEFAULT now()
);

然后应用它:

psql solana_aml -f schema.sql

其中有五个细节在做实际工作:

  • blocks 为每个交付的区块获得一行,即使是没有匹配交易的区块。 这就是使完整性可检查的原因。仅匹配交易的表无法告诉你一段安静的时间意味着"什么都没发生"还是"我们没在看"。
  • 存储了 previous_blockhash,而不仅仅是 slot。 Solana 会跳过 slot,所以缺失的 slot 编号证明不了什么。区块哈希无论中间跳过了多少个 slot 都彼此链接,这就是完整性检查所利用的。
  • transactionsscreened_accounts 上的 ON DELETE CASCADE 是整个纠正机制。删除一个 blocks 行会带上它的交易及其筛查结果,所以替换被替换的 slot 是一条单一语句。
  • alerts 有意没有外键和级联。 已经发送给分析师的页面是一个历史事实,重组不会撤销它。将警报保持在级联之外是将纠正数据集与抑制重复页面分开的原因,否则上述两种机制会混淆。
  • watchlist_versions 存储实际地址,而不仅仅是版本标签。 记录上的版本字符串是一个指针,指向虚无的指针证明不了什么。归档每个修订版本让你一年后能回答"那个列表包含什么",也是使验证查询中的结论可以独立重新推导的原因。

筛查结果存在于 screened_accounts 而不是 transactions 上,因为筛查是按地址而不是按交易的。一笔转账为发送方和接收方各产生一行,每行携带自己的结论和 list_version

这些表没有解决的是随时间增长的数据量。 它们被写成正确且易于推理,而不是在规模上最优,其中两张因不同原因稳定增长。blocks 以链产生区块的速率增长,无论你的程序有多安静,每天大约 200,000 行,这是这里最便宜的表。screened_accounts 随你的程序活动乘以每笔交易地址数增长,它是第一个变大变多的。

这在多大程度上重要完全取决于你被监控的程序看到多少流量,诚实的范围从微不足道到真正的工程问题。因为 AML 记录保存通常以年计,它会复合而不是趋于平稳。

生产部署可能希望某些组合:按 slot 或按月进行范围分区、对较旧分区进行压缩和使用更便宜的存储、与你实际监管窗口相关的保留策略,以及可能比每地址一行更紧凑的筛查结果表示。有两个约束值得带入该设计:缩小筛查对象是一个覆盖决策而不是存储优化,因为它使行的缺失变得模糊,而 watchlist_versions 是唯一不能随意修剪的表,因为丢弃修订版本会使引用它的每个结论失效。

规模调整、分区和归档都不在本指南范围内。在它处理生产流量之前,根据你自己的数据量和保留义务来制定这些。

将观察名单加载到 Key-Value Store

过滤器在没有网络访问的沙箱中运行,所以它不能调用你的数据库、内部服务或制裁信息源。Key-Value Store 是过滤器可以读取的唯一可变状态,这就是使它成为独立于你的代码变化的列表的合适家园的原因。

那里有两个键:列表本身,以及标识它当前持有的列表修订版本的字符串。过滤器读取两者。

那个版本字符串是值得小心的部分。单独来看它只是一个标签,当有人问列表在给定 slot 持有什么时,指向你不再拥有的内容的标签证明不了什么。所以下面的同步任务同时做三件事:它从内容派生版本而不是断言它,在该版本下归档内容,然后才发布。

将你的信息源地址以扁平 JSON 地址字符串数组的形式放在 watchlist.json 中,例如 ["Addr1...", "Addr2..."],然后从与摄取服务相同的项目中运行:

// sync-watchlist.js
const crypto = require('crypto');
const fs = require('fs');
const { Pool } = require('pg');

const pool = new Pool({ connectionString: process.env.DATABASE_URL });
const KV = 'https://api.quicknode.com/kv/rest/v1';
const WATCHLIST_KEY = 'sanctioned-addresses';
const LIST_VERSION_KEY = 'sanctions-list-version';
const FEED_ID = 'ofac-sdn';

async function kv(method, path, body) {
  const res = await fetch(`${KV}${path}`, {
    method,
    headers: {
      'Content-Type': 'application/json',
      'x-api-key': process.env.QN_API_KEY,
    },
    body: body ? JSON.stringify(body) : undefined,
  });
  const text = await res.text();
  if (!res.ok) throw new Error(`${method} ${path} -> ${res.status} ${text}`);
  return text ? JSON.parse(text) : null;
}

// 缺失的键是预期状态,而不是失败,所以 404 作为 null 返回,而 403 仍然抛出。
// 把两者混为一谈会把错误的 API 密钥变成看起来已发布但实际上没有的观察名单。
async function kvGet(path) {
  const res = await fetch(`${KV}${path}`, {
    headers: { 'x-api-key': process.env.QN_API_KEY },
  });
  if (res.status === 404) return null;
  const text = await res.text();
  if (!res.ok) throw new Error(`GET ${path} -> ${res.status} ${text}`);
  return text ? JSON.parse(text) : null;
}

async function main() {
  // 替换为你的信息源解析器。它只需要返回一个地址数组。
  const raw = JSON.parse(fs.readFileSync('watchlist.json', 'utf8'));
  const addresses = [...new Set(raw)].sort();

  // 版本是从内容中派生的,所以它不能与内容不一致。
  // 手写的日期可以:有人在不改变列表的情况下提升它,或者
  // 改变列表但忘记提升它。内容哈希不会撒谎。
  const contentHash = crypto
    .createHash('sha256')
    .update(addresses.join('\n'))
    .digest('hex');
  const today = new Date().toISOString().slice(0, 10);
  const listVersion = `${FEED_ID}-${today}.${contentHash.slice(0, 8)}`;

  // Key-Value Store 现在实际持有的内容。针对这个而不是归档进行协调,
  // 是让一个半完成的运行自我修复的方式。
  const liveItems = (await kvGet(`/lists/${WATCHLIST_KEY}`))?.data?.items ?? null;
  const liveVersion = (await kvGet(`/values/${LIST_VERSION_KEY}`))?.data?.value ?? null;

  const client = await pool.connect();
  try {
    const { rows: [previous] } = await client.query(
      `SELECT list_version, content_hash, addresses
       FROM watchlist_versions ORDER BY published_at DESC LIMIT 1`,
    );

    // 仅在归档、实时列表和实时版本都一致时才跳过。
    // 仅检查归档会把中途失败的运行视为已完成并永不重试。
    const listMatches = liveItems
      && liveItems.length === addresses.length
      && [...liveItems].sort().join('\n') === addresses.join('\n');

    if (previous?.content_hash === contentHash
        && liveVersion === previous.list_version
        && listMatches) {
      console.log(`no change, still ${previous.list_version}`);
      return;
    }

    // 1. 先归档,使版本在任何记录引用它之前解析为真实内容。
    await client.query(
      `INSERT INTO watchlist_versions (list_version, content_hash, addresses)
       VALUES ($1, $2, $3)
       ON CONFLICT (list_version) DO NOTHING`,
      [listVersion, contentHash, addresses],
    );

    // 2. 将实时列表对齐,作为对列表实际内容的差异。
    //    对从未创建的列表 PATCH 会返回 404。
    let addItems = addresses;
    let removeItems = [];

    if (liveItems) {
      const before = new Set(liveItems);
      const after = new Set(addresses);
      addItems = addresses.filter((a) => !before.has(a));
      removeItems = liveItems.filter((a) => !after.has(a));
      if (addItems.length || removeItems.length) {
        await kv('PATCH', `/lists/${WATCHLIST_KEY}`, { addItems, removeItems });
      }
    } else {
      await kv('POST', '/lists', { key: WATCHLIST_KEY, items: addresses });
    }

    // 3. 最后发布版本。见下面关于顺序的说明。
    await kv('POST', '/values', { key: LIST_VERSION_KEY, value: listVersion });

    console.log(
      `published ${listVersion}:`,
      `+${addItems.length} -${removeItems.length}, ${addresses.length} total`,
    );
  } finally {
    client.release();
    await pool.end();
  }
}

main().catch((err) => {
  console.error(err);
  process.exit(1);
});

只要你发布信息源就运行它:

node --env-file=.env sync-watchlist.js

其中有四件事是刻意为之。

这个任务针对 Key-Value Store 而不是归档进行协调。 归档自行提交,所以在 KV 调用时失败的一次运行会留下一个版本行,但没有发布任何内容。仅从归档决定是否有工作要做,会把一个部分失败的运行读作已完成的同步,并在每次后续运行时跳过它,使列表永久为空而任务报告没有变化。先读取实时列表和版本意味着下一次运行会修复它。

版本是内容哈希,而不是日期。 ofac-sdn-2026-07-30.a3f9c1b2 是从排序后的地址派生的,所以它不能偏离它描述的内容,在没有变化时重新运行任务是一个空操作而不是一个虚假的新版本。日期前缀是为人类准备的;后缀是可验证的部分。

归档在列表发布之前写入。 反过来做就会有一个窗口,记录引用了一个解析为虚无的版本。

版本值最后发布。 Key-Value Store 没有跨越列表和值的事务,所以有一个短暂的窗口,两者不一致。最后发布版本意味着在窗口期间筛查是当前的,而标签滞后,这是更安全的失败:你宁可捕获一个新列出的地址并有一个簿记差异,也不愿基于过时的列表进行筛查。无论哪种方式,差异都是可检测的,"验证记录完整"下的连贯性查询就是找到它的。

注意两个接口之间的大小写差异,这很容易绊倒人:REST API 接受 addItemsremoveItems,而过滤器内部的 qnUpsertList 方法接受 add_itemsremove_items

因为列表作用于你的 Quicknode 账户而不是单个 Stream,同一个观察名单服务于你运行的每个 Stream 和任何需要通过 REST API 检查成员资格的外部服务。

从真实来源填充列表。 从实际信息源填充 watchlist.json,无论是对 OFAC SDN 列表 的解析器还是商业提供商,永远不要使用你没有从中获取的地址。为了在接入信息源之前冒烟测试警报路径,放入一个你期望在自己的匹配交易中看到的地址,之后移除它。将你添加的任何内容视为关于真实方的声明,因为这就是你的日志中的制裁命中稍后会被读取的方式。

编写 Stream 过滤器

在仪表板中创建你的 Stream,选择 Solana主网,选择 区块 数据集,并选择"在流式传输前修改负载"。将默认的 main 函数替换为以下内容:

// Solana 区块数据集,jsonParsed 编码。要求批大小为 1(见下面的说明)。
const MONITORED_PROGRAM = 'jupoNjAxXgZ4rjzxzPMP4oxduvQsQtZzyknqvzYNrNu';
const WATCHLIST_KEY = 'sanctioned-addresses';
const LIST_VERSION_KEY = 'sanctions-list-version';

async function main(stream) {
  const meta = stream.metadata || {};
  const startSlot = meta.batch_start_range;
  const endSlot = meta.batch_end_range;

  // 每批次一次读取。每条记录都携带这个,所以筛查结果引用它所依据的确切修订版本。
  // 修订版本本身归档在 watchlist_versions 中,这就是引用解析到的内容。
  // 如果键缺失,我们标记 'unversioned' 而不是抛出异常,所以配置错误会降级为
  // 连贯性查询标记的行,而不是暂停的 Stream。
  const listVersion = (await qnLib.qnGetValue(LIST_VERSION_KEY)) || 'unversioned';

  const blocks = [];
  const matched = [];

  for (const block of stream.data) {
    blocks.push({
      slot: startSlot,
      blockhash: block.blockhash,
      previousBlockhash: block.previousBlockhash,
      blockTime: block.blockTime,
    });

    for (const tx of block.transactions || []) {
      const keys = (tx.transaction?.message?.accountKeys || []).map((k) => k.pubkey || k);
      if (!keys.includes(MONITORED_PROGRAM)) continue;

      matched.push({
        signature: tx.transaction.signatures[0],
        signer: keys[0],
        feeLamports: tx.meta?.fee ?? null,
        succeeded: tx.meta?.err === null,
        accounts: accountsThatMovedValue(tx, keys),
      });
    }
  }

  // 批次中的每个地址在单次往返中被筛查。结果与我们发送的地址索引对齐。
  const addresses = [...new Set(matched.flatMap((m) => m.accounts.map((a) => a.address)))];
  const hits = addresses.length
    ? await qnLib.qnContainsListItems(WATCHLIST_KEY, addresses)
    : [];
  const sanctioned = new Set(addresses.filter((_, i) => hits[i]));

  const transactions = [];
  const screened = [];
  const alerts = [];

  for (const m of matched) {
    transactions.push({
      signature: m.signature,
      slot: startSlot,
      signer: m.signer,
      feeLamports: m.feeLamports,
      succeeded: m.succeeded,
    });

    for (const account of m.accounts) {
      const isSanctioned = sanctioned.has(account.address);

      screened.push({
        signature: m.signature,
        slot: startSlot,
        address: account.address,
        solDelta: account.solDelta,
        sanctioned: isSanctioned,
        listVersion,
      });

      if (isSanctioned) {
        alerts.push({
          alertKey: `${m.signature}:${account.address}`,
          signature: m.signature,
          slot: startSlot,
          address: account.address,
          listVersion,
        });
      }
    }
  }

  // 始终返回负载,永远不要返回 null。在没有匹配交易的区块上,
  // 数组是空的,但 `blocks` 仍然携带完整性检查所依赖的哈希链接。
  return {
    batch: {
      startSlot,
      endSlot,
      reorgedSlots: meta.blocks_reorged || [],
    },
    blocks,
    transactions,
    screened,
    alerts,
  };
}

// 余额实际变化的地址。SOL 变动来自余额数组。
// Token 变动是针对 token 账户报告的,所以需要筛查的所有者是地址,而不是 token 账户本身。
function accountsThatMovedValue(tx, keys) {
  const pre = tx.meta?.preBalances || [];
  const post = tx.meta?.postBalances || [];
  const moved = new Map();

  keys.forEach((address, i) => {
    const delta = (post[i] ?? 0) - (pre[i] ?? 0);
    if (delta !== 0) moved.set(address, delta);
  });

  const tokenAccounts = new Map();
  for (const balance of tx.meta?.preTokenBalances || []) {
    tokenAccounts.set(balance.accountIndex, {
      owner: balance.owner,
      before: balance.uiTokenAmount?.amount ?? '0',
      after: '0',
    });
  }
  for (const balance of tx.meta?.postTokenBalances || []) {
    const entry = tokenAccounts.get(balance.accountIndex)
      || { owner: balance.owner, before: '0', after: '0' };
    entry.owner = entry.owner || balance.owner;
    entry.after = balance.uiTokenAmount?.amount ?? '0';
    tokenAccounts.set(balance.accountIndex, entry);
  }
  for (const entry of tokenAccounts.values()) {
    if (entry.owner && entry.before !== entry.after && !moved.has(entry.owner)) {
      moved.set(entry.owner, null);
    }
  }

  return [...moved].map(([address, solDelta]) => ({ address, solDelta }));
}

将相同的代码保存为项目目录中的 filter.js。仪表板持有 Stream 运行的副本,但"配置 Stream"和"关闭已知缺口"下的 REST API 调用都从磁盘读取它。

将"测试区块"字段设置为最近的 slot,然后点击"▶️ 运行测试"。使用上面的 SPL Token 程序,你将匹配大多数区块,这使其成为一个有用的冒烟测试。看到它工作后,换入你自己的程序地址。

测试会填充 stream.metadata,所以 batch_start_range 解析为你测试的 slot,startSlotendSlot 返回等于它。测试无法演练的唯一部分是纠正路径:blocks_reorged 仅在实际上替换链历史的交付中出现,所以 reorgedSlots 在每个测试中都是空数组。

此过滤器假定 jsonParsed 编码

为了保持代码简短,此过滤器直接从 message.accountKeys 读取账户,并信任该数组持有交易中的每个账户,索引与 preBalancespostBalances 对齐。这在 jsonParsed 下成立,其中地址查找表账户被合并到 accountKeys 中,每个条目是携带 pubkey 的对象。

这个选择有值得在生产流量之前理解的后果。在 json 编码下,accountKeys 只持有静态声明的键,作为裸字符串,查找表账户单独出现在 meta.loadedAddresses 中,在静态键之后按可写然后只读排序。仅针对 accountKeys 编写的代码会错过每个通过查找表到达的程序,并且永远不会筛查从查找表加载的地址。两种失败都不会引发错误,在最近的 Mainnet 区块上,这个缺口覆盖了超过一半触及 SPL Token 程序的交易,在筛查系统中这是最糟糕的错误方式。

如果你流式传输除 jsonParsed 之外的任何编码,重新组装完整账户列表并保持它与余额数组对齐是你需要处理的。

那个过滤器中的七个决定值得说明,因为每个都是这个代码的明显版本出错的地方。

它永远不会返回 null,即使没有交易匹配。 大多数过滤器在没有匹配时返回 null,这完全跳过交付。这里这会破坏完整性检查,所以值得精确说明过滤器在安静区块上返回什么:一个空的 transactions 数组,但一个填充的 blocks 数组。区块记录是哈希链接,这是你不能承受丢失的东西。

要明白为什么,取三个连续区块,其中只有第一个和第三个触及你的程序。跳过中间那个,你的表持有 slot 100(哈希 A)和 slot 102(父哈希 B)。这些行在链上不再是相邻的,所以连贯性查询将 BA 比较,发现它们不同,并在 slot 102 报告一个缺口。实际上没有缺失任何东西。因为大多数区块不会触及你的程序,这个误报在几乎每个安静时段都会触发,而一个不断触发的完整性警报是你不再读的警报。

存储每个区块不会污染日志,因为 blockstransactions 是分开的表:完整性主干保持密集,而交易日志保持稀疏。成本是交付量而不是 API 积分,因为积分按处理的区块消耗,无论过滤器返回什么。在 Solana 上,这大约是每秒 2.5 次小交付。如果这些流量对你的用例不值得,你接受的权衡是放弃独立的完整性验证并依赖 Streams 的顺序交付。

slot 来自 metadata.batch_start_range,而不是来自区块体。 Solana 的 getBlock 响应不包含它自己的 slot 编号。诱人的替代品是 parentSlot + 1,它只在没有 slot 在此区块之前被跳过时是正确的,当有 slot 被跳过时会静默地差一个或多个。批次区间是权威的。这也是为什么批大小必须为 1:在该大小下,batch_start_rangebatch_end_range 是同一个 slot,每次交付对应恰好一个区块。

accountKeys 上匹配能捕获 CPI。 检查 instructions[].programId 只找到顶层调用,所以通过中间人到达你的程序的交易会溜过去。交易中任何地方调用的每个程序,包括通过跨程序调用,都必须出现在该交易的账户键中,所以这个检查覆盖两种情况。

批次区间和重组标志被复制到负载体中。 Streams 也以 HTTP 标头发送元数据,但一旦过滤器塑造负载,体就是你过滤器返回的任何内容。将区间和 blocks_reorged 放在负载内是给你带内完整性信号的方式:交付说明它覆盖什么以及是否是纠正,而不需要你的摄取服务依赖标头解析。

筛查在每次交付中进行一次批量查找,而不是每地址一次。 过滤器收集批次中的每个地址,去重,并进行单个 qnContainsListItems 调用,返回一个与发送内容对齐的布尔数组。文档对此强调,原因是延迟:嵌套循环中的逐地址调用会把一次往返变成数百次,而 Stream 在过滤器返回之前无法前进。

筛查在余额变化上运行,而不是在指令解析上。 过滤器筛查任何 SOL 余额移动的地址,加上任何余额移动的 token 账户的所有者。从 meta 中的余额数组工作而不是解码指令意味着它不需要理解你的程序的指令布局,并且不能被异常调用路径规避:如果价值移动了,余额就改变了,地址就被筛查了。Token 余额是针对 token 账户而不是钱包报告的,这就是代码解析 owner 而不是筛查 token 账户地址的原因,因为 token 账户不是制裁名单点名的方。

筛查结论在交付前计算并随记录一起传输。 这是可审计结果与重构结果之间的区别。下游没有任何东西必须重新推导地址是否被列出,也没有任何东西必须猜测哪个修订版本适用,因为两者都已经在行上了。

构建摄取服务

在你之前设置的项目中创建 server.js

// server.js
const crypto = require('crypto');
const express = require('express');
const { Pool } = require('pg');

const pool = new Pool({ connectionString: process.env.DATABASE_URL });
const SECURITY_TOKEN = process.env.QN_SECURITY_TOKEN;

const app = express();
app.use(express.raw({ type: 'application/json', limit: '50mb' }));

function verifySignature(req) {
  const nonce = req.get('X-QN-Nonce');
  const timestamp = req.get('X-QN-Timestamp');
  const signature = req.get('X-QN-Signature');
  if (!nonce || !timestamp || !signature) return false;

  const expected = crypto
    .createHmac('sha256', SECURITY_TOKEN)
    .update(nonce + timestamp + req.body.toString('utf8'))
    .digest('hex');

  const a = Buffer.from(expected, 'hex');
  const b = Buffer.from(signature, 'hex');
  return a.length === b.length && crypto.timingSafeEqual(a, b);
}

app.post('/streams', async (req, res) => {
  if (!verifySignature(req)) {
    return res.status(401).send('invalid signature');
  }

  const { batch, blocks, transactions, screened, alerts } = JSON.parse(
    req.body.toString('utf8'),
  );
  const client = await pool.connect();

  try {
    await client.query('BEGIN');

    // reorgedSlots 是这个交付替换了链历史的权威信号。
    // 单独的非空删除只意味着我们之前见过这些 slot,普通重试也是如此。
    const reorgedSlots = batch.reorgedSlots || [];
    const isCorrection = reorgedSlots.length > 0;

    // 这个 slot 区间中任何已存储的内容都被刚到达的内容取代。
    // 先移除它就是使纠正成为纠正的原因:只存在于被替换区块上的行消失,而不是
    // 坐在规范行旁边的表中。重组的 slot 与交付区间一起被删除,因为新的规范链
    // 可能对之前有区块的 slot 根本没有区块,而从未被重新交付的 slot 否则会保留其被替换的行。
    const { rowCount: rewritten } = await client.query(
      'DELETE FROM blocks WHERE slot BETWEEN $1 AND $2 OR slot = ANY($3::bigint[])',
      [batch.startSlot, batch.endSlot, reorgedSlots],
    );

    for (const b of blocks) {
      await client.query(
        `INSERT INTO blocks (slot, blockhash, previous_blockhash, block_time, corrected_at)
         VALUES ($1, $2, $3, to_timestamp($4), CASE WHEN $5 THEN now() END)
         ON CONFLICT (slot) DO NOTHING`,
        [b.slot, b.blockhash, b.previousBlockhash, b.blockTime, isCorrection],
      );
    }

    for (const t of transactions) {
      await client.query(
        `INSERT INTO transactions (signature, slot, signer, fee_lamports, succeeded)
         VALUES ($1, $2, $3, $4, $5)
         ON CONFLICT (signature, slot) DO NOTHING`,
        [t.signature, t.slot, t.signer, t.feeLamports, t.succeeded],
      );
    }

    for (const s of screened) {
      await client.query(
        `INSERT INTO screened_accounts
           (signature, slot, address, sol_delta, sanctioned, list_version)
         VALUES ($1, $2, $3, $4, $5, $6)
         ON CONFLICT (signature, slot, address) DO NOTHING`,
        [s.signature, s.slot, s.address, s.solDelta, s.sanctioned, s.listVersion],
      );
    }

    // 警报插入就是去重门。alert_key 是主键,所以已经警报过的签名和地址对会冲突
    // 并返回空。只有不存在的行会返回,也只有那些行会变成页面。
    const toDispatch = [];
    for (const a of alerts) {
      const { rows } = await client.query(
        `INSERT INTO alerts (alert_key, signature, slot, address, list_version)
         VALUES ($1, $2, $3, $4, $5)
         ON CONFLICT (alert_key) DO NOTHING
         RETURNING alert_key`,
        [a.alertKey, a.signature, a.slot, a.address, a.listVersion],
      );
      if (rows.length > 0) toDispatch.push(a);
    }

    await client.query('COMMIT');

    if (isCorrection) {
      console.log(
        `corrected slots ${batch.startSlot}-${batch.endSlot}`,
        `reorged: ${JSON.stringify(reorgedSlots)}`,
      );
    } else if (rewritten > 0) {
      console.log(`re-delivery of slots ${batch.startSlot}-${batch.endSlot}, rewrote in place`);
    }

    // 仅在提交后分发。就你未能存储的记录向分析师发出页面,比晚一点发出更糟。
    for (const a of toDispatch) {
      try {
        await dispatchAlert(a);
        await pool.query(
          'UPDATE alerts SET dispatched_at = now() WHERE alert_key = $1',
          [a.alertKey],
        );
      } catch (err) {
        // 该行已经提交,所以警报不会丢失。扫描 dispatched_at IS NULL 的行并在带外重试。
        console.error('alert dispatch failed, row retained', a.alertKey, err);
      }
    }

    res.sendStatus(200);
  } catch (err) {
    // 如果失败的是连接本身,ROLLBACK 也会抛出。吞掉它,以便下面的 500 仍然发出,
    // 而不是请求挂起直到交付超时。
    await client.query('ROLLBACK').catch(() => {});
    console.error('ingest failed, refusing to acknowledge', err);
    // 任何非 2xx 都会使 Streams 重试然后暂停。永远不要确认你没有实际写入的交付。
    res.sendStatus(500);
  } finally {
    client.release();
  }
});

// 替换为你的分页路径:PagerDuty、Slack、案件管理队列。
// 让它保持在数据库事务之外。
async function dispatchAlert(alert) {
  console.log(
    ` SANCTIONS HIT ${alert.address} in ${alert.signature}`,
    `at slot ${alert.slot} (list ${alert.listVersion})`,
  );
}

app.listen(3000, () => console.log('ingest listening on http://localhost:3000'));

从 Stream 的设置选项卡中获取安全Token填入 .env 中的 QN_SECURITY_TOKEN,然后运行它:

node --env-file=.env server.js

然后暴露它:

ngrok http 3000

那个服务有趣的地方在于它有多小。一个 DELETE 后跟插入,在一个事务内,处理了本来需要各自代码路径的三种情况:

  • 重组纠正。 Streams 用规范数据重新交付受影响的 slot。删除丢弃被替换的版本,包括任何只存在于旧区块上的交易,插入写入新的。仅 upsert 无法做到这一点,因为 upsert 无法移除纠正不再包含的行。删除 reorgedSlots 中指定的 slot 以及交付区间,覆盖了新规范链对之前有区块的 slot 没有区块的情况,这种情况永远不会被重新交付,所以否则会永远保留其被替换的行。
  • 重试交付。 如果你的服务成功写入但确认从未到达 Streams,同一批次会再次到达。删除区间并重写它会产生相同的表,所以重复是无害的。这就是去重需求,通过使写入幂等而不是通过跟踪已见内容来满足。
  • 重叠重放。 覆盖你已经拥有的 slot 的历史 Stream 会就地重写它们而不是与它们冲突,这就是使填补缺口的重放可以在不先检查的情况下安全运行的原因。

注意删除和 corrected_at 标记由不同信号驱动,混淆它们是一个容易犯的错误。删除在每次交付时运行,因为这就是使写入幂等的原因。corrected_at 时间戳仅在 reorgedSlots 非空时设置,因为那是链历史实际改变的唯一情况。改为以删除为键设置时间戳会在每次普通重试时标记"已纠正",使你在稍后审查日志时无法区分网络故障和重组。

警报位于所有这一切之外,这就是将它们保留在自己的没有级联的表中的意义。blocks 删除自由地重写数据集,如果需要可以一遍又一遍,而 alerts 累积。因为 alert_key 是主键,插入本身就做了抑制:已经发出页面的签名和地址对会冲突,RETURNING 产生空,不会发出第二个页面。因此,纠正记录和抑制重复页面发生在不同层,不能互相干扰。

两个值得明确陈述的后果。移除交易的重组不会撤销已经发出的警报,而且不应该:警报是你所采取行动的记录,alerts 行保留其 slot,以便你稍后可以对照 blocksscreened_accounts 进行核对。这是运行一个有意义最新区块延迟的原因之一,因为避免对被替换的交易发出页面的最便宜方式不是那么早看到它。因为分发发生在提交之后而不是在它内部,分页中断会让你留下一个已提交的警报,其 dispatched_at 仍然是 NULL,这是一个你可以扫描的查询,而不是一个丢失的页面。

image.png

决定警报的去向

dispatchAlert 有意保留为存根,因为这是数据管道停止和合规流程开始的地方。管道的工作以产生持久、去重的命中而结束。该命中发生什么由你决定,这取决于你的警报量、必须多快处理,以及你的程序是否要求谁审查过的记录。

一些选项,大致从最中断到最不中断:

  • 待命分页,通过类似 PagerDuty 或 Opsgenie 的方式。仅当有人确实必须在几分钟内行动时才合适。记住交易已经是最终的,所以存在的任何紧迫性来自下游操作,如冻结账户或切断未来活动,而不是来自交易本身。
  • 监控聊天频道,在 Slack 或 Teams 中。低摩擦且常见,但它不记录谁看了什么,而这往往正是审查人员会问的。
  • 案件管理,无论是通用工单系统还是专门构建的 AML 案件工具。设置工作更多,通常最合适,因为它为每次命中产生一个分配、一个调查员和一个文件化的处理结果。这个产物通常是监管机构想看到的,而不仅仅是之前的通知。
  • 队列或事件总线,如 SQS、Kafka 或 Pub/Sub。当多个系统需要同一个命中时,或者当你希望以后能够将警报重放到新消费者时,值得使用。
  • 完全不推送。 在低量时这是可辩护的。alerts 表已经是一个工作队列,定期审查 dispatched_at IS NULL 的行可能本身就满足程序。如果这是选择,写下来作为程序而不是把它留作遗漏。

无论你选择什么,上面形态的两个属性值得保留。将路由留在数据库事务之外,这样分页中断永远不会回滚已存储的记录。将 alerts 表视为事实来源而不是通知,这样更改目的地是一个函数的更改,不会有命中依赖于可能没有发生的交付。

一个步骤通常属于这里,本示例没有尝试:在路由之前将被标记的地址与你自己的 KYC 记录关联,这样警报到达时命名内部客户而不是裸地址。那个映射存在于你的系统中而不是链上,这就是它位于此管道之外的原因。

返回 500 是正确的行为。 捕获错误、记录并返回 200 让管道继续前进是很诱人的。不要这样做。2xx 告诉 Streams 该 slot 已安全存储并让它前进,这会把数据库问题转化为记录中一个永久的洞。返回 500 使 Streams 重试然后暂停 Stream,让你留下一个可以看到并恢复的停止管道。

配置 Stream

回到你的 Stream 配置,将目的地设置为 Webhook,使用你的 ngrok URL 加路径,例如 https://your-subdomain.ngrok-free.app/streams。然后设置这些选项:

设置 原因
数据集 区块 完整区块,如 getBlock 返回的
批大小 1 使 batch_start_range 成为区块的 slot
弹性批次 在尖端附近保持批大小为 1
重组时重新流式传输 纠正需求的确定性一半
最新区块延迟 32 概率性一半,大约落后尖端 13 秒
压缩 在摄取服务中少处理一件事

点击"▶️ 测试目的地",然后创建 Stream。几秒钟内 blocks 中应该开始出现行。

验证记录完整

这是将日志转化为证据的查询。每个区块的 previous_blockhash 应该等于它之前区块的 blockhash

SELECT slot AS break_at_slot,
       previous_blockhash AS expected_parent,
       prior_hash         AS actual_prior_block
FROM (
  SELECT slot,
         previous_blockhash,
         LAG(blockhash) OVER (ORDER BY slot) AS prior_hash
  FROM blocks
) linked
WHERE prior_hash IS NOT NULL
  AND previous_blockhash <> prior_hash
ORDER BY slot;

没有行意味着你持有的区块形成一条不间断的链,在第一个和最后一个之间没有缺失。任何行都是缺口的远边:break_at_slot 处的区块期望一个你没有的父区块。

这正是 slot 编号检查不起作用的地方。跳过的 slot 根本不会产生区块,所以连续的区块通过哈希跨越它们保持链接,slot 编号的跳跃不是任何事情的证据。哈希连贯性完全忽略 slot 编号,回答你真正关心的那个问题,即你是否持有存在的每个区块。

image.png

这个查询不会告诉你两件事,值得明确。它无法检测区间任一端点的缺口,因为没有远端的东西可以链接,所以将它与你 Stream 开始的 slot 和当前链尖端配对。它也没有说明区块内的单个交易是否被正确过滤,这是与完整性分开的问题。

按计划运行它并对任何行发出警报。它很便宜,当有人问你的监控系统在给定 slot 看到了什么时,它是最接近直接答案的东西。

还有三个查询值得手头保留。第一个直接回答筛查来源问题,这是审计人员抽样你的日志时实际上会问的。因为被引用的修订版本已归档,查询不仅报告存储的结论,它重新推导它:

-- 我们对此地址得出了什么结论,针对哪个修订版本,而
-- 该修订版本是否实际支持该结论?
SELECT s.slot, s.signature, s.address, s.sanctioned, s.list_version,
       v.content_hash,
       (s.address = ANY (v.addresses)) AS listed_in_that_revision,
       b.block_time, b.corrected_at
FROM screened_accounts s
JOIN blocks b USING (slot)
LEFT JOIN watchlist_versions v ON v.list_version = s.list_version
WHERE s.address = '9kwU8PYhsmRfgS3nwnzT3TvnDeuvdbMAXqWsri2X8rAU'
ORDER BY s.slot;

第二个是连贯性检查,也是需要按计划运行的那个。它捕获筛查状态出错的两种方式:归档列表不支持的结论,以及引用从未归档的修订版本的记录:

SELECT s.slot, s.signature, s.address, s.sanctioned, s.list_version,
       CASE WHEN v.list_version IS NULL THEN '没有归档的修订版本'
            ELSE '结论与归档列表不一致' END AS problem
FROM screened_accounts s
LEFT JOIN watchlist_versions v ON v.list_version = s.list_version
WHERE v.list_version IS NULL
   OR s.sanctioned <> (s.address = ANY (v.addresses))
ORDER BY s.slot;

零行意味着表中的每个结论都可以从归档列表独立重现。这比"我们标了一个版本"更强的陈述,也是能在有人故意试图在你的过程中戳洞时幸存下来的那一个。非零行指向一个特定、可修复的问题:通常是观察名单同步中的发布顺序窗口,或因为过滤器在 sanctions-list-version 存在之抢跑而被标记为 unversioned 的行。两者都通过重放受影响的 slot 来解决。

第三个找到被记录但从未交付的警报,这是分页中断的恢复路径:

SELECT alert_key, address, signature, slot, list_version, raised_at
FROM alerts
WHERE dispatched_at IS NULL
ORDER BY raised_at;

每行存储 list_version 而不是全局解析它,是使这种方法工作的原因。当列表变化时,旧记录保持它们当时得到的结论,所以新列出的地址不会追溯性地重写你声称知道的内容。然后用更新的列表重新筛查过去的窗口是一个刻意的填补缺口重放,而不是一个静默的副作用。

关闭已知缺口

当连贯性检查找到中断时,或者当你需要一个早于 Stream 的回溯窗口时,在明确的 slot 区间上创建第二个 Stream 并指向同一个 webhook:

curl -X POST "https://api.quicknode.com/streams/rest/v1/streams" \
  -H "Content-Type: application/json" \
  -H "x-api-key: $QN_API_KEY" \
  -d "{
    \"name\": \"solana-aml-replay-436190000-436191000\",
    \"network\": \"solana-mainnet\",
    \"dataset\": \"block\",
    \"region\": \"usa_east\",
    \"filter_function\": \"$FILTER_B64\",
    \"start_range\": 436190000,
    \"end_range\": 436191000,
    \"dataset_batch_size\": 1,
    \"elastic_batch_enabled\": false,
    \"destination\": \"webhook\",
    \"destination_attributes\": {
      \"url\": \"https://your-subdomain.ngrok-free.app/streams\",
      \"compression\": \"none\",
      \"max_retry\": 3,
      \"retry_interval_sec\": 5,
      \"post_timeout_sec\": 30
    },
    \"status\": \"active\"
  }"

相同的过滤器,相同的目的地,相同的表,所以恢复的记录在形态上与实时捕获的记录相同,没有单独的回填路径需要推理。因为摄取服务重写它被给予的任何区间,你可以在部分持有的 slot 上运行重放,而不必先检查哪些。这里省略了 fix_block_reorgs,因为固定的历史区间已经是最终的。Stream 在 end_range 处自行结束。

一件需要刻意注意的事:重放针对当前的观察名单进行筛查,而不是这些 slot 最初产生时的情况。写入恢复行的 list_version 将是今天的,这是正确的,并且当重放的原因是新列出的地址时正是你想要的。这确实意味着在你已经持有的 slot 上重放是重新筛查,而不仅仅是修复,所以该区间以前的结论会被覆盖。如果你需要保留早期的结论,在重放之前为区间快照 screened_accounts

之后重新运行连贯性查询。中断应该消失了。

Solana 重放需要付费计划

免费试用账户只能创建跟随链尖端的 Solana Streams。历史 slot 区间(缺口关闭所依赖的)在付费计划上可用。Streams 按处理的区块计费 API 积分,从与 RPC 相同的池中提取,过滤不会降低该成本,因为区块仍然必须被获取和处理。在开始之前,使用 Streams 计费文档中的计算器估算重放区间。

总结

这个架构产生的是一个在底层链变化时自纠正的 Solana 交易日志,以及一个基于实际交易数据触发而不是相信程序自身执行已经工作的制裁命中。使它可靠的属性是不引人注目的:重新交付是幂等的,每个筛查结果都携带产生它的列表版本,纠正留下时间戳,失败的写入会停止管道而不是跳过 slot。

这是一个监控程序的一层,而不是整个程序。警报阈值、案件管理、调查员工作流、报告、保留计划以及所有这些背后的政策都在此管道之外,确认完成的系统满足你的义务只有你和你的合规团队能做的工作。

审计人员还会询问你控件下面的基础设施。Quicknode 维持 SOC 1 Type 2、SOC 2 Type 2 和 ISO/IEC 27001 认证。详细信息在安全页面上。

后续步骤

  • 决定观察名单来自哪里,无论是对 OFAC SDN 列表 的解析器还是商业提供商,以及同步频率
  • 规划被标记地址如何与你的 KYC 记录关联,这样警报到达时带有内部身份
  • 阅读如何流式传输 Solana 程序数据了解围绕程序塑造 Solana Stream 的其他方式,包括将捕获范围扩大到一次多个程序

有问题或想讨论合规数据管道?加入我们的 Discord、关注 @Quicknode、或直接联系我们。

常见问题

Solana 有链重组吗,为什么 AML 监控需要处理它?

Solana 不会重组已最终确定的区块,但在最终确定之前,集群可能在一个竞争分叉上达成共识,所以在尖端附近观察到的区块可能不在规范链上。slot 也可能被跳过,这使缺失的 slot 编号有歧义:它可能正常,也可能是你记录中的一个缺口。对于 AML 监控,两种情况都重要,因为包含非规范交易的记录与缺少规范交易的记录一样是问题。Streams 通过将新区块的父哈希与它之前流式传输的哈希进行比较来检测分叉,然后重新交付带有 reorgs 和 blocks_reorged 元数据标记的受影响区间。

如果程序已经在链上阻止了受制裁地址,为什么还要外部筛查交易?

两个原因。链上黑名单的时效性取决于其最后一次部署,而 Key-Value Store 中的外部列表可以随着你的制裁信息源发布而频繁更新。而且攻击者可能完全绕过了程序的检查,通过一个未被标记的钱包、一个中间跳转或程序自身逻辑中的一个缺口。筛查实际交易数据独立于程序执行是否工作。与它正在检查的事物共享失败模式的控制不是控制。

重新流式传输的 slot 会让我的合规团队对同一笔交易收到两次警报吗?

不会,只要警报去重与数据纠正保持分离。这种方法给每个警报一个签名和地址的键,使它成为故意排除在重组级联之外的 alerts 表的主键,并且只在插入实际创建行时发出页面。纠正按需要的频率重写数据集,而警报累积,所以两者不能干扰。去重不会撤销重组后来移除的交易上的警报,这是运行一个有意义的 keep_distance_from_tip 的原因之一。

在 Solana 上运行 AML 交易监控需要多少钱?

Streams 按处理的区块计费 API 积分,使用取决于网络和数据集的乘数,并从与 RPC 相同的积分池中提取。过滤器不会降低该成本,因为区块仍然必须被获取和处理以应用过滤器逻辑,尽管它们大幅减少带宽和下游存储。因为 Solana 产生区块很快,在激活之前使用 Streams 计费文档中的计算器进行估算。此管道依赖的重组时重新流式传输选项需要付费计划。

以这种方式记录 Solana 交易是否使我的系统 AML 合规?

不。这涵盖了数据层的一个部分:捕获完整、自纠正的交易记录并对照观察名单进行筛查。监控系统是否满足 AML 义务取决于你自己的要求、政策、阈值、调查和报告工作流以及保留规则,这些都不是数据管道决定的。将此视为一个适应示例,并与你的合规团队和法律顾问确认完成的系统。

我可以直接使用数据库目的地而不是 webhook 和摄取服务吗?

可以。Streams 交付到 webhook、S3 兼容存储、PostgreSQL、Azure Blob Storage 和 Kafka,所以数据库可以直接接收记录,中间不需要你自己的服务。权衡是纠正处理:直接目的地让你依赖 upsert 来吸收重新交付的行,而 webhook 和小型摄取服务让你在一个事务内移除被取代的行、记录纠正应用的时间并分发警报。

  • 原文链接: quicknode.com/guides/sol...
  • 鸿途知科网 AI 助手,为大家转译优秀英文文章,如有翻译不通的地方,还请包涵~
版权声明

本文仅代表作者观点,不代表区块链技术网立场。
本文系作者授权本站发表,未经许可,不得转载。

发表评论:

◎欢迎参与讨论,请在这里发表您的看法、交流您的观点。

热门