referrerpolicy=no-referrer-when-downgrade

pallet_revive_eth_rpc/
block_sync.rs

1// This file is part of Substrate.
2
3// Copyright (C) Parity Technologies (UK) Ltd.
4// SPDX-License-Identifier: Apache-2.0
5
6// Licensed under the Apache License, Version 2.0 (the "License");
7// you may not use this file except in compliance with the License.
8// You may obtain a copy of the License at
9//
10// 	http://www.apache.org/licenses/LICENSE-2.0
11//
12// Unless required by applicable law or agreed to in writing, software
13// distributed under the License is distributed on an "AS IS" BASIS,
14// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
15// See the License for the specific language governing permissions and
16// limitations under the License.
17
18//! Historic block syncing logic for the Ethereum JSON-RPC server.
19
20use crate::{
21	BlockInfoProvider,
22	client::{Client, ClientError, GapFillRequest, SubstrateBlockNumber},
23};
24use pallet_revive::evm::H256;
25use tokio::sync::mpsc;
26
27const LOG_TARGET: &str = "eth-rpc::block-sync";
28
29/// Trait for types that can be used as keys in the `sync_state` table.
30pub trait SyncStateKey: std::fmt::Display {}
31
32/// Labels used to track sync progress in the `sync_state` table.
33#[derive(Debug, Clone, Copy, derive_more::Display)]
34pub enum SyncLabel {
35	/// Lowest synced block. Only decreases.
36	#[display(fmt = "sync-tail")]
37	Tail,
38	/// Highest synced block. Absent means no sync has started.
39	/// During backfill: upper boundary being filled.
40	/// After backfill: advanced by the finalized-block subscription.
41	#[display(fmt = "sync-head")]
42	Head,
43}
44
45/// Chain metadata stored in the `sync_state` table.
46#[derive(Debug, Clone, Copy, derive_more::Display)]
47pub enum ChainMetadata {
48	/// Genesis block hash โ€” used for chain identity verification.
49	#[display(fmt = "chain-genesis")]
50	Genesis,
51	/// Auto-discovered first EVM block on the chain.
52	#[display(fmt = "chain-first-evm-block")]
53	FirstEvmBlock,
54}
55
56impl SyncStateKey for SyncLabel {}
57impl SyncStateKey for ChainMetadata {}
58
59/// Sync checkpoint persisted in the `sync_state` table to allow resuming after a restart.
60#[derive(Debug, Clone, Copy, Eq, PartialEq)]
61pub struct SyncCheckpoint {
62	pub block_number: SubstrateBlockNumber,
63	pub block_hash: Option<H256>,
64}
65
66impl SyncCheckpoint {
67	/// Create a checkpoint with a known block hash.
68	pub fn new(block_number: SubstrateBlockNumber, block_hash: H256) -> Self {
69		Self { block_number, block_hash: Some(block_hash) }
70	}
71
72	/// Create a checkpoint with only a block number (no hash).
73	pub fn from_number(block_number: SubstrateBlockNumber) -> Self {
74		Self { block_number, block_hash: None }
75	}
76}
77
78/// How often (in blocks) the backward sync checkpoints are persisted to the database.
79const BLOCK_INTERVAL: u32 = 128;
80
81/// Options for [`Client::sync_backward_range`].
82struct BackwardSyncRange {
83	from: SubstrateBlockNumber,
84	to: SubstrateBlockNumber,
85	/// Set `Head` label after syncing the first block.
86	set_head: bool,
87	/// Checkpoint `Tail` label periodically and at end.
88	checkpoint_tail: bool,
89	/// When true, persist the first EVM block boundary if a non-EVM block is encountered.
90	persist_first_evm_block: bool,
91}
92
93impl Client {
94	/// Verify that the stored genesis hash matches the connected chain.
95	async fn validate_chain_identity(&self) -> Result<H256, ClientError> {
96		let genesis_hash: H256 = self.api().genesis_hash();
97
98		if let Some(checkpoint) =
99			self.receipt_provider().get_sync_label(ChainMetadata::Genesis).await?
100		{
101			if let Some(stored) = checkpoint.block_hash {
102				if stored != genesis_hash {
103					return Err(ClientError::ChainMismatch);
104				}
105			}
106		}
107
108		Ok(genesis_hash)
109	}
110
111	/// Verify that a stored boundary block still exists on the finalized chain.
112	async fn verify_boundary(&self, checkpoint: &SyncCheckpoint) -> Result<(), ClientError> {
113		let num = checkpoint.block_number;
114		let hash = checkpoint.block_hash;
115		match (num, hash) {
116			(_, None) => {
117				log::error!(target: LOG_TARGET,
118					"Boundary #{num}: missing stored hash");
119				Err(ClientError::SyncBoundaryMismatch)
120			},
121			(_, Some(stored_hash)) => {
122				let block = self.block_provider().block_by_number(num).await?.ok_or_else(|| {
123					log::error!(target: LOG_TARGET,
124						"Boundary #{num}: block not found on chain \
125						 (node may have pruned it โ€” use an archive node with --eth-pruning archive)");
126					ClientError::SyncBoundaryMismatch
127				})?;
128				if block.block_hash() != stored_hash {
129					log::error!(target: LOG_TARGET,
130						"Boundary #{num}: hash mismatch โ€” stored {stored_hash:?}, \
131						 chain {:?}", block.block_hash());
132					return Err(ClientError::SyncBoundaryMismatch);
133				}
134				Ok(())
135			},
136		}
137	}
138
139	/// Checkpoint the given sync label to the DB.
140	async fn checkpoint_sync_label(&self, label: SyncLabel, num: SubstrateBlockNumber, hash: H256) {
141		let cp = SyncCheckpoint::new(num, hash);
142		let result = match label {
143			SyncLabel::Head => self.receipt_provider().advance_sync_label(label, cp).await,
144			SyncLabel::Tail => self.receipt_provider().recede_sync_label(label, cp).await,
145		};
146		if let Err(err) = result {
147			log::warn!(target: LOG_TARGET, "Failed to update sync_label[{label}]: {err:?}");
148		}
149	}
150
151	/// Backward sync historical blocks from the latest finalized block to the first EVM block.
152	/// Resumes from the last checkpoint if a previous sync was interrupted.
153	/// Fatal errors (chain/DB mismatch) are propagated; transient errors are swallowed
154	/// to avoid taking down the RPC server.
155	pub async fn sync_backward(&self) -> Result<(), ClientError> {
156		log::info!(target: LOG_TARGET,
157			"๐Ÿ”„ Historical block sync enabled. \
158			 For a complete sync, the connected node should be an archive node.");
159		match self.sync_backward_inner().await {
160			Ok(()) => Ok(()),
161			Err(err) if err.is_chain_validation_error() => Err(err),
162			Err(err) => {
163				log::error!(target: LOG_TARGET, "๐Ÿ—„๏ธ Sync stopped due to {err}.");
164				Ok(())
165			},
166		}
167	}
168
169	async fn sync_backward_inner(&self) -> Result<(), ClientError> {
170		let genesis_hash = self.validate_chain_identity().await?;
171		let latest_finalized_block = self.latest_finalized_block().await;
172		let latest_finalized = SyncCheckpoint::new(
173			latest_finalized_block.block_number(),
174			latest_finalized_block.block_hash(),
175		);
176
177		// Store genesis (idempotent).
178		self.receipt_provider()
179			.set_sync_label(ChainMetadata::Genesis, SyncCheckpoint::new(0, genesis_hash))
180			.await?;
181
182		let (head, tail) = tokio::try_join!(
183			self.receipt_provider().get_sync_label(SyncLabel::Head),
184			self.receipt_provider().get_sync_label(SyncLabel::Tail),
185		)?;
186
187		match (tail, head) {
188			(Some(tail), Some(head)) => {
189				// Verify boundary hashes still match the finalized chain.
190				tokio::try_join!(self.verify_boundary(&tail), self.verify_boundary(&head),)?;
191				self.sync_backward_resume(tail, head, latest_finalized).await?;
192			},
193			(Some(_), None) => {
194				log::warn!(target: LOG_TARGET,
195					"๐Ÿ—„๏ธ Tail exists without Head โ€” possible partial corruption, \
196					 starting fresh sync from #{}", latest_finalized.block_number);
197				self.sync_backward_fresh(latest_finalized.block_number).await?;
198			},
199			_ => {
200				log::info!(target: LOG_TARGET,
201					"๐Ÿ—„๏ธ Fresh sync: syncing backward from #{}", latest_finalized.block_number);
202				self.sync_backward_fresh(latest_finalized.block_number).await?;
203			},
204		}
205
206		self.mark_backfill_complete();
207
208		log::info!(target: LOG_TARGET, "๐Ÿ—„๏ธ Historic sync complete");
209		Ok(())
210	}
211
212	/// Backward sync from `latest_finalized` down to the first EVM block.
213	async fn sync_backward_fresh(
214		&self,
215		latest_finalized: SubstrateBlockNumber,
216	) -> Result<(), ClientError> {
217		let first_evm = self.receipt_provider().first_evm_block().unwrap_or(0);
218		self.sync_backward_range(BackwardSyncRange {
219			from: latest_finalized,
220			to: first_evm,
221			set_head: true,
222			checkpoint_tail: true,
223			persist_first_evm_block: true,
224		})
225		.await
226	}
227
228	/// Resume backward sync by filling the top gap (new blocks) and bottom gap (backfill).
229	async fn sync_backward_resume(
230		&self,
231		tail: SyncCheckpoint,
232		head: SyncCheckpoint,
233		latest_finalized: SyncCheckpoint,
234	) -> Result<(), ClientError> {
235		log::info!(target: LOG_TARGET,
236			"๐Ÿ—„๏ธ Resuming sync: DB has blocks #{}..#{}, chain head is #{}",
237			tail.block_number, head.block_number, latest_finalized.block_number);
238
239		let top_gap = async {
240			// Top gap: sync from latest_finalized down to head + 1.
241			if head.block_number < latest_finalized.block_number {
242				self.sync_backward_range(BackwardSyncRange {
243					from: latest_finalized.block_number,
244					to: head.block_number.saturating_add(1),
245					set_head: false,
246					checkpoint_tail: false,
247					persist_first_evm_block: false,
248				})
249				.await?;
250
251				// Mark top gap complete so a restart won't redo it.
252				self.receipt_provider()
253					.advance_sync_label(SyncLabel::Head, latest_finalized)
254					.await?;
255			}
256			Ok::<_, ClientError>(())
257		};
258
259		let bottom_gap = async {
260			// Bottom gap: sync from tail - 1 down to the first EVM block.
261			let first_evm = self.receipt_provider().first_evm_block().unwrap_or(0);
262			if tail.block_number > first_evm {
263				self.sync_backward_range(BackwardSyncRange {
264					from: tail.block_number.saturating_sub(1),
265					to: first_evm,
266					set_head: false,
267					checkpoint_tail: true,
268					persist_first_evm_block: true,
269				})
270				.await?;
271			} else {
272				log::debug!(target: LOG_TARGET, "๐Ÿ—„๏ธ No backward gap to fill");
273			}
274			Ok::<_, ClientError>(())
275		};
276
277		tokio::try_join!(top_gap, bottom_gap)?;
278
279		Ok(())
280	}
281
282	/// Backward sync from block `from` down to block `to` (inclusive).
283	/// Stops early if a non-EVM block is discovered (auto-discovery of first EVM block).
284	async fn sync_backward_range(
285		&self,
286		BackwardSyncRange {
287			from,
288			to,
289			set_head,
290			checkpoint_tail,
291			persist_first_evm_block,
292		}: BackwardSyncRange,
293	) -> Result<(), ClientError> {
294		if from < to {
295			log::debug!(target: LOG_TARGET,	"โฌ‡๏ธ Backward sync: nothing to sync (#{from}..#{to})");
296			return Ok(());
297		}
298
299		log::info!(target: LOG_TARGET, "โฌ‡๏ธ Backward sync: #{from} down to #{to}");
300
301		let mut block = self
302			.block_provider()
303			.block_by_number(from)
304			.await?
305			.ok_or(ClientError::BlockNotFound)?;
306
307		let mut blocks_synced = 0u64;
308		let mut last_synced: Option<(SubstrateBlockNumber, H256)> = None;
309		let at_checkpoint =
310			|synced: u64| synced <= 1 || synced.is_multiple_of(u64::from(BLOCK_INTERVAL));
311
312		let loop_result: Result<(), ClientError> = loop {
313			let block_number = block.block_number();
314			let block_hash = block.block_hash();
315
316			// A block whose runtime does not expose `eth_block_hash` predates pallet-revive and
317			// is treated exactly like a block without an EVM hash: it marks the end of the
318			// backward sync.
319			let ethereum_hash = match self.runtime_api(block_hash).await {
320				Ok(runtime_api) => {
321					match runtime_api.eth_block_hash(pallet_revive::evm::U256::from(block_number)) {
322						Some(future) => match future.await {
323							Ok(hash) => hash,
324							Err(err) => {
325								log::error!(target: LOG_TARGET,	"โš ๏ธ eth_block_hash failed for #{block_number}: {err:?}, stopping");
326								break Err(err);
327							},
328						},
329						None => None,
330					}
331				},
332				Err(err) => {
333					log::error!(target: LOG_TARGET,	"โš ๏ธ eth_block_hash failed for #{block_number}: {err:?}, stopping");
334					break Err(err);
335				},
336			};
337
338			match ethereum_hash {
339				Some(hash) => {
340					if let Err(err) =
341						self.receipt_provider().insert_block_receipts_past(&block, &hash).await
342					{
343						log::error!(target: LOG_TARGET,
344							"โš ๏ธ Insert failed for #{block_number}: {err:?}, stopping");
345						break Err(err);
346					}
347
348					last_synced = Some((block_number, block_hash));
349					blocks_synced += 1;
350
351					if blocks_synced == 1 && set_head {
352						self.checkpoint_sync_label(SyncLabel::Head, block_number, block_hash).await;
353					}
354
355					if at_checkpoint(blocks_synced) {
356						log::debug!(target: LOG_TARGET,
357							"โฌ‡๏ธ Backward sync progress: #{block_number} ({blocks_synced} blocks synced)");
358						if checkpoint_tail {
359							self.checkpoint_sync_label(SyncLabel::Tail, block_number, block_hash)
360								.await;
361						}
362					}
363				},
364				None => {
365					if persist_first_evm_block {
366						let first_evm_block = block_number.saturating_add(1);
367						log::debug!(target: LOG_TARGET,
368							"๐Ÿ” No EVM hash at #{block_number}, setting first_evm_block to #{first_evm_block}");
369						if let Err(err) =
370							self.receipt_provider().set_first_evm_block(first_evm_block).await
371						{
372							log::warn!(target: LOG_TARGET, "Failed to persist first-evm-block: {err:?}");
373						}
374					} else {
375						log::debug!(target: LOG_TARGET,
376							"๐Ÿ” No EVM hash at #{block_number}, skipping first EVM block update");
377					}
378
379					break Ok(());
380				},
381			}
382
383			if block_number > to {
384				let parent_hash = block.block_header().await?.parent_hash;
385				match self
386					.block_provider()
387					.block_by_hash(&parent_hash)
388					.await
389					.map_err(Into::into)
390					.and_then(|opt| opt.ok_or(ClientError::BlockNotFound))
391				{
392					Ok(b) => block = b,
393					Err(err) => {
394						log::error!(target: LOG_TARGET,
395							"โš ๏ธ Could not fetch parent of #{block_number}: {err:?}, stopping");
396						break Err(err);
397					},
398				}
399			} else {
400				break Ok(());
401			}
402		};
403
404		// Checkpoint the last synced block if it wasn't already at a checkpoint interval.
405		if loop_result.is_ok() && checkpoint_tail && !at_checkpoint(blocks_synced) {
406			if let Some((num, hash)) = last_synced {
407				self.checkpoint_sync_label(SyncLabel::Tail, num, hash).await;
408			}
409		}
410
411		log::info!(target: LOG_TARGET,
412			"โฌ‡๏ธ Backward sync: {blocks_synced} blocks synced \
413			 (requested #{from}..#{to})");
414
415		loop_result
416	}
417
418	/// Run the background subscription gap filler, processing requests sequentially.
419	pub(crate) async fn run_subscription_gap_filler(&self, mut rx: mpsc::Receiver<GapFillRequest>) {
420		log::info!(target: LOG_TARGET, "๐Ÿ”„ Subscription gap filler started");
421
422		while let Some(GapFillRequest { from_inclusive, to_inclusive }) = rx.recv().await {
423			log::info!(target: LOG_TARGET, "๐Ÿ”„ Subscription gap filler: processing #{from_inclusive} down to #{to_inclusive}");
424			if let Err(err) = self
425				.sync_backward_range(BackwardSyncRange {
426					from: from_inclusive,
427					to: to_inclusive,
428					set_head: false,
429					checkpoint_tail: false,
430					persist_first_evm_block: false,
431				})
432				.await
433			{
434				log::error!(target: LOG_TARGET, "๐Ÿ”„ Subscription gap fill failed for #{from_inclusive}..#{to_inclusive}: {err:?}");
435			} else {
436				log::info!(target: LOG_TARGET, "๐Ÿ”„ Subscription gap filler: done with #{from_inclusive}..#{to_inclusive}");
437			}
438			// Mark done unconditionally โ€” mirrors how subscribe_new_blocks handles
439			// callback errors: log and move on.
440			self.subscription_gap_queue().mark_done();
441		}
442
443		log::info!(target: LOG_TARGET, "๐Ÿ”„ Subscription gap filler stopped");
444	}
445}