1use crate::{
18 Address, BlockInfoProvider, BlockNumberOrTag, Bytes, ChainMetadata, ClientError, Filter,
19 FilterBlockOption, Log, ReceiptExtractor, ReceiptInfo, SubxtBlockInfoProvider, SyncLabel,
20 SyncStateKey,
21 block_sync::SyncCheckpoint,
22 client::{SubstrateBlock, SubstrateBlockNumber},
23};
24use futures::future::OptionFuture;
25use pallet_revive::evm::TransactionSigned;
26use sp_core::{H256, U256};
27use sqlx::{QueryBuilder, Row, Sqlite, SqlitePool, query};
28use std::{
29 collections::{BTreeMap, HashMap},
30 sync::Arc,
31};
32use tokio::sync::Mutex;
33
34const LOG_TARGET: &str = "eth-rpc::receipt_provider";
35const MAX_LOG_RESULTS: usize = 10_000;
36
37fn parse_log_row(row: sqlx::sqlite::SqliteRow) -> Result<Log, sqlx::Error> {
39 let block_hash: Vec<u8> = row.try_get("block_hash")?;
40 let transaction_index: i64 = row.try_get("transaction_index")?;
41 let log_index: i64 = row.try_get("log_index")?;
42 let address: Vec<u8> = row.try_get("address")?;
43 let block_number: i64 = row.try_get("block_number")?;
44 let transaction_hash: Vec<u8> = row.try_get("transaction_hash")?;
45 let topic_0: Option<Vec<u8>> = row.try_get("topic_0")?;
46 let topic_1: Option<Vec<u8>> = row.try_get("topic_1")?;
47 let topic_2: Option<Vec<u8>> = row.try_get("topic_2")?;
48 let topic_3: Option<Vec<u8>> = row.try_get("topic_3")?;
49 let data: Option<Vec<u8>> = row.try_get("data")?;
50
51 let topics = [topic_0, topic_1, topic_2, topic_3]
52 .iter()
53 .filter_map(|t| t.as_ref().map(|t| H256::from_slice(t)))
54 .collect::<Vec<_>>();
55
56 Ok(Log {
57 address: Address::from_slice(&address),
58 block_hash: H256::from_slice(&block_hash),
59 block_number: U256::from(block_number as u64),
60 data: data.map(Bytes::from),
61 log_index: U256::from(log_index as u64),
62 topics,
63 transaction_hash: H256::from_slice(&transaction_hash),
64 transaction_index: U256::from(transaction_index as u64),
65 removed: false,
66 })
67}
68
69#[derive(Clone)]
71pub struct DbContext {
72 pool: SqlitePool,
73 max_variable_number: usize,
75 tx_insert_chunk_size: usize,
77 log_insert_chunk_size: usize,
79}
80
81impl DbContext {
82 pub const DEFAULT_MAX_VARIABLE_NUMBER: usize = 999;
84 const TX_HASH_COLUMNS: usize = 3;
86 const LOG_COLUMNS: usize = 11;
88
89 pub fn new(pool: SqlitePool, max_variable_number: usize) -> Self {
90 assert!(
91 max_variable_number >= Self::LOG_COLUMNS,
92 "SQLite max_variable_number ({max_variable_number}) must be >= {}",
93 Self::LOG_COLUMNS
94 );
95 Self {
96 pool,
97 max_variable_number,
98 tx_insert_chunk_size: max_variable_number / Self::TX_HASH_COLUMNS,
99 log_insert_chunk_size: max_variable_number / Self::LOG_COLUMNS,
100 }
101 }
102}
103
104#[derive(Clone)]
106pub struct ReceiptProvider<B: BlockInfoProvider = SubxtBlockInfoProvider> {
107 db_ctx: DbContext,
109 block_provider: B,
111 receipt_extractor: ReceiptExtractor,
113 keep_latest_n_blocks: Option<usize>,
115 block_number_to_hashes: Arc<Mutex<BTreeMap<SubstrateBlockNumber, BlockHashMap>>>,
117}
118
119#[derive(Clone, Debug, PartialEq, Eq)]
121struct BlockHashMap {
122 substrate_hash: H256,
123 ethereum_hash: H256,
124}
125
126impl BlockHashMap {
127 fn new(substrate_hash: H256, ethereum_hash: H256) -> Self {
128 Self { substrate_hash, ethereum_hash }
129 }
130}
131
132pub trait BlockInfo {
136 fn hash(&self) -> H256;
138 fn number(&self) -> SubstrateBlockNumber;
140}
141
142impl BlockInfo for SubstrateBlock {
143 fn hash(&self) -> H256 {
144 self.block_hash()
145 }
146 fn number(&self) -> SubstrateBlockNumber {
147 self.block_number()
148 }
149}
150
151pub const MAX_CACHED_BLOCKS: usize = 256;
153
154macro_rules! upsert_sync_label {
157 ($pool:expr, $op:literal, $label:expr, $checkpoint:expr) => {{
158 let label_str = $label.to_string();
159 let block_number = $checkpoint.block_number as i64;
160 let block_hash = $checkpoint.block_hash.map(|h| h.as_bytes().to_vec());
161 query!(
162 "INSERT INTO sync_state (label, block_number, block_hash)
163 VALUES ($1, $2, $3)
164 ON CONFLICT(label) DO UPDATE
165 SET block_number = excluded.block_number, block_hash = excluded.block_hash
166 WHERE sync_state.block_number " +
167 $op + " excluded.block_number
168 ",
169 label_str,
170 block_number,
171 block_hash
172 )
173 .execute($pool)
174 .await?;
175 }};
176}
177
178async fn insert_block_mapping<'e, E: sqlx::Executor<'e, Database = Sqlite>>(
179 executor: E,
180 block_map: &BlockHashMap,
181) -> Result<sqlx::sqlite::SqliteQueryResult, sqlx::Error> {
182 let ethereum_hash_ref = block_map.ethereum_hash.as_ref();
183 let substrate_hash_ref = block_map.substrate_hash.as_ref();
184 query!(
185 r#"
186 INSERT OR REPLACE INTO eth_to_substrate_blocks (ethereum_block_hash, substrate_block_hash)
187 VALUES ($1, $2)
188 "#,
189 ethereum_hash_ref,
190 substrate_hash_ref,
191 )
192 .execute(executor)
193 .await
194}
195
196impl<B: BlockInfoProvider> ReceiptProvider<B> {
197 pub async fn new(
199 db_ctx: DbContext,
200 block_provider: B,
201 receipt_extractor: ReceiptExtractor,
202 keep_latest_n_blocks: Option<usize>,
203 ) -> Result<Self, ClientError> {
204 sqlx::migrate!()
205 .run(&db_ctx.pool)
206 .await
207 .map_err(|e| sqlx::Error::Migrate(e.into()))?;
208
209 let provider = Self {
210 db_ctx,
211 block_provider,
212 receipt_extractor,
213 keep_latest_n_blocks,
214 block_number_to_hashes: Default::default(),
215 };
216 provider.restore_first_evm_block().await?;
217
218 Ok(provider)
219 }
220
221 pub fn is_before_earliest_block(&self, at: &BlockNumberOrTag) -> bool {
223 match at {
224 BlockNumberOrTag::Number(block_number) => {
225 self.receipt_extractor.is_before_first_evm_block(*block_number)
226 },
227 BlockNumberOrTag::Latest |
228 BlockNumberOrTag::Finalized |
229 BlockNumberOrTag::Safe |
230 BlockNumberOrTag::Earliest |
231 BlockNumberOrTag::Pending => false,
232 }
233 }
234
235 pub fn first_evm_block(&self) -> Option<SubstrateBlockNumber> {
237 self.receipt_extractor.first_evm_block()
238 }
239
240 pub async fn set_first_evm_block(
242 &self,
243 block_number: SubstrateBlockNumber,
244 ) -> Result<(), ClientError> {
245 self.receipt_extractor.set_first_evm_block(block_number);
246 self.set_sync_label(ChainMetadata::FirstEvmBlock, SyncCheckpoint::from_number(block_number))
247 .await
248 }
249
250 async fn restore_first_evm_block(&self) -> Result<(), ClientError> {
252 let Some(evm_first) =
253 self.get_sync_label(ChainMetadata::FirstEvmBlock).await?.map(|c| c.block_number)
254 else {
255 return Ok(());
256 };
257
258 let has_evm_hash = |block_number: SubstrateBlockNumber| async move {
259 match self.block_provider.block_by_number(block_number).await.ok().flatten() {
260 Some(block) => self
261 .receipt_extractor
262 .get_ethereum_block_hash(&block.hash(), block_number)
263 .await
264 .is_some(),
265 None => false,
266 }
267 };
268
269 let current_has_evm = has_evm_hash(evm_first).await;
271 let predecessor_has_evm =
272 if evm_first > 0 { has_evm_hash(evm_first - 1).await } else { false };
273
274 if !current_has_evm || predecessor_has_evm {
275 log::warn!(target: LOG_TARGET,
276 "๐๏ธ Stored first-evm-block=#{evm_first} is stale \
277 (has_evm={current_has_evm}, predecessor_has_evm={predecessor_has_evm}), \
278 clearing.");
279 if let Err(e) = self.delete_sync_label(ChainMetadata::FirstEvmBlock).await {
280 log::error!(target: LOG_TARGET,
281 "๐๏ธ Failed to clear stale first-evm-block from DB: {e:?}");
282 }
283 } else {
284 self.receipt_extractor.set_first_evm_block(evm_first);
285 }
286 Ok(())
287 }
288
289 pub async fn find_transaction(&self, transaction_hash: &H256) -> Option<(H256, usize)> {
291 let transaction_hash_bytes = transaction_hash.as_ref();
292 let result = query!(
293 r#"
294 SELECT block_hash, transaction_index
295 FROM transaction_hashes
296 WHERE transaction_hash = $1
297 "#,
298 transaction_hash_bytes
299 )
300 .fetch_optional(&self.db_ctx.pool)
301 .await
302 .inspect_err(|err| {
303 log::trace!(target: LOG_TARGET,
304 "find_transaction: DB query failed for tx {transaction_hash:?}: {err:?}");
305 })
306 .ok()?
307 .or_else(|| {
308 log::trace!(target: LOG_TARGET,
309 "find_transaction: tx {transaction_hash:?} not found in DB");
310 None
311 })?;
312
313 let block_hash = H256::from_slice(&result.block_hash[..]);
314 let transaction_index = result.transaction_index.try_into().ok()?;
315 Some((block_hash, transaction_index))
316 }
317
318 pub async fn get_substrate_hash(&self, ethereum_block_hash: &H256) -> Option<H256> {
320 let ethereum_hash = ethereum_block_hash.as_ref();
321 let result = query!(
322 r#"
323 SELECT substrate_block_hash
324 FROM eth_to_substrate_blocks
325 WHERE ethereum_block_hash = $1
326 "#,
327 ethereum_hash
328 )
329 .fetch_optional(&self.db_ctx.pool)
330 .await
331 .inspect_err(|e| {
332 log::error!(target: LOG_TARGET, "failed to get block mapping for ethereum block {ethereum_block_hash:?}, err: {e:?}");
333 })
334 .ok()?
335 .or_else(||{
336 log::trace!(target: LOG_TARGET, "No block mapping found for ethereum block: {ethereum_block_hash:?}");
337 None
338 })?;
339
340 log::trace!(target: LOG_TARGET, "Get block mapping ethereum block: {:?} -> substrate block: {ethereum_block_hash:?}", H256::from_slice(&result.substrate_block_hash[..]));
341
342 Some(H256::from_slice(&result.substrate_block_hash[..]))
343 }
344
345 pub async fn get_ethereum_hash(&self, substrate_block_hash: &H256) -> Option<H256> {
347 let substrate_hash = substrate_block_hash.as_ref();
348 let result = query!(
349 r#"
350 SELECT ethereum_block_hash
351 FROM eth_to_substrate_blocks
352 WHERE substrate_block_hash = $1
353 "#,
354 substrate_hash
355 )
356 .fetch_optional(&self.db_ctx.pool)
357 .await
358 .inspect_err(|e| {
359 log::error!(target: LOG_TARGET, "failed to get block mapping for substrate block {substrate_block_hash:?}, err: {e:?}");
360 })
361 .ok()?
362 .or_else(||{
363 log::trace!(target: LOG_TARGET, "No block mapping found for substrate block: {substrate_block_hash:?}");
364 None
365 })?;
366
367 log::trace!(target: LOG_TARGET, "Get block mapping substrate block: {substrate_block_hash:?} -> ethereum block: {:?}", H256::from_slice(&result.ethereum_block_hash[..]));
368
369 Some(H256::from_slice(&result.ethereum_block_hash[..]))
370 }
371
372 async fn remove(&self, block_mappings: &[BlockHashMap]) -> Result<(), ClientError> {
374 if block_mappings.is_empty() {
375 return Ok(());
376 }
377 log::debug!(target: LOG_TARGET, "Removing block hashes: {block_mappings:?}");
378
379 let mut db_tx = self.db_ctx.pool.begin().await?;
380
381 for chunk in block_mappings.chunks(self.db_ctx.max_variable_number) {
382 let placeholders = vec!["?"; chunk.len()].join(", ");
383 let sql_tx =
384 format!("DELETE FROM transaction_hashes WHERE block_hash in ({placeholders})");
385 let sql_logs = format!("DELETE FROM logs WHERE block_hash in ({placeholders})");
386 let sql_mappings = format!(
387 "DELETE FROM eth_to_substrate_blocks WHERE substrate_block_hash in ({placeholders})"
388 );
389
390 let mut delete_tx_query = sqlx::query(&sql_tx);
391 let mut delete_logs_query = sqlx::query(&sql_logs);
392 let mut delete_mappings_query = sqlx::query(&sql_mappings);
393
394 for block_map in chunk {
395 delete_tx_query = delete_tx_query.bind(block_map.substrate_hash.as_ref());
396 delete_logs_query = delete_logs_query.bind(block_map.ethereum_hash.as_ref());
397 delete_mappings_query =
398 delete_mappings_query.bind(block_map.substrate_hash.as_ref());
399 }
400
401 delete_tx_query.execute(&mut *db_tx).await?;
402 delete_logs_query.execute(&mut *db_tx).await?;
403 delete_mappings_query.execute(&mut *db_tx).await?;
404 }
405
406 db_tx.commit().await?;
407 Ok(())
408 }
409
410 pub async fn get_sync_label(
412 &self,
413 label: impl SyncStateKey,
414 ) -> Result<Option<SyncCheckpoint>, ClientError> {
415 let label_str = label.to_string();
416 let row = query!(
417 r#"
418 SELECT block_number, block_hash
419 FROM sync_state
420 WHERE label = $1
421 "#,
422 label_str
423 )
424 .fetch_optional(&self.db_ctx.pool)
425 .await?;
426
427 match row {
428 Some(row) => {
429 let block_number: SubstrateBlockNumber =
430 row.block_number.try_into().map_err(|_| {
431 sqlx::Error::Decode(
432 format!("block_number {} overflows u32", row.block_number).into(),
433 )
434 })?;
435 Ok(Some(SyncCheckpoint {
436 block_number,
437 block_hash: row
438 .block_hash
439 .filter(|b| b.len() == 32)
440 .map(|b| H256::from_slice(&b)),
441 }))
442 },
443 None => Ok(None),
444 }
445 }
446
447 pub async fn set_sync_label(
449 &self,
450 label: impl SyncStateKey,
451 checkpoint: SyncCheckpoint,
452 ) -> Result<(), ClientError> {
453 let label_str = label.to_string();
454 let block_number = checkpoint.block_number as i64;
455 let block_hash = checkpoint.block_hash.map(|h| h.as_bytes().to_vec());
456 query!(
457 r#"
458 INSERT OR REPLACE INTO sync_state (label, block_number, block_hash)
459 VALUES ($1, $2, $3)
460 "#,
461 label_str,
462 block_number,
463 block_hash,
464 )
465 .execute(&self.db_ctx.pool)
466 .await?;
467 Ok(())
468 }
469
470 pub async fn delete_sync_label(&self, label: impl SyncStateKey) -> Result<(), ClientError> {
472 let label_str = label.to_string();
473 query!(
474 r#"
475 DELETE FROM sync_state WHERE label = $1
476 "#,
477 label_str,
478 )
479 .execute(&self.db_ctx.pool)
480 .await?;
481 Ok(())
482 }
483
484 pub async fn advance_sync_label(
488 &self,
489 label: SyncLabel,
490 checkpoint: SyncCheckpoint,
491 ) -> Result<(), ClientError> {
492 upsert_sync_label!(&self.db_ctx.pool, "<", label, checkpoint);
493 Ok(())
494 }
495
496 pub async fn recede_sync_label(
500 &self,
501 label: SyncLabel,
502 checkpoint: SyncCheckpoint,
503 ) -> Result<(), ClientError> {
504 upsert_sync_label!(&self.db_ctx.pool, ">", label, checkpoint);
505 Ok(())
506 }
507
508 pub async fn get_processed_eth_block_hash(
510 &self,
511 block_number: SubstrateBlockNumber,
512 substrate_hash: H256,
513 ) -> Option<H256> {
514 self.block_number_to_hashes
515 .lock()
516 .await
517 .get(&block_number)
518 .filter(|entry| entry.substrate_hash == substrate_hash)
519 .map(|entry| entry.ethereum_hash)
520 }
521
522 pub async fn receipts_from_block(
524 &self,
525 block: &SubstrateBlock,
526 ethereum_hash: H256,
527 ) -> Result<Vec<(TransactionSigned, ReceiptInfo)>, ClientError> {
528 self.receipt_extractor
529 .extract_from_block_with_eth_hash(block, ethereum_hash)
530 .await
531 }
532
533 pub async fn block_receipts(
535 &self,
536 block: &SubstrateBlock,
537 ) -> Result<Vec<(TransactionSigned, ReceiptInfo)>, ClientError> {
538 self.receipt_extractor.extract_from_block(block).await
539 }
540
541 pub async fn insert_block_receipts_past(
544 &self,
545 block: &SubstrateBlock,
546 ethereum_hash: &H256,
547 ) -> Result<(), ClientError> {
548 let receipts = self
549 .receipt_extractor
550 .extract_from_block_with_eth_hash(block, *ethereum_hash)
551 .await?;
552 self.insert_into_db(block, &receipts, ethereum_hash).await?;
553 Ok(())
554 }
555
556 pub async fn insert_block_receipts(
558 &self,
559 block: &SubstrateBlock,
560 receipts: &[(TransactionSigned, ReceiptInfo)],
561 ethereum_hash: &H256,
562 ) -> Result<(), ClientError> {
563 self.insert(block, receipts, ethereum_hash).await
564 }
565
566 async fn insert(
571 &self,
572 block: &impl BlockInfo,
573 receipts: &[(TransactionSigned, ReceiptInfo)],
574 ethereum_hash: &H256,
575 ) -> Result<(), ClientError> {
576 let block_map = BlockHashMap::new(block.hash(), *ethereum_hash);
577 self.prune_blocks(block.number(), &block_map).await?;
578 self.insert_into_db(block, receipts, ethereum_hash).await?;
579 Ok(())
580 }
581
582 async fn prune_blocks(
584 &self,
585 block_number: SubstrateBlockNumber,
586 block_map: &BlockHashMap,
587 ) -> Result<(), ClientError> {
588 let mut to_remove = Vec::new();
589 let mut block_number_to_hash = self.block_number_to_hashes.lock().await;
590
591 match block_number_to_hash.insert(block_number, block_map.clone()) {
593 Some(old_block_map) if &old_block_map != block_map => {
594 to_remove.push(old_block_map);
595
596 let mut next_block_number = block_number.saturating_add(1);
599 while let Some(old_block_map) = block_number_to_hash.remove(&next_block_number) {
600 to_remove.push(old_block_map);
601 next_block_number = next_block_number.saturating_add(1);
602 }
603 },
604 _ => {},
605 }
606
607 if let Some(keep_latest_n_blocks) = self.keep_latest_n_blocks {
608 while block_number_to_hash.len() > keep_latest_n_blocks {
611 if let Some((_, block_map)) = block_number_to_hash.pop_first() {
613 to_remove.push(block_map);
614 }
615 }
616 } else {
617 while block_number_to_hash.len() > MAX_CACHED_BLOCKS {
620 block_number_to_hash.pop_first();
621 }
622 }
623
624 drop(block_number_to_hash);
626
627 if !to_remove.is_empty() {
628 log::trace!(target: LOG_TARGET, "Pruning old blocks: {to_remove:?}");
629 self.remove(&to_remove).await?;
630 }
631
632 Ok(())
633 }
634
635 async fn insert_into_db(
637 &self,
638 block: &impl BlockInfo,
639 receipts: &[(TransactionSigned, ReceiptInfo)],
640 ethereum_hash: &H256,
641 ) -> Result<(), ClientError> {
642 let block_number = block.number() as i64;
643 let substrate_block_hash = block.hash();
644 let substrate_hash_ref = substrate_block_hash.as_ref();
645 let ethereum_hash_ref = ethereum_hash.as_ref();
646
647 log::trace!(target: LOG_TARGET, "Inserting receipts for block #{block_number} ethereum: {ethereum_hash:?} substrate: {substrate_block_hash:?}");
648
649 let result = sqlx::query!(
652 r#"SELECT EXISTS(SELECT 1 FROM eth_to_substrate_blocks WHERE substrate_block_hash = $1) AS "exists!:bool""#, substrate_hash_ref
653 )
654 .fetch_one(&self.db_ctx.pool)
655 .await?;
656
657 if result.exists {
660 log::trace!(target: LOG_TARGET,
661 "Skipping receipt insert for block #{block_number} ({substrate_block_hash:?}): \
662 mapping already exists. ETH hash: {ethereum_hash:?}, receipts count: {count}",
663 count = receipts.len(),
664 );
665 return Ok(());
666 }
667
668 let mut db_tx = self.db_ctx.pool.begin().await?;
669
670 for chunk in receipts.chunks(self.db_ctx.tx_insert_chunk_size) {
671 let mut query_builder = QueryBuilder::<Sqlite>::new(
672 "INSERT OR REPLACE INTO transaction_hashes (transaction_hash, block_hash, transaction_index) ",
673 );
674 query_builder.push_values(chunk, |mut row, (_, receipt)| {
675 row.push_bind(receipt.transaction_hash.as_ref() as &[u8])
676 .push_bind(substrate_hash_ref)
677 .push_bind(receipt.transaction_index.as_u32() as i32);
678 });
679 query_builder.build().execute(&mut *db_tx).await?;
680 }
681
682 let all_logs: Vec<(i32, &[u8], &Log)> = receipts
683 .iter()
684 .flat_map(|(_, receipt)| {
685 let tx_index = receipt.transaction_index.as_u32() as i32;
686 let tx_hash: &[u8] = receipt.transaction_hash.as_ref();
687 receipt.logs.iter().map(move |log| (tx_index, tx_hash, log))
688 })
689 .collect();
690
691 for chunk in all_logs.chunks(self.db_ctx.log_insert_chunk_size) {
692 let mut query_builder = QueryBuilder::<Sqlite>::new(
693 "INSERT OR REPLACE INTO logs(block_hash, transaction_index, log_index, address, block_number, transaction_hash, topic_0, topic_1, topic_2, topic_3, data) ",
694 );
695 query_builder.push_values(chunk, |mut row, (tx_index, tx_hash, log)| {
696 row.push_bind(ethereum_hash_ref)
697 .push_bind(*tx_index)
698 .push_bind(log.log_index.as_u32() as i32)
699 .push_bind(log.address.as_ref() as &[u8])
700 .push_bind(block_number)
701 .push_bind(*tx_hash)
702 .push_bind(log.topics.first().map(|v| &v[..]))
703 .push_bind(log.topics.get(1).map(|v| &v[..]))
704 .push_bind(log.topics.get(2).map(|v| &v[..]))
705 .push_bind(log.topics.get(3).map(|v| &v[..]))
706 .push_bind(log.data.as_ref().map(|v| &v.0[..]));
707 });
708 query_builder.build().execute(&mut *db_tx).await?;
709 }
710
711 let block_map = BlockHashMap::new(substrate_block_hash, *ethereum_hash);
712 insert_block_mapping(&mut *db_tx, &block_map).await?;
713
714 db_tx.commit().await?;
715 log::trace!(target: LOG_TARGET, "Inserted {} receipts for block #{block_number} ethereum: {ethereum_hash:?} substrate: {substrate_block_hash:?}", receipts.len());
716
717 Ok(())
718 }
719
720 pub async fn logs<Fut>(
724 &self,
725 filter: Option<Filter>,
726 resolve_block_number: impl Fn(BlockNumberOrTag) -> Fut,
727 ) -> anyhow::Result<Vec<Log>>
728 where
729 Fut: Future<Output = anyhow::Result<U256>>,
730 {
731 let mut qb = QueryBuilder::<Sqlite>::new("SELECT logs.* FROM logs WHERE 1=1");
732 let filter = filter.unwrap_or_default();
733
734 match filter.block_option {
735 FilterBlockOption::AtBlockHash(hash) => {
736 qb.push(" AND block_hash = ").push_bind(hash.as_slice().to_vec());
737 },
738 FilterBlockOption::Range { from_block, to_block } => {
739 let from_block =
740 OptionFuture::from(from_block.map(&resolve_block_number)).await.transpose()?;
741 let to_block =
742 OptionFuture::from(to_block.map(&resolve_block_number)).await.transpose()?;
743
744 let latest_block = U256::from(self.block_provider.latest_block_number().await);
746
747 match (from_block, to_block) {
748 (Some(block), _) | (_, Some(block)) if block > latest_block => {
749 anyhow::bail!("block number exceeds latest block");
750 },
751 (Some(from_block), Some(to_block)) if from_block > to_block => {
752 anyhow::bail!("invalid block range params");
753 },
754 (Some(from_block), Some(to_block)) if from_block == to_block => {
755 qb.push(" AND block_number = ").push_bind(from_block.as_u64() as i64);
756 },
757 (Some(from_block), Some(to_block)) => {
758 qb.push(" AND block_number BETWEEN ")
759 .push_bind(from_block.as_u64() as i64)
760 .push(" AND ")
761 .push_bind(to_block.as_u64() as i64);
762 },
763 (Some(from_block), None) => {
764 qb.push(" AND block_number >= ").push_bind(from_block.as_u64() as i64);
765 },
766 (None, Some(to_block)) => {
767 qb.push(" AND block_number <= ").push_bind(to_block.as_u64() as i64);
768 },
769 (None, None) => {
770 qb.push(" AND block_number = ").push_bind(latest_block.as_u64() as i64);
771 },
772 }
773 },
774 }
775
776 if !filter.address.is_empty() {
777 qb.push(" AND address IN (");
778 let mut separated = qb.separated(", ");
779 for addr in filter.address {
780 separated.push_bind(addr.as_slice().to_vec());
781 }
782 separated.push_unseparated(")");
783 }
784
785 for (i, topic) in filter.topics.into_iter().enumerate() {
786 if topic.is_empty() {
787 continue;
788 }
789
790 qb.push(format_args!(" AND topic_{i} IN ("));
791 let mut separated = qb.separated(", ");
792 for hash in topic {
793 separated.push_bind(hash.as_slice().to_vec());
794 }
795 separated.push_unseparated(")");
796 }
797
798 qb.push(" LIMIT ").push_bind(MAX_LOG_RESULTS as i64);
799
800 let logs = qb.build().try_map(parse_log_row).fetch_all(&self.db_ctx.pool).await?;
801
802 if logs.len() == MAX_LOG_RESULTS {
803 log::warn!(
804 target: LOG_TARGET,
805 "Log query hit limit of {MAX_LOG_RESULTS}; results may be truncated",
806 );
807 }
808
809 Ok(logs)
810 }
811
812 pub async fn logs_by_block_number(
814 &self,
815 block_number: SubstrateBlockNumber,
816 ethereum_hash: H256,
817 ) -> Result<Vec<Log>, ClientError> {
818 let mut query_builder =
819 QueryBuilder::<Sqlite>::new("SELECT logs.* FROM logs WHERE block_number = ");
820 query_builder
821 .push_bind(block_number as i64)
822 .push(" AND block_hash = ")
823 .push_bind(ethereum_hash.as_bytes().to_vec())
824 .push(" ORDER BY log_index LIMIT ")
825 .push_bind(MAX_LOG_RESULTS as i64);
826
827 let logs = query_builder
828 .build()
829 .try_map(parse_log_row)
830 .fetch_all(&self.db_ctx.pool)
831 .await?;
832
833 if logs.len() == MAX_LOG_RESULTS {
834 log::warn!(
835 target: LOG_TARGET,
836 "Log query for block {block_number} hit limit of {MAX_LOG_RESULTS}; results may be truncated",
837 );
838 }
839
840 Ok(logs)
841 }
842
843 pub async fn receipts_count_per_block(&self, block_hash: &H256) -> Option<usize> {
845 let block_hash = block_hash.as_ref();
846 let row = query!(
847 r#"
848 SELECT COUNT(*) as count
849 FROM transaction_hashes
850 WHERE block_hash = $1
851 "#,
852 block_hash
853 )
854 .fetch_one(&self.db_ctx.pool)
855 .await
856 .ok()?;
857
858 let count = row.count as usize;
859 Some(count)
860 }
861
862 pub async fn block_transaction_hashes(
864 &self,
865 block_hash: &H256,
866 ) -> Option<HashMap<usize, H256>> {
867 let block_hash = block_hash.as_ref();
868 let rows = query!(
869 r#"
870 SELECT transaction_index, transaction_hash
871 FROM transaction_hashes
872 WHERE block_hash = $1
873 "#,
874 block_hash
875 )
876 .map(|row| {
877 let transaction_index = row.transaction_index as usize;
878 let transaction_hash = H256::from_slice(&row.transaction_hash);
879 (transaction_index, transaction_hash)
880 })
881 .fetch_all(&self.db_ctx.pool)
882 .await
883 .ok()?;
884
885 Some(rows.into_iter().collect())
886 }
887
888 pub async fn receipt_by_block_hash_and_index(
890 &self,
891 block_hash: &H256,
892 transaction_index: usize,
893 ) -> Option<ReceiptInfo> {
894 let block = self.block_provider.block_by_hash(block_hash).await.ok()??;
895 let (_, receipt) = self
896 .receipt_extractor
897 .extract_from_transaction(&block, transaction_index)
898 .await
899 .ok()?;
900 Some(receipt)
901 }
902
903 pub async fn receipt_by_hash(&self, transaction_hash: &H256) -> Option<ReceiptInfo> {
905 let (block_hash, transaction_index) = self.find_transaction(transaction_hash).await?;
906
907 let block = match self.block_provider.block_by_hash(&block_hash).await {
908 Ok(Some(b)) => b,
909 Ok(None) => {
910 log::trace!(target: LOG_TARGET,
911 "receipt_by_hash: block {block_hash:?} not available from node (pruned?) for tx {transaction_hash:?}");
912 return None;
913 },
914 Err(err) => {
915 log::trace!(target: LOG_TARGET,
916 "receipt_by_hash: failed to fetch block {block_hash:?} for tx {transaction_hash:?}: {err:?}");
917 return None;
918 },
919 };
920
921 match self.receipt_extractor.extract_from_transaction(&block, transaction_index).await {
922 Ok((_, receipt)) => Some(receipt),
923 Err(err) => {
924 log::trace!(target: LOG_TARGET,
925 "receipt_by_hash: extraction failed for tx {transaction_hash:?} in block {block_hash:?}: {err:?}");
926 None
927 },
928 }
929 }
930
931 pub async fn signed_tx_by_hash(&self, transaction_hash: &H256) -> Option<TransactionSigned> {
933 let (block_hash, transaction_index) = self.find_transaction(transaction_hash).await?;
934
935 let block = self.block_provider.block_by_hash(&block_hash).await.ok()??;
936 let (signed_tx, _) = self
937 .receipt_extractor
938 .extract_from_transaction(&block, transaction_index)
939 .await
940 .inspect_err(|err| {
941 log::trace!(target: LOG_TARGET,
942 "signed_tx_by_hash: extraction failed for tx {transaction_hash:?} \
943 in block {block_hash:?}: {err:?}");
944 })
945 .ok()?;
946 Some(signed_tx)
947 }
948}
949
950#[cfg(test)]
951mod tests {
952 use super::*;
953 use crate::{
954 ReceiptInfo,
955 test::{MockBlockInfo, MockBlockInfoProvider},
956 };
957 use alloy_primitives::{Address as AlloyAddress, B256};
958 use pallet_revive::evm::TransactionSigned;
959 use pretty_assertions::assert_eq;
960 use sp_core::{H160, H256};
961 use sqlx::SqlitePool;
962
963 async fn count(pool: &SqlitePool, table: &str, block_hash: Option<H256>) -> usize {
964 let count: i64 = match block_hash {
965 None => {
966 sqlx::query_scalar(&format!("SELECT COUNT(*) FROM {table}"))
967 .fetch_one(pool)
968 .await
969 },
970 Some(hash) => {
971 sqlx::query_scalar(&format!("SELECT COUNT(*) FROM {table} WHERE block_hash = ?"))
972 .bind(hash.as_ref())
973 .fetch_one(pool)
974 .await
975 },
976 }
977 .unwrap();
978
979 count as _
980 }
981
982 fn mock_provider() -> ReceiptProvider<MockBlockInfoProvider> {
983 ReceiptProvider {
984 db_ctx: DbContext::new(
985 SqlitePool::connect_lazy("sqlite::memory:").unwrap(),
986 DbContext::DEFAULT_MAX_VARIABLE_NUMBER,
987 ),
988 block_provider: MockBlockInfoProvider {},
989 receipt_extractor: ReceiptExtractor::new_mock(),
990 keep_latest_n_blocks: None,
991 block_number_to_hashes: Default::default(),
992 }
993 }
994
995 fn mock_resolve_block_number_with_latest(
997 latest: u64,
998 ) -> impl Fn(BlockNumberOrTag) -> std::future::Ready<anyhow::Result<U256>> {
999 move |block: BlockNumberOrTag| {
1000 std::future::ready(Ok(match block {
1001 BlockNumberOrTag::Number(v) => U256::from(v),
1002 BlockNumberOrTag::Earliest => U256::zero(),
1003 _ => U256::from(latest),
1005 }))
1006 }
1007 }
1008
1009 impl ReceiptProvider<MockBlockInfoProvider> {
1010 fn with_db_ctx(mut self, db_ctx: DbContext) -> Self {
1011 self.db_ctx = db_ctx;
1012 self
1013 }
1014
1015 fn with_extractor(mut self, extractor: ReceiptExtractor) -> Self {
1016 self.receipt_extractor = extractor;
1017 self
1018 }
1019
1020 fn with_keep_latest(mut self, n: Option<usize>) -> Self {
1021 self.keep_latest_n_blocks = n;
1022 self
1023 }
1024 }
1025
1026 async fn setup_sqlite_provider(pool: SqlitePool) -> ReceiptProvider<MockBlockInfoProvider> {
1027 mock_provider()
1028 .with_db_ctx(DbContext::new(pool, DbContext::DEFAULT_MAX_VARIABLE_NUMBER))
1029 .with_keep_latest(Some(10))
1030 }
1031
1032 #[sqlx::test]
1033 async fn test_insert_remove(pool: SqlitePool) -> anyhow::Result<()> {
1034 let provider = setup_sqlite_provider(pool).await;
1035 let block = MockBlockInfo { hash: H256::default(), number: 0 };
1036 let receipts = vec![(
1037 TransactionSigned::default(),
1038 ReceiptInfo {
1039 logs: vec![Log { block_hash: block.hash, ..Default::default() }],
1040 ..Default::default()
1041 },
1042 )];
1043 let ethereum_hash = H256::from([1_u8; 32]);
1044 let block_map = BlockHashMap::new(block.hash(), ethereum_hash);
1045
1046 provider.insert(&block, &receipts, ðereum_hash).await?;
1047 let row = provider.find_transaction(&receipts[0].1.transaction_hash).await;
1048 assert_eq!(row, Some((block.hash, 0)));
1049 assert_eq!(count(&provider.db_ctx.pool, "transaction_hashes", Some(block.hash())).await, 1);
1050 assert_eq!(count(&provider.db_ctx.pool, "logs", Some(ethereum_hash)).await, 1);
1051
1052 provider.remove(&[block_map]).await?;
1053 assert_eq!(count(&provider.db_ctx.pool, "transaction_hashes", Some(block.hash())).await, 0);
1054 assert_eq!(count(&provider.db_ctx.pool, "logs", Some(ethereum_hash)).await, 0);
1055 Ok(())
1056 }
1057
1058 #[sqlx::test]
1059 async fn test_prune(pool: SqlitePool) -> anyhow::Result<()> {
1060 let provider = setup_sqlite_provider(pool).await;
1061 let n = provider.keep_latest_n_blocks.unwrap();
1062
1063 for i in 0..2 * n {
1064 let block = MockBlockInfo { hash: H256::from([i as u8; 32]), number: i as _ };
1065 let transaction_hash = H256::from([i as u8; 32]);
1066 let receipts = vec![(
1067 TransactionSigned::default(),
1068 ReceiptInfo {
1069 transaction_hash,
1070 logs: vec![Log {
1071 block_hash: block.hash,
1072 transaction_hash,
1073 ..Default::default()
1074 }],
1075 ..Default::default()
1076 },
1077 )];
1078 let ethereum_hash = H256::from([(i + 1) as u8; 32]);
1079 provider.insert(&block, &receipts, ðereum_hash).await?;
1080 }
1081 assert_eq!(count(&provider.db_ctx.pool, "transaction_hashes", None).await, n);
1082 assert_eq!(count(&provider.db_ctx.pool, "logs", None).await, n);
1083 assert_eq!(count(&provider.db_ctx.pool, "eth_to_substrate_blocks", None).await, n);
1084 assert_eq!(provider.block_number_to_hashes.lock().await.len(), n);
1085
1086 return Ok(());
1087 }
1088
1089 #[sqlx::test]
1090 async fn test_fork(pool: SqlitePool) -> anyhow::Result<()> {
1091 let provider = setup_sqlite_provider(pool).await;
1092
1093 let build_block = |seed, number| {
1094 let block = MockBlockInfo { hash: H256::from([seed; 32]), number };
1095 let transaction_hash = H256::from([seed; 32]);
1096 let receipts = vec![(
1097 TransactionSigned::default(),
1098 ReceiptInfo {
1099 transaction_hash,
1100 logs: vec![Log {
1101 block_hash: block.hash,
1102 transaction_hash,
1103 ..Default::default()
1104 }],
1105 ..Default::default()
1106 },
1107 )];
1108 let ethereum_hash = H256::from([seed + 1; 32]);
1109
1110 (block, receipts, ethereum_hash)
1111 };
1112
1113 let (block0, receipts, ethereum_hash_0) = build_block(0, 0);
1115 provider.insert(&block0, &receipts, ðereum_hash_0).await?;
1116 let (block1, receipts, ethereum_hash_1) = build_block(1, 1);
1117 provider.insert(&block1, &receipts, ðereum_hash_1).await?;
1118 let (block2, receipts, ethereum_hash_2) = build_block(2, 2);
1119 provider.insert(&block2, &receipts, ðereum_hash_2).await?;
1120 let (block3, receipts, ethereum_hash_3) = build_block(3, 3);
1121 provider.insert(&block3, &receipts, ðereum_hash_3).await?;
1122
1123 assert_eq!(count(&provider.db_ctx.pool, "transaction_hashes", None).await, 4);
1124 assert_eq!(count(&provider.db_ctx.pool, "logs", None).await, 4);
1125 assert_eq!(count(&provider.db_ctx.pool, "eth_to_substrate_blocks", None).await, 4);
1126 assert_eq!(
1127 provider.block_number_to_hashes.lock().await.clone(),
1128 [
1129 (0, BlockHashMap::new(block0.hash, ethereum_hash_0)),
1130 (1, BlockHashMap::new(block1.hash, ethereum_hash_1)),
1131 (2, BlockHashMap::new(block2.hash, ethereum_hash_2)),
1132 (3, BlockHashMap::new(block3.hash, ethereum_hash_3))
1133 ]
1134 .into(),
1135 );
1136
1137 let (fork_block, receipts, ethereum_hash_fork) = build_block(4, 1);
1139 provider.insert(&fork_block, &receipts, ðereum_hash_fork).await?;
1140
1141 assert_eq!(count(&provider.db_ctx.pool, "transaction_hashes", None).await, 2);
1142 assert_eq!(count(&provider.db_ctx.pool, "logs", None).await, 2);
1143 assert_eq!(count(&provider.db_ctx.pool, "eth_to_substrate_blocks", None).await, 2);
1144
1145 assert_eq!(
1146 provider.block_number_to_hashes.lock().await.clone(),
1147 [
1148 (0, BlockHashMap::new(block0.hash, ethereum_hash_0)),
1149 (1, BlockHashMap::new(fork_block.hash, ethereum_hash_fork))
1150 ]
1151 .into(),
1152 );
1153
1154 return Ok(());
1155 }
1156
1157 #[sqlx::test]
1158 async fn test_reorg_same_transaction_hash(pool: SqlitePool) -> anyhow::Result<()> {
1159 let provider = setup_sqlite_provider(pool).await;
1160
1161 let tx_hash = H256::from([42u8; 32]);
1163
1164 let block_a = MockBlockInfo { hash: H256::from([1u8; 32]), number: 1 };
1166 let ethereum_hash_a = H256::from([2u8; 32]);
1167 let receipts_a = vec![(
1168 TransactionSigned::default(),
1169 ReceiptInfo {
1170 transaction_hash: tx_hash,
1171 transaction_index: U256::from(0),
1172 ..Default::default()
1173 },
1174 )];
1175
1176 provider.insert(&block_a, &receipts_a, ðereum_hash_a).await?;
1177
1178 let (found_hash, _) = provider.find_transaction(&tx_hash).await.unwrap();
1180 assert_eq!(found_hash, block_a.hash);
1181
1182 provider.block_number_to_hashes.lock().await.clear();
1184
1185 let block_b = MockBlockInfo { hash: H256::from([3u8; 32]), number: 1 };
1187 let ethereum_hash_b = H256::from([4u8; 32]);
1188 let receipts_b = vec![(
1189 TransactionSigned::default(),
1190 ReceiptInfo {
1191 transaction_hash: tx_hash, transaction_index: U256::from(0),
1193 ..Default::default()
1194 },
1195 )];
1196
1197 provider.insert(&block_b, &receipts_b, ðereum_hash_b).await?;
1199
1200 let (found_hash, _) = provider.find_transaction(&tx_hash).await.unwrap();
1202 assert_eq!(found_hash, block_b.hash);
1203
1204 Ok(())
1205 }
1206
1207 #[sqlx::test]
1208 async fn test_receipts_count_per_block(pool: SqlitePool) -> anyhow::Result<()> {
1209 let provider = setup_sqlite_provider(pool).await;
1210 let block = MockBlockInfo { hash: H256::default(), number: 0 };
1211 let receipts = vec![
1212 (
1213 TransactionSigned::default(),
1214 ReceiptInfo { transaction_hash: H256::from([0u8; 32]), ..Default::default() },
1215 ),
1216 (
1217 TransactionSigned::default(),
1218 ReceiptInfo { transaction_hash: H256::from([1u8; 32]), ..Default::default() },
1219 ),
1220 ];
1221 let ethereum_hash = H256::from([2u8; 32]);
1222
1223 provider.insert(&block, &receipts, ðereum_hash).await?;
1224 let count = provider.receipts_count_per_block(&block.hash).await;
1225 assert_eq!(count, Some(2));
1226 Ok(())
1227 }
1228
1229 #[sqlx::test]
1230 async fn test_query_logs(pool: SqlitePool) -> anyhow::Result<()> {
1231 let provider = setup_sqlite_provider(pool).await;
1232 let block1 = MockBlockInfo { hash: H256::from([1u8; 32]), number: 1 };
1233 let block2 = MockBlockInfo { hash: H256::from([2u8; 32]), number: 2 };
1234 let ethereum_hash1 = H256::from([3u8; 32]);
1235 let ethereum_hash2 = H256::from([4u8; 32]);
1236 let log1 = Log {
1237 block_hash: ethereum_hash1,
1238 block_number: block1.number.into(),
1239 address: H160::from([1u8; 20]),
1240 topics: vec![H256::from([1u8; 32]), H256::from([2u8; 32])],
1241 data: Some(vec![0u8; 32].into()),
1242 transaction_hash: H256::default(),
1243 transaction_index: U256::from(1),
1244 log_index: U256::from(1),
1245 ..Default::default()
1246 };
1247 let log2 = Log {
1248 block_hash: ethereum_hash2,
1249 block_number: block2.number.into(),
1250 address: H160::from([2u8; 20]),
1251 topics: vec![H256::from([2u8; 32]), H256::from([3u8; 32])],
1252 transaction_hash: H256::from([1u8; 32]),
1253 transaction_index: U256::from(2),
1254 log_index: U256::from(1),
1255 ..Default::default()
1256 };
1257
1258 provider
1259 .insert(
1260 &block1,
1261 &vec![(
1262 TransactionSigned::default(),
1263 ReceiptInfo {
1264 logs: vec![log1.clone()],
1265 transaction_hash: log1.transaction_hash,
1266 transaction_index: log1.transaction_index,
1267 ..Default::default()
1268 },
1269 )],
1270 ðereum_hash1,
1271 )
1272 .await?;
1273 provider
1274 .insert(
1275 &block2,
1276 &vec![(
1277 TransactionSigned::default(),
1278 ReceiptInfo {
1279 logs: vec![log2.clone()],
1280 transaction_hash: log2.transaction_hash,
1281 transaction_index: log2.transaction_index,
1282 ..Default::default()
1283 },
1284 )],
1285 ðereum_hash2,
1286 )
1287 .await?;
1288
1289 let resolve_block_number = mock_resolve_block_number_with_latest(block2.number.into());
1290
1291 let logs = provider.logs(None, &resolve_block_number).await?;
1293 assert_eq!(logs, vec![log2.clone()]);
1294
1295 let logs = provider
1297 .logs(Some(Filter::new().from_block(log2.block_number.as_u64())), &resolve_block_number)
1298 .await?;
1299 assert_eq!(logs, vec![log2.clone()]);
1300
1301 let logs = provider
1303 .logs(Some(Filter::new().from_block(BlockNumberOrTag::Latest)), &resolve_block_number)
1304 .await?;
1305 assert_eq!(logs, vec![log2.clone()]);
1306
1307 let logs = provider
1309 .logs(Some(Filter::new().to_block(log1.block_number.as_u64())), &resolve_block_number)
1310 .await?;
1311 assert_eq!(logs, vec![log1.clone()]);
1312
1313 let logs = provider
1315 .logs(
1316 Some(Filter::new().at_block_hash(B256::from(log1.block_hash.0))),
1317 &resolve_block_number,
1318 )
1319 .await?;
1320 assert_eq!(logs, vec![log1.clone()]);
1321
1322 let logs = provider
1324 .logs(
1325 Some(
1326 Filter::new()
1327 .from_block(BlockNumberOrTag::Earliest)
1328 .address(AlloyAddress::from(log1.address.0)),
1329 ),
1330 &resolve_block_number,
1331 )
1332 .await?;
1333 assert_eq!(logs, vec![log1.clone()]);
1334
1335 let logs = provider
1337 .logs(
1338 Some(Filter::new().from_block(BlockNumberOrTag::Earliest).address(vec![
1339 AlloyAddress::from(log1.address.0),
1340 AlloyAddress::from(log2.address.0),
1341 ])),
1342 &resolve_block_number,
1343 )
1344 .await?;
1345 assert_eq!(logs, vec![log1.clone(), log2.clone()]);
1346
1347 let logs = provider
1349 .logs(
1350 Some(
1351 Filter::new()
1352 .from_block(BlockNumberOrTag::Earliest)
1353 .event_signature(B256::from(log1.topics[0].0)),
1354 ),
1355 &resolve_block_number,
1356 )
1357 .await?;
1358 assert_eq!(logs, vec![log1.clone()]);
1359
1360 let logs = provider
1362 .logs(
1363 Some(
1364 Filter::new()
1365 .from_block(BlockNumberOrTag::Earliest)
1366 .event_signature(B256::from(log1.topics[0].0))
1367 .topic1(B256::from(log1.topics[1].0)),
1368 ),
1369 &resolve_block_number,
1370 )
1371 .await?;
1372 assert_eq!(logs, vec![log1.clone()]);
1373
1374 let logs = provider
1376 .logs(
1377 Some(Filter::new().from_block(BlockNumberOrTag::Earliest).event_signature(vec![
1378 B256::from(log1.topics[0].0),
1379 B256::from(log2.topics[0].0),
1380 ])),
1381 &resolve_block_number,
1382 )
1383 .await?;
1384 assert_eq!(logs, vec![log1.clone(), log2.clone()]);
1385
1386 let logs = provider
1388 .logs(
1389 Some(
1390 Filter::new()
1391 .from_block(BlockNumberOrTag::Earliest)
1392 .to_block(BlockNumberOrTag::Latest)
1393 .address(vec![
1394 AlloyAddress::from(log1.address.0),
1395 AlloyAddress::from(log2.address.0),
1396 ])
1397 .event_signature(vec![
1398 B256::from(log1.topics[0].0),
1399 B256::from(log2.topics[0].0),
1400 ]),
1401 ),
1402 &resolve_block_number,
1403 )
1404 .await?;
1405 assert_eq!(logs, vec![log1.clone(), log2.clone()]);
1406
1407 for tag in [
1408 BlockNumberOrTag::Latest,
1409 BlockNumberOrTag::Finalized,
1410 BlockNumberOrTag::Safe,
1411 BlockNumberOrTag::Pending,
1412 ] {
1413 let logs = provider
1414 .logs(Some(Filter::new().from_block(tag)), &resolve_block_number)
1415 .await?;
1416 assert_eq!(logs, vec![log2.clone()], "tag {tag:?} should resolve to the latest block");
1417 }
1418
1419 let logs = provider
1420 .logs(
1421 Some(
1422 Filter::new()
1423 .from_block(log1.block_number.as_u64())
1424 .to_block(log1.block_number.as_u64()),
1425 ),
1426 &resolve_block_number,
1427 )
1428 .await?;
1429 assert_eq!(logs, vec![log1.clone()], "from == to selects the single block");
1430
1431 let result = provider
1432 .logs(
1433 Some(
1434 Filter::new()
1435 .from_block(log2.block_number.as_u64())
1436 .to_block(log1.block_number.as_u64()),
1437 ),
1438 &resolve_block_number,
1439 )
1440 .await;
1441 assert!(result.is_err(), "from > to should be rejected");
1442
1443 let result = provider
1444 .logs(
1445 Some(Filter::new().from_block(log2.block_number.as_u64() + 100)),
1446 &resolve_block_number,
1447 )
1448 .await;
1449 assert!(result.is_err(), "block number beyond latest should be rejected");
1450
1451 let failing = |_: BlockNumberOrTag| {
1452 std::future::ready(Err::<U256, anyhow::Error>(anyhow::anyhow!("resolver failed")))
1453 };
1454 let result = provider
1455 .logs(Some(Filter::new().from_block(BlockNumberOrTag::Latest)), &failing)
1456 .await;
1457 assert!(result.is_err(), "a resolver error should propagate");
1458
1459 Ok(())
1460 }
1461
1462 #[sqlx::test]
1463 async fn test_block_mapping_insert_get(pool: SqlitePool) -> anyhow::Result<()> {
1464 let provider = setup_sqlite_provider(pool).await;
1465 let ethereum_hash = H256::from([1u8; 32]);
1466 let substrate_hash = H256::from([2u8; 32]);
1467 let block_map = BlockHashMap::new(substrate_hash, ethereum_hash);
1468
1469 insert_block_mapping(&provider.db_ctx.pool, &block_map).await?;
1471
1472 let resolved = provider.get_substrate_hash(ðereum_hash).await;
1474 assert_eq!(resolved, Some(substrate_hash));
1475
1476 let resolved = provider.get_ethereum_hash(&substrate_hash).await;
1478 assert_eq!(resolved, Some(ethereum_hash));
1479
1480 Ok(())
1481 }
1482
1483 #[sqlx::test]
1484 async fn test_block_mapping_remove(pool: SqlitePool) -> anyhow::Result<()> {
1485 let provider = setup_sqlite_provider(pool).await;
1486 let ethereum_hash1 = H256::from([1u8; 32]);
1487 let ethereum_hash2 = H256::from([2u8; 32]);
1488 let substrate_hash1 = H256::from([3u8; 32]);
1489 let substrate_hash2 = H256::from([4u8; 32]);
1490 let block_map1 = BlockHashMap::new(substrate_hash1, ethereum_hash1);
1491 let block_map2 = BlockHashMap::new(substrate_hash2, ethereum_hash2);
1492
1493 insert_block_mapping(&provider.db_ctx.pool, &block_map1).await?;
1495 insert_block_mapping(&provider.db_ctx.pool, &block_map2).await?;
1496
1497 assert_eq!(
1499 provider.get_substrate_hash(&block_map1.ethereum_hash).await,
1500 Some(block_map1.substrate_hash)
1501 );
1502 assert_eq!(
1503 provider.get_substrate_hash(&block_map2.ethereum_hash).await,
1504 Some(block_map2.substrate_hash)
1505 );
1506
1507 provider.remove(&[block_map1]).await?;
1509
1510 assert_eq!(provider.get_substrate_hash(ðereum_hash1).await, None);
1512 assert_eq!(provider.get_substrate_hash(ðereum_hash2).await, Some(substrate_hash2));
1513
1514 Ok(())
1515 }
1516
1517 #[sqlx::test]
1518 async fn test_block_mapping_pruning_integration(pool: SqlitePool) -> anyhow::Result<()> {
1519 let provider = setup_sqlite_provider(pool).await;
1520 let ethereum_hash = H256::from([1u8; 32]);
1521 let substrate_hash = H256::from([2u8; 32]);
1522 let block_map = BlockHashMap::new(substrate_hash, ethereum_hash);
1523
1524 insert_block_mapping(&provider.db_ctx.pool, &block_map).await?;
1526 assert_eq!(
1527 provider.get_substrate_hash(&block_map.ethereum_hash).await,
1528 Some(block_map.substrate_hash)
1529 );
1530
1531 provider.remove(&[block_map.clone()]).await?;
1533
1534 assert_eq!(provider.get_substrate_hash(&block_map.ethereum_hash).await, None);
1536
1537 Ok(())
1538 }
1539
1540 #[sqlx::test]
1541 async fn test_logs_with_ethereum_block_hash_mapping(pool: SqlitePool) -> anyhow::Result<()> {
1542 let provider = setup_sqlite_provider(pool).await;
1543 let ethereum_hash = H256::from([1u8; 32]);
1544 let substrate_hash = H256::from([2u8; 32]);
1545 let block_number = 1u64;
1546
1547 let log = Log {
1549 block_hash: ethereum_hash,
1550 block_number: block_number.into(),
1551 address: H160::from([1u8; 20]),
1552 topics: vec![H256::from([1u8; 32])],
1553 transaction_hash: H256::from([3u8; 32]),
1554 transaction_index: U256::from(0),
1555 log_index: U256::from(0),
1556 data: Some(vec![0u8; 32].into()),
1557 ..Default::default()
1558 };
1559
1560 let block = MockBlockInfo { hash: substrate_hash, number: block_number };
1562 let receipts = vec![(
1563 TransactionSigned::default(),
1564 ReceiptInfo {
1565 logs: vec![log.clone()],
1566 transaction_hash: log.transaction_hash,
1567 transaction_index: log.transaction_index,
1568 ..Default::default()
1569 },
1570 )];
1571 provider.insert(&block, &receipts, ðereum_hash).await?;
1572
1573 let logs = provider
1575 .logs(
1576 Some(Filter::new().at_block_hash(B256::from(ethereum_hash.0))),
1577 mock_resolve_block_number_with_latest(block.number.into()),
1578 )
1579 .await?;
1580 assert_eq!(logs, vec![log]);
1581
1582 Ok(())
1583 }
1584
1585 #[sqlx::test]
1586 async fn test_mapping_count(pool: SqlitePool) -> anyhow::Result<()> {
1587 let provider = setup_sqlite_provider(pool).await;
1588
1589 assert_eq!(count(&provider.db_ctx.pool, "eth_to_substrate_blocks", None).await, 0);
1591
1592 let block_map1 = BlockHashMap::new(H256::from([1u8; 32]), H256::from([2u8; 32]));
1593 let block_map2 = BlockHashMap::new(H256::from([3u8; 32]), H256::from([4u8; 32]));
1594
1595 insert_block_mapping(&provider.db_ctx.pool, &block_map1).await?;
1597 insert_block_mapping(&provider.db_ctx.pool, &block_map2).await?;
1598
1599 assert_eq!(count(&provider.db_ctx.pool, "eth_to_substrate_blocks", None).await, 2);
1600
1601 provider.remove(&[block_map1]).await?;
1603 assert_eq!(count(&provider.db_ctx.pool, "eth_to_substrate_blocks", None).await, 1);
1604
1605 Ok(())
1606 }
1607
1608 #[sqlx::test]
1609 async fn restore_first_evm_block_clears_stale(pool: SqlitePool) -> anyhow::Result<()> {
1610 let provider = setup_sqlite_provider(pool).await;
1611
1612 provider
1614 .set_sync_label(ChainMetadata::FirstEvmBlock, SyncCheckpoint::from_number(42))
1615 .await?;
1616
1617 provider.restore_first_evm_block().await?;
1620
1621 assert_eq!(provider.first_evm_block(), None);
1623
1624 let checkpoint = provider.get_sync_label(ChainMetadata::FirstEvmBlock).await?;
1626 assert!(checkpoint.is_none());
1627 Ok(())
1628 }
1629
1630 #[sqlx::test]
1631 async fn advance_sync_label_only_increases(pool: SqlitePool) -> anyhow::Result<()> {
1632 let provider = setup_sqlite_provider(pool).await;
1633 let hash_a = H256::repeat_byte(0xAA);
1634 let hash_b = H256::repeat_byte(0xBB);
1635
1636 provider
1638 .advance_sync_label(SyncLabel::Head, SyncCheckpoint::new(100, hash_a))
1639 .await?;
1640 let checkpoint = provider.get_sync_label(SyncLabel::Head).await?.unwrap();
1641 assert_eq!((checkpoint.block_number, checkpoint.block_hash), (100, Some(hash_a)));
1642
1643 provider
1645 .advance_sync_label(SyncLabel::Head, SyncCheckpoint::new(200, hash_b))
1646 .await?;
1647 let checkpoint = provider.get_sync_label(SyncLabel::Head).await?.unwrap();
1648 assert_eq!((checkpoint.block_number, checkpoint.block_hash), (200, Some(hash_b)));
1649
1650 provider
1652 .advance_sync_label(SyncLabel::Head, SyncCheckpoint::new(50, hash_a))
1653 .await?;
1654 provider
1655 .advance_sync_label(SyncLabel::Head, SyncCheckpoint::new(200, hash_a))
1656 .await?;
1657 let checkpoint = provider.get_sync_label(SyncLabel::Head).await?.unwrap();
1658 assert_eq!((checkpoint.block_number, checkpoint.block_hash), (200, Some(hash_b)));
1659
1660 Ok(())
1661 }
1662
1663 #[sqlx::test]
1664 async fn recede_sync_label_only_decreases(pool: SqlitePool) -> anyhow::Result<()> {
1665 let provider = setup_sqlite_provider(pool).await;
1666 let hash_a = H256::repeat_byte(0xAA);
1667 let hash_b = H256::repeat_byte(0xBB);
1668
1669 provider
1671 .recede_sync_label(SyncLabel::Tail, SyncCheckpoint::new(100, hash_a))
1672 .await?;
1673 let checkpoint = provider.get_sync_label(SyncLabel::Tail).await?.unwrap();
1674 assert_eq!((checkpoint.block_number, checkpoint.block_hash), (100, Some(hash_a)));
1675
1676 provider
1678 .recede_sync_label(SyncLabel::Tail, SyncCheckpoint::new(50, hash_b))
1679 .await?;
1680 let checkpoint = provider.get_sync_label(SyncLabel::Tail).await?.unwrap();
1681 assert_eq!((checkpoint.block_number, checkpoint.block_hash), (50, Some(hash_b)));
1682
1683 provider
1685 .recede_sync_label(SyncLabel::Tail, SyncCheckpoint::new(200, hash_a))
1686 .await?;
1687 provider
1688 .recede_sync_label(SyncLabel::Tail, SyncCheckpoint::new(50, hash_a))
1689 .await?;
1690 let checkpoint = provider.get_sync_label(SyncLabel::Tail).await?.unwrap();
1691 assert_eq!((checkpoint.block_number, checkpoint.block_hash), (50, Some(hash_b)));
1692
1693 Ok(())
1694 }
1695
1696 #[tokio::test]
1697 async fn is_before_earliest_block_edge_cases() {
1698 let extractor = ReceiptExtractor::new_mock();
1700 extractor.set_first_evm_block(10);
1701 let provider = mock_provider().with_extractor(extractor);
1702
1703 let huge = BlockNumberOrTag::Number(u64::MAX);
1704 assert!(!provider.is_before_earliest_block(&huge));
1705
1706 let just_over = BlockNumberOrTag::Number(u32::MAX as u64 + 1);
1707 assert!(!provider.is_before_earliest_block(&just_over));
1708
1709 let provider = mock_provider();
1711 assert!(!provider.is_before_earliest_block(&BlockNumberOrTag::Number(0)));
1712 assert!(!provider.is_before_earliest_block(&BlockNumberOrTag::Number(1_000_000)));
1713
1714 assert!(!provider.is_before_earliest_block(&BlockNumberOrTag::Latest));
1716 }
1717
1718 #[sqlx::test]
1719 async fn persistent_mode_caps_in_memory_map(pool: SqlitePool) -> anyhow::Result<()> {
1720 let provider = mock_provider()
1722 .with_db_ctx(DbContext::new(pool, DbContext::DEFAULT_MAX_VARIABLE_NUMBER));
1723
1724 let start_block: u64 = 1;
1726 let n = MAX_CACHED_BLOCKS + 1;
1727 let end_block = start_block + n as u64;
1728 for i in start_block..end_block {
1729 let block = MockBlockInfo { hash: H256::from_low_u64_be(i), number: i as _ };
1730 let receipts = vec![(
1731 TransactionSigned::default(),
1732 ReceiptInfo {
1733 transaction_hash: H256::from_low_u64_be(i),
1734 logs: vec![Log {
1735 block_hash: block.hash,
1736 transaction_hash: H256::from_low_u64_be(i),
1737 ..Default::default()
1738 }],
1739 ..Default::default()
1740 },
1741 )];
1742 let ethereum_hash = H256::from_low_u64_be(i + 1);
1743 provider.insert(&block, &receipts, ðereum_hash).await?;
1744 }
1745
1746 let map = provider.block_number_to_hashes.lock().await;
1748 assert_eq!(map.len(), MAX_CACHED_BLOCKS);
1749
1750 assert!(!map.contains_key(&1));
1752 assert!(map.contains_key(&2));
1753 assert!(map.contains_key(&(MAX_CACHED_BLOCKS as u64 + 1)));
1754 drop(map);
1755
1756 assert_eq!(count(&provider.db_ctx.pool, "eth_to_substrate_blocks", None).await, n);
1758
1759 Ok(())
1760 }
1761
1762 fn make_hash(i: usize, fill: u8) -> H256 {
1763 let mut hash = [fill; 32];
1764 hash[..8].copy_from_slice(&i.to_le_bytes());
1765 H256::from(hash)
1766 }
1767
1768 fn make_receipts(
1769 tx_offset: usize,
1770 n_tx: usize,
1771 n_logs: usize,
1772 ) -> Vec<(TransactionSigned, ReceiptInfo)> {
1773 let mut receipts = Vec::with_capacity(n_tx);
1774
1775 for i in 0..n_tx {
1776 let transaction_hash = make_hash(tx_offset + i, 0x00);
1777
1778 let mut logs = Vec::with_capacity(n_logs);
1779 for j in 0..n_logs {
1780 logs.push(Log { transaction_hash, log_index: U256::from(j), ..Default::default() });
1781 }
1782
1783 receipts.push((
1784 TransactionSigned::default(),
1785 ReceiptInfo {
1786 transaction_hash,
1787 transaction_index: U256::from(i),
1788 logs,
1789 ..Default::default()
1790 },
1791 ));
1792 }
1793
1794 receipts
1795 }
1796
1797 async fn assert_receipts_inserted(
1798 provider: &ReceiptProvider<MockBlockInfoProvider>,
1799 block: &MockBlockInfo,
1800 ethereum_hash: &H256,
1801 receipts: &[(TransactionSigned, ReceiptInfo)],
1802 ) {
1803 let mut expected_logs = 0;
1804 for (_, receipt) in receipts {
1805 assert_eq!(
1806 provider.find_transaction(&receipt.transaction_hash).await,
1807 Some((block.hash(), receipt.transaction_index.as_u32() as usize))
1808 );
1809 expected_logs += receipt.logs.len();
1810 }
1811 assert_eq!(count(&provider.db_ctx.pool, "logs", Some(*ethereum_hash)).await, expected_logs);
1812 }
1813
1814 #[sqlx::test]
1815 async fn test_bulk_insert(pool: SqlitePool) -> anyhow::Result<()> {
1816 let provider = setup_sqlite_provider(pool).await.with_keep_latest(None);
1817 let tx_chunk = provider.db_ctx.tx_insert_chunk_size;
1818 let log_chunk = provider.db_ctx.log_insert_chunk_size;
1819
1820 let cases = [
1821 (tx_chunk, 1), (tx_chunk + 1, log_chunk), (1000, 3), ];
1825
1826 let mut tx_offset = 0;
1827 for (i, (n_tx, n_logs)) in cases.into_iter().enumerate() {
1828 let block = MockBlockInfo { hash: make_hash(i, 0x00), number: i as u64 + 1 };
1829 let ethereum_hash = make_hash(i, 0xff);
1830 let receipts = make_receipts(tx_offset, n_tx, n_logs);
1831 tx_offset += n_tx;
1832 provider.insert(&block, &receipts, ðereum_hash).await?;
1833 assert_receipts_inserted(&provider, &block, ðereum_hash, &receipts).await;
1834 }
1835 Ok(())
1836 }
1837
1838 #[sqlx::test]
1839 async fn test_duplicate_insert_succeeds(pool: SqlitePool) -> anyhow::Result<()> {
1840 let provider = setup_sqlite_provider(pool).await.with_keep_latest(None);
1841 let block = MockBlockInfo { hash: make_hash(0, 0xAA), number: 1 };
1842 let ethereum_hash = make_hash(0, 0xBB);
1843 let receipts = make_receipts(0, 5, 3);
1844
1845 provider.insert_into_db(&block, &receipts, ðereum_hash).await?;
1847 assert_eq!(count(&provider.db_ctx.pool, "transaction_hashes", None).await, 5);
1848 assert_eq!(count(&provider.db_ctx.pool, "logs", Some(ethereum_hash)).await, 15);
1849 assert_eq!(count(&provider.db_ctx.pool, "eth_to_substrate_blocks", None).await, 1);
1850
1851 sqlx::query("DELETE FROM eth_to_substrate_blocks")
1853 .execute(&provider.db_ctx.pool)
1854 .await?;
1855 assert_eq!(count(&provider.db_ctx.pool, "eth_to_substrate_blocks", None).await, 0);
1856
1857 provider.insert_into_db(&block, &receipts, ðereum_hash).await?;
1859
1860 assert_eq!(count(&provider.db_ctx.pool, "transaction_hashes", None).await, 5);
1862 assert_eq!(count(&provider.db_ctx.pool, "logs", Some(ethereum_hash)).await, 15);
1863 assert_eq!(count(&provider.db_ctx.pool, "eth_to_substrate_blocks", None).await, 1);
1864
1865 Ok(())
1866 }
1867
1868 #[sqlx::test]
1869 async fn test_insert_empty_receipts(pool: SqlitePool) -> anyhow::Result<()> {
1870 let provider = setup_sqlite_provider(pool).await.with_keep_latest(None);
1871 let block = MockBlockInfo { hash: H256::from([1u8; 32]), number: 1 };
1872 let ethereum_hash = H256::from([2u8; 32]);
1873
1874 provider.insert(&block, &[], ðereum_hash).await?;
1875
1876 assert_eq!(count(&provider.db_ctx.pool, "eth_to_substrate_blocks", None).await, 1);
1878 assert_eq!(count(&provider.db_ctx.pool, "transaction_hashes", None).await, 0);
1879 assert_eq!(count(&provider.db_ctx.pool, "logs", None).await, 0);
1880
1881 let receipts = make_receipts(0, 3, 2);
1883 provider.insert(&block, &receipts, ðereum_hash).await?;
1884 assert_eq!(count(&provider.db_ctx.pool, "eth_to_substrate_blocks", None).await, 1);
1885 assert_eq!(count(&provider.db_ctx.pool, "transaction_hashes", None).await, 0);
1886 assert_eq!(count(&provider.db_ctx.pool, "logs", None).await, 0);
1887
1888 Ok(())
1889 }
1890
1891 #[sqlx::test]
1892 async fn test_bulk_delete(pool: SqlitePool) -> anyhow::Result<()> {
1893 let db_ctx = DbContext::new(pool, DbContext::LOG_COLUMNS);
1895 let provider = mock_provider().with_db_ctx(db_ctx).with_keep_latest(None);
1896
1897 let n_blocks = 25;
1898 let n_tx_per_block = 5;
1899 let n_logs_per_receipt = 3;
1900 let mut block_mappings = Vec::new();
1901
1902 for i in 0..n_blocks {
1903 let block = MockBlockInfo { hash: make_hash(i, 0xAA), number: i as u64 + 1 };
1904 let ethereum_hash = make_hash(i, 0xBB);
1905 let receipts = make_receipts(i * n_tx_per_block, n_tx_per_block, n_logs_per_receipt);
1906 provider.insert_into_db(&block, &receipts, ðereum_hash).await?;
1907 block_mappings.push(BlockHashMap::new(block.hash, ethereum_hash));
1908 }
1909
1910 assert_eq!(
1911 count(&provider.db_ctx.pool, "transaction_hashes", None).await,
1912 n_blocks * n_tx_per_block
1913 );
1914 assert_eq!(
1915 count(&provider.db_ctx.pool, "logs", None).await,
1916 n_blocks * n_tx_per_block * n_logs_per_receipt
1917 );
1918 assert_eq!(count(&provider.db_ctx.pool, "eth_to_substrate_blocks", None).await, n_blocks);
1919
1920 provider.remove(&block_mappings).await?;
1921
1922 assert_eq!(count(&provider.db_ctx.pool, "transaction_hashes", None).await, 0);
1923 assert_eq!(count(&provider.db_ctx.pool, "logs", None).await, 0);
1924 assert_eq!(count(&provider.db_ctx.pool, "eth_to_substrate_blocks", None).await, 0);
1925
1926 Ok(())
1927 }
1928
1929 #[sqlx::test]
1930 async fn test_get_processed_eth_block_hash(pool: SqlitePool) -> anyhow::Result<()> {
1931 let provider = setup_sqlite_provider(pool).await;
1932 let block = MockBlockInfo { hash: H256::from([0xAA; 32]), number: 10 };
1933 let ethereum_hash = H256::from([0xBB; 32]);
1934 let receipts = vec![(TransactionSigned::default(), ReceiptInfo::default())];
1935
1936 assert!(provider.get_processed_eth_block_hash(10, block.hash).await.is_none());
1938
1939 provider.insert(&block, &receipts, ðereum_hash).await?;
1941 assert_eq!(
1942 provider.get_processed_eth_block_hash(10, block.hash).await,
1943 Some(ethereum_hash)
1944 );
1945
1946 assert!(
1948 provider
1949 .get_processed_eth_block_hash(10, H256::from([0xCC; 32]))
1950 .await
1951 .is_none()
1952 );
1953
1954 assert!(provider.get_processed_eth_block_hash(11, block.hash).await.is_none());
1956
1957 Ok(())
1958 }
1959
1960 #[sqlx::test]
1961 async fn test_logs_by_block_number(pool: SqlitePool) -> anyhow::Result<()> {
1962 let provider = setup_sqlite_provider(pool).await;
1963 let substrate_hash = H256::from([0xAA; 32]);
1964 let tx_hash = H256::from([0xBB; 32]);
1965 let block = MockBlockInfo { hash: substrate_hash, number: 42 };
1966 let ethereum_hash = H256::from([0xCC; 32]);
1967
1968 let log0 = Log {
1969 block_hash: ethereum_hash,
1970 block_number: U256::from(42),
1971 transaction_hash: tx_hash,
1972 log_index: U256::from(0),
1973 address: H160::from([0x01; 20]),
1974 ..Default::default()
1975 };
1976 let log1 = Log {
1977 block_hash: ethereum_hash,
1978 block_number: U256::from(42),
1979 transaction_hash: tx_hash,
1980 log_index: U256::from(1),
1981 address: H160::from([0x02; 20]),
1982 ..Default::default()
1983 };
1984
1985 let receipts = vec![(
1986 TransactionSigned::default(),
1987 ReceiptInfo {
1988 transaction_hash: tx_hash,
1989 block_hash: ethereum_hash,
1990 logs: vec![log0.clone(), log1.clone()],
1991 ..Default::default()
1992 },
1993 )];
1994
1995 let logs = provider.logs_by_block_number(42, ethereum_hash).await?;
1997 assert!(logs.is_empty());
1998
1999 provider.insert(&block, &receipts, ðereum_hash).await?;
2000
2001 let logs = provider.logs_by_block_number(42, ethereum_hash).await?;
2003 assert_eq!(logs.len(), 2);
2004 assert_eq!(logs[0].address, log0.address);
2005 assert_eq!(logs[1].address, log1.address);
2006 assert_eq!(logs[0].log_index, U256::from(0));
2007 assert_eq!(logs[1].log_index, U256::from(1));
2008
2009 let logs = provider.logs_by_block_number(43, ethereum_hash).await?;
2011 assert!(logs.is_empty());
2012
2013 let logs = provider.logs_by_block_number(42, H256::from([0xDD; 32])).await?;
2015 assert!(logs.is_empty());
2016
2017 Ok(())
2018 }
2019}