> ## Documentation Index
> Fetch the complete documentation index at: https://docs.leafage.chaintable.com/llms.txt
> Use this file to discover all available pages before exploring further.

# 数据链路

> 一个区块从写节点执行到可被 eth_call 查询的完整路径，包含 reorg、追赶和冷启动

本页跟着一个区块走完全程。每个阶段都标注了对应的代码位置，便于定位实现。

## 端到端时序

```mermaid theme={null}
sequenceDiagram
    participant N as P2P 网络
    participant G as 写节点
    participant S3 as S3（内/外桶）
    participant K as Kafka 内部 topic
    participant L as leafage-evm
    participant C as consistency-checker
    participant E as etcd
    participant P as nodex-proxy

    N->>G: 新区块
    G->>G: 执行区块，tracer 收集<br/>traces / events / stateDiff
    G->>S3: OnCommit 并发上传<br/>Header、StateDiff、BlockFile、Validation
    G->>K: writeBlockAndSetHead 发布<br/>BlockChangeNotification（仅 Leader）

    par 状态摄入
        K-->>L: 通知
        L->>S3: 拉取 Header + StateDiff
        L->>L: 推入 StateTree 头部
    and 一致性校验
        K-->>C: 同一条通知
        C->>L: 轮询 eth_blockNumber
        C->>S3: 读 / 重写 BlockValidation
        C->>E: 写节点状态 + lastBlockNumber
        C->>C: 发布外部通知
    end

    E-->>P: watch 事件
    P->>L: 路由 eth_call
```

## 阶段 1：执行与追踪

写节点通过 P2P 正常同步和执行区块。pipeline tracer 以 `tracing.Hooks` 的形式挂在 EVM 上，在执行过程中收集数据。

| Hook                    | 时机     | 收集内容                          |
| ----------------------- | ------ | ----------------------------- |
| `OnBlockStart`          | 区块执行开始 | 初始化 `BlockCtx`，创建 `BlockFile` |
| `OnTxStart` / `OnTxEnd` | 每笔交易   | 交易元数据，展平调用树                   |
| `OnEnter` / `OnExit`    | 每个调用帧  | 构建调用树，标记失败传播                  |
| `OnOpcode`              | 每条指令   | 识别 `SSTORE`，预加载 prestate      |
| `OnLog`                 | 事件产生   | 记录事件及其所属 trace                |
| `OnCommit`              | 状态提交后  | 生成 `BlockStorageDiff`，触发上传    |

`OnCommit` 和 `OnBlockDBStart` 是 Leafage 对上游 Geth 的扩展，定义在 `core/tracing/hooks.go`，由 `core/blockchain.go` 在 `ProcessBlock` 中分发。

<Note>
  StateDiff 有两种获取方式：默认从 `OnCommit` 拿到的 StateDB 提交集直接转换；配置 `enable_prestate_tracer` 后改用 prestate tracer 在指令级采集。前者更准确也更省开销，后者用于没有 commit hook 的客户端。
</Note>

## 阶段 2：序列化与分发

`OnCommit` 里并发执行四路上传（`tracer/pipeline_tracer.go`）：

| 数据              | 桶   | 键                                             | 格式          | 消费者                                |
| --------------- | --- | --------------------------------------------- | ----------- | ---------------------------------- |
| Header          | 内部桶 | `{chainID}[/{version}]/{blockHash}/block`     | JSON + gzip | leafage-evm                        |
| StateDiff       | 内部桶 | `{chainID}[/{version}]/{stateRoot}/stateDiff` | RLP         | leafage-evm                        |
| BlockFile       | 外部桶 | `{chainID}[/{version}]/{blockHash}`           | JSON + gzip | 外部消费者                              |
| BlockValidation | 外部桶 | `{chainID}[/{version}]/{height}/{blockHash}`  | JSON + gzip | consistency-checker、leafage-evm 追赶 |

上传完成后，写节点在确定新的 canonical head 时（`core/blockchain.go` 的 `writeBlockAndSetHead`）向 Kafka 发布通知。发布前会用 `getCommonAncestor` 对比上一条已发布的区块与当前 head：

* 只有新块 → `changeType: 1`，`newBlocks` 按高度升序
* 存在被丢弃的分支 → `changeType: 2`，同时带 `dropBlocks` 和 `newBlocks`

只有 Leader 实例发布 Kafka 消息，备节点仍然上传 S3。Leader 由 pipeline 的 etcd 选举决定。

<Warning>
  S3 上传和 Kafka 发布不是一个原子操作。消费者收到通知时对应的 S3 对象通常已就绪，但实现上必须容忍短暂的 404 并重试——leafage-evm 和 consistency-checker 都带退避重试。
</Warning>

## 阶段 3：状态摄入

leafage-evm 的 `KafkaUpdater`（`bin/leafage-evm/src/updater/kafka_updater.rs`）消费通知：

<Steps>
  <Step title="解析通知">
    从 `newBlocks` 拿到区块哈希与父哈希。
  </Step>

  <Step title="并行拉取 S3">
    按哈希取 Header，按 state root 取 StateDiff。父块与当前块 state root 相同时跳过 diff 拉取，视为空差异。
  </Step>

  <Step title="更新 StateTree">
    调用 `tree.update_block(block_info, block_diff)`，把新差异层挂到链表头部。
  </Step>

  <Step title="提交 offset">
    落盘成功后写入 `offset_dir`，用于崩溃恢复。
  </Step>
</Steps>

StateTree 是一个差异层链表：每层只保存相对父层的变更。查询时从头部向下遍历，命中即返回；全部未命中则落到 RocksDB。

## 阶段 4：终结化

当区块深度超过 `--diff-depth-limit`（默认 64）时，最老的一层被刷写到 RocksDB 并从内存移除。

* **State 节点**：直接覆盖写 `address` / `address || slot`，只保留最新值。
* **Archive 节点**：双写。`address || block_num` 保留历史版本，`address || u64::MAX` 作为最新值快速路径。历史查询用 RocksDB 的 `seek_for_prev` 定位不大于目标高度的最近版本。

## 阶段 5：查询服务

```mermaid theme={null}
flowchart LR
    C["客户端"] -->|"POST /:chainId"| P["nodex-proxy"]
    P -->|"latest / 距链头 ≤64"| S["State 节点"]
    P -->|"距链头 >64"| A["Archive 节点"]
    P -->|"错误码 -39008"| N["Native 节点"]
    S -->|"-39006 重试"| A
```

nodex-proxy 解析请求里的区块参数，结合 etcd 中的 `lastBlockNumber` 决定节点池，再按权重或轮询选出具体节点。leafage-evm 在本地状态上用 revm 执行并返回结果。

两个自动故障转移路径：

| 错误码      | 含义                                    | 重试目标                 |
| -------- | ------------------------------------- | -------------------- |
| `-39006` | `StateBlockNotFound`，State 节点没有该高度的状态 | Archive 节点           |
| `-39008` | `CosmosPrecompile`，需要原链能力             | Native 节点（路径改写为 `/`） |

## 并行链路：一致性校验

consistency-checker 消费同一条 Kafka 通知，但不在 leafage-evm 的数据路径上。它的处理顺序（`check/check.go` 的 `Process`）：

<Steps>
  <Step title="消息校验与去重">
    重复消息或已处理过的消息走对齐后直接推进，避免重试死锁。
  </Step>

  <Step title="预取 S3">
    在等待副本的同时并行预取 BlockValidation 和同高度的键列表。
  </Step>

  <Step title="等待副本收敛">
    以 `check_interval_ms`（默认 20ms）轮询所有副本的 `eth_blockNumber`，直到 `ready_ratio`（默认 0.8）的副本追上该高度，或超过 `check_timeout_ms`（默认 2000ms）判失败。
  </Step>

  <Step title="写入本地 Pebble DB">
    双索引：`h{hash}` → BlockInfo，`n{number}` → hash，供 JSON-RPC 查询。
  </Step>

  <Step title="标记分叉">
    列出 S3 外部桶中同高度的所有 BlockValidation 对象，把非规范的重写为 `is_fork: true`。
  </Step>

  <Step title="发布外部通知">
    先发 drop 通知，再发新块通知，写入外部 Kafka topic。
  </Step>

  <Step title="发布后复核">
    用一次新鲜的 LIST 再检查一遍同高度的分叉标记，捕捉等待副本期间新上传的对象。
  </Step>
</Steps>

同时，每一轮轮询的结果都会写回 etcd：节点的 `stateType`（1 追上 / 2 落后 / 3 离线）以及链的 `lastBlockNumber`。这就是 nodex-proxy 的路由依据。

## Reorg 的三处处理

Leafage 在三个层面分别处理链重组，理解它们的分工可以避免误判问题出在哪一层。

<AccordionGroup>
  <Accordion title="写节点：产生 dropBlocks">
    `writeBlockAndSetHead` 用 `getCommonAncestor` 比较上一条已发布的区块和新 head，把从共同祖先到旧 head 的路径作为 `dropBlocks`，到新 head 的路径作为 `newBlocks`，以 `changeType: 2` 一次性发出。
  </Accordion>

  <Accordion title="leafage-evm：分叉层与安全回补">
    StateTree 用两张表索引差异层：`hash_diff_map` 保存所有块（含分叉块），`num_diff_map` 只跟踪规范链。因此按哈希查询仍可访问分叉状态，按高度查询始终走规范链。

    S3 追赶时靠近链头的区块可能落在错误的分支上。`--catchup-safe-depth` 指定链头附近多少个区块改为沿 Kafka 通知里的精确 parent-hash 链回补，而不是按高度索引查。该值应大于目标链的最大 reorg 深度（例如 Moonriver 设 64），`0` 表示禁用。
  </Accordion>

  <Accordion title="consistency-checker：标记与外部通知">
    收到 `changeType: 2` 时先重写 `dropBlocks` 对应的 BlockValidation 为 `is_fork: true`，再扫描同高度的其他对象。另有周期巡检（`fork_scan_interval_sec`，默认 60 秒，回看 `fork_scan_lookback` 个高度）兜底。外部消费者通过 `OuterBlockChangeNotification` 的 `is_fork` 字段感知。
  </Accordion>
</AccordionGroup>

## 冷启动与追赶

leafage-evm 启动时按持久化的 Kafka offset 决定路径：

```text theme={null}
读取 offset_dir 中的 offset
  ├─ offset >= Kafka 最低水位 → 从 offset 继续消费
  └─ offset 缺失或已过期      → 先从 S3 追赶，再从最新位置开始消费
```

S3 追赶按高度逐块回补：先在外部桶用前缀 `{chainID}[/{version}]/{height}/` 列出该高度的对象，取 `is_fork` 为 false 的那个哈希，再按哈希去内部桶取 Header 和 StateDiff。批大小由 `--init-task-queue-size`（默认 256）控制。

<Tip>
  这里是 consistency-checker 与 leafage-evm 的隐式契约：分叉标记写在外部桶的 BlockValidation 上，leafage-evm 追赶时依赖它来选出规范分支。如果 checker 停摆，同高度存在多个对象时追赶会退化到逐个读取判断。
</Tip>

全新节点的常规恢复方式是下载 RocksDB 快照后再追赶 Kafka，通常在分钟级完成，而不是从创世块重放。

## 可观测点

排查“数据没到”时，按链路顺序检查这些指标：

| 环节                  | 指标                                                                               | 组件                  |
| ------------------- | -------------------------------------------------------------------------------- | ------------------- |
| 区块执行完成              | `LatestBlockNumber`                                                              | pipeline            |
| S3 上传完成             | `LatestUploadedBlockNumber`，与上一项的差值反映上传积压                                        | pipeline            |
| Kafka 发布耗时          | `BlockPushTimer`                                                                 | pipeline            |
| 副本等待时长              | `pipeline_replica_ready_wait_seconds`、`pipeline_replica_ready_timeouts_total`    | consistency-checker |
| 已确认高度               | `pipeline_block_num`                                                             | consistency-checker |
| 写节点到外部 Kafka 的端到端延迟 | `pipeline_block_ingress_to_outer_kafka_seconds`                                  | consistency-checker |
| 分叉标记异常              | `pipeline_fork_scan_rewrites_total`、`pipeline_drop_block_rewrite_failures_total` | consistency-checker |
| 查询失败率               | `jrpcx_rpc_calls_failed`                                                         | nodex-proxy         |
