1use 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
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 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 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 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 self.subscription_gap_queue().mark_done();
430 }
431
432 log::info!(target: LOG_TARGET, "๐ Subscription gap filler stopped");
433 }
434}