1use 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
29pub trait SyncStateKey: std::fmt::Display {}
31
32#[derive(Debug, Clone, Copy, derive_more::Display)]
34pub enum SyncLabel {
35 #[display(fmt = "sync-tail")]
37 Tail,
38 #[display(fmt = "sync-head")]
42 Head,
43}
44
45#[derive(Debug, Clone, Copy, derive_more::Display)]
47pub enum ChainMetadata {
48 #[display(fmt = "chain-genesis")]
50 Genesis,
51 #[display(fmt = "chain-first-evm-block")]
53 FirstEvmBlock,
54}
55
56impl SyncStateKey for SyncLabel {}
57impl SyncStateKey for ChainMetadata {}
58
59#[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 pub fn new(block_number: SubstrateBlockNumber, block_hash: H256) -> Self {
69 Self { block_number, block_hash: Some(block_hash) }
70 }
71
72 pub fn from_number(block_number: SubstrateBlockNumber) -> Self {
74 Self { block_number, block_hash: None }
75 }
76}
77
78const BLOCK_INTERVAL: u32 = 128;
80
81struct BackwardSyncRange {
83 from: SubstrateBlockNumber,
84 to: SubstrateBlockNumber,
85 set_head: bool,
87 checkpoint_tail: bool,
89 persist_first_evm_block: bool,
91}
92
93impl Client {
94 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 self.subscription_gap_queue().mark_done();
441 }
442
443 log::info!(target: LOG_TARGET, "๐ Subscription gap filler stopped");
444 }
445}