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, runtime_api::RuntimeApi},
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			let ethereum_hash = match RuntimeApi::new((*block).clone())
317				.eth_block_hash(pallet_revive::evm::U256::from(block_number))
318				.await
319			{
320				Ok(h) => h,
321				Err(err) => {
322					log::error!(target: LOG_TARGET, "โš ๏ธ eth_block_hash failed for #{block_number}: {err:?}, stopping");
323					break Err(err);
324				},
325			};
326
327			match ethereum_hash {
328				Some(hash) => {
329					if let Err(err) =
330						self.receipt_provider().insert_block_receipts_past(&block, &hash).await
331					{
332						log::error!(target: LOG_TARGET,
333							"โš ๏ธ Insert failed for #{block_number}: {err:?}, stopping");
334						break Err(err);
335					}
336
337					last_synced = Some((block_number, block_hash));
338					blocks_synced += 1;
339
340					if blocks_synced == 1 && set_head {
341						self.checkpoint_sync_label(SyncLabel::Head, block_number, block_hash).await;
342					}
343
344					if at_checkpoint(blocks_synced) {
345						log::debug!(target: LOG_TARGET,
346							"โฌ‡๏ธ Backward sync progress: #{block_number} ({blocks_synced} blocks synced)");
347						if checkpoint_tail {
348							self.checkpoint_sync_label(SyncLabel::Tail, block_number, block_hash)
349								.await;
350						}
351					}
352				},
353				None => {
354					if persist_first_evm_block {
355						let first_evm_block = block_number.saturating_add(1);
356						log::debug!(target: LOG_TARGET,
357							"๐Ÿ” No EVM hash at #{block_number}, setting first_evm_block to #{first_evm_block}");
358						if let Err(err) =
359							self.receipt_provider().set_first_evm_block(first_evm_block).await
360						{
361							log::warn!(target: LOG_TARGET, "Failed to persist first-evm-block: {err:?}");
362						}
363					} else {
364						log::debug!(target: LOG_TARGET,
365							"๐Ÿ” No EVM hash at #{block_number}, skipping first EVM block update");
366					}
367
368					break Ok(());
369				},
370			}
371
372			if block_number > to {
373				let parent_hash = block.block_header().await?.parent_hash;
374				match self
375					.block_provider()
376					.block_by_hash(&parent_hash)
377					.await
378					.map_err(Into::into)
379					.and_then(|opt| opt.ok_or(ClientError::BlockNotFound))
380				{
381					Ok(b) => block = b,
382					Err(err) => {
383						log::error!(target: LOG_TARGET,
384							"โš ๏ธ Could not fetch parent of #{block_number}: {err:?}, stopping");
385						break Err(err);
386					},
387				}
388			} else {
389				break Ok(());
390			}
391		};
392
393		// Checkpoint the last synced block if it wasn't already at a checkpoint interval.
394		if loop_result.is_ok() && checkpoint_tail && !at_checkpoint(blocks_synced) {
395			if let Some((num, hash)) = last_synced {
396				self.checkpoint_sync_label(SyncLabel::Tail, num, hash).await;
397			}
398		}
399
400		log::info!(target: LOG_TARGET,
401			"โฌ‡๏ธ Backward sync: {blocks_synced} blocks synced \
402			 (requested #{from}..#{to})");
403
404		loop_result
405	}
406
407	/// Run the background subscription gap filler, processing requests sequentially.
408	pub(crate) async fn run_subscription_gap_filler(&self, mut rx: mpsc::Receiver<GapFillRequest>) {
409		log::info!(target: LOG_TARGET, "๐Ÿ”„ Subscription gap filler started");
410
411		while let Some(GapFillRequest { from_inclusive, to_inclusive }) = rx.recv().await {
412			log::info!(target: LOG_TARGET, "๐Ÿ”„ Subscription gap filler: processing #{from_inclusive} down to #{to_inclusive}");
413			if let Err(err) = self
414				.sync_backward_range(BackwardSyncRange {
415					from: from_inclusive,
416					to: to_inclusive,
417					set_head: false,
418					checkpoint_tail: false,
419					persist_first_evm_block: false,
420				})
421				.await
422			{
423				log::error!(target: LOG_TARGET, "๐Ÿ”„ Subscription gap fill failed for #{from_inclusive}..#{to_inclusive}: {err:?}");
424			} else {
425				log::info!(target: LOG_TARGET, "๐Ÿ”„ Subscription gap filler: done with #{from_inclusive}..#{to_inclusive}");
426			}
427			// Mark done unconditionally โ€” mirrors how subscribe_new_blocks handles
428			// callback errors: log and move on.
429			self.subscription_gap_queue().mark_done();
430		}
431
432		log::info!(target: LOG_TARGET, "๐Ÿ”„ Subscription gap filler stopped");
433	}
434}