referrerpolicy=no-referrer-when-downgrade

pallet_revive_eth_rpc/
receipt_provider.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.
17use 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
37/// Parse a SQLite row from the `logs` table into a [`Log`].
38fn 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/// SQLite connection pool with precomputed bulk-insert chunk sizes.
70#[derive(Clone)]
71pub struct DbContext {
72	pool: SqlitePool,
73	/// Max bound parameters per query.
74	max_variable_number: usize,
75	/// Chunk size for bulk INSERT into `transaction_hashes`.
76	tx_insert_chunk_size: usize,
77	/// Chunk size for bulk INSERT into `logs`.
78	log_insert_chunk_size: usize,
79}
80
81impl DbContext {
82	/// Conservative default for `SQLITE_LIMIT_VARIABLE_NUMBER`; SQLite >=3.32 uses 32766.
83	pub const DEFAULT_MAX_VARIABLE_NUMBER: usize = 999;
84	/// Columns in the `transaction_hashes` table.
85	const TX_HASH_COLUMNS: usize = 3;
86	/// Columns in the `logs` table.
87	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/// ReceiptProvider stores transaction receipts and logs in a SQLite database.
105#[derive(Clone)]
106pub struct ReceiptProvider<B: BlockInfoProvider = SubxtBlockInfoProvider> {
107	/// The database pool.
108	db_ctx: DbContext,
109	/// The block provider used to fetch blocks, and reconstruct receipts.
110	block_provider: B,
111	/// A means to extract receipts from extrinsics.
112	receipt_extractor: ReceiptExtractor,
113	/// When `Some`, old blocks will be pruned.
114	keep_latest_n_blocks: Option<usize>,
115	/// A Map of the latest block numbers to block hashes.
116	block_number_to_hashes: Arc<Mutex<BTreeMap<SubstrateBlockNumber, BlockHashMap>>>,
117}
118
119/// Substrate block to Ethereum block mapping
120#[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
132/// Provides information about a block,
133/// This is an abstraction on top of [`SubstrateBlock`] that can't be mocked in tests.
134/// Can be removed once <https://github.com/paritytech/subxt/issues/1883> is fixed.
135pub trait BlockInfo {
136	/// Returns the block hash.
137	fn hash(&self) -> H256;
138	/// Returns the block number.
139	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
151/// Maximum number of entries kept in the block to hash map.
152pub const MAX_CACHED_BLOCKS: usize = 256;
153
154/// Upsert a sync label row, updating only when the existing `block_number`
155/// compares with `$op` against the new value. `$op` must be `"<"` or `">"`.
156macro_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	/// Create a new `ReceiptProvider` with the given database URL and block provider.
198	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	/// Returns `true` if the block is before the auto-discovered `first_evm_block`.
222	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	/// The auto-discovered first EVM block, or `None` if not yet discovered.
236	pub fn first_evm_block(&self) -> Option<SubstrateBlockNumber> {
237		self.receipt_extractor.first_evm_block()
238	}
239
240	/// Set the auto-discovered first EVM block (in-memory + persisted to DB).
241	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	/// Restore `first_evm_block` from DB, clearing it if the boundary has shifted.
251	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		// Stale if evm_first no longer has an EVM hash, or its predecessor now does.
270		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	// Get block hash and transaction index by transaction hash
290	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	/// Get the Substrate block hash for the given Ethereum block hash.
319	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	/// Get the Ethereum block hash for the given Substrate block hash.
346	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	/// Deletes older records from the database.
373	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	/// Read a sync label entry.
411	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	/// Upsert a sync label entry.
448	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	/// Delete a sync label entry.
471	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	/// Atomically update a sync label entry only if the new block number is strictly higher.
485	///
486	/// Inserts the row if it doesn't exist yet.
487	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	/// Atomically update a sync label entry only if the new block number is lower.
497	///
498	/// Inserts the row if it doesn't exist yet.
499	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	/// Look up the ethereum block hash for a previously processed block from the in-memory cache.
509	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	/// Fetch receipts from the given block, using a pre-fetched ethereum block hash.
523	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	/// Fetch all receipts for the given block.
534	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	/// Like [`Self::insert_block_receipts`] but writes only to the DB (no cache update).
542	/// Used for historic sync where fork detection is unnecessary.
543	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	/// Insert pre-extracted receipts and update the block cache (with fork detection).
557	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	/// Insert receipts into the provider, updating the in-memory block cache for fork detection.
567	///
568	/// Note: Can be merged into `insert_block_receipts` once <https://github.com/paritytech/subxt/issues/1883> is fixed and subxt let
569	/// us create Mock `SubstrateBlock`
570	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	/// Handle fork detection (always) and DB pruning (temporary mode only).
583	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		// Fork? - If inserting the same block number with a different hash, remove the old ones.
592		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				// Now loop through the blocks that were building on top of the old fork and remove
597				// them.
598				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			// If we have more blocks than we should keep, remove the oldest ones by count
609			// (not by block number range, to handle gaps correctly)
610			while block_number_to_hash.len() > keep_latest_n_blocks {
611				// Remove the block with the smallest number (first in BTreeMap)
612				if let Some((_, block_map)) = block_number_to_hash.pop_first() {
613					to_remove.push(block_map);
614				}
615			}
616		} else {
617			// Evict oldest entries to prevent unbounded growth.
618			// Forks deeper than MAX_CACHED_BLOCKS(256) are unlikely.
619			while block_number_to_hash.len() > MAX_CACHED_BLOCKS {
620				block_number_to_hash.pop_first();
621			}
622		}
623
624		// Release the lock.
625		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	/// Insert receipts into the database without updating the in-memory block cache.
636	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		// Check if mapping already exists (eg. added when processing best block and we are now
650		// processing finalized block)
651		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		// Assuming that if no mapping exists then no relevant entries in transaction_hashes and
658		// logs exist
659		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	/// Get logs that match the given filter.
721	///
722	/// `resolve_block_number` converts a [`BlockNumberOrTag`] to a concrete block number.
723	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				// Read the latest block *after* resolving the tags.
745				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	/// Fetch all logs for a given block from the database.
813	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	/// Get the number of receipts per block.
844	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	/// Return all transaction hashes for the given block hash.
863	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	/// Get the receipt for the given block hash and transaction index.
889	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	/// Get the receipt for the given transaction hash.
904	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	/// Get the signed transaction for the given transaction hash.
932	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	/// Test resolver that handles Latest โ†’ `latest` and Earliest โ†’ 0.
996	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				// The mock only tracks `latest`, so the remaining tags resolve to it.
1004				_ => 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, &ethereum_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, &ethereum_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		// Build 4 blocks on consecutive heights: 0,1,2,3.
1114		let (block0, receipts, ethereum_hash_0) = build_block(0, 0);
1115		provider.insert(&block0, &receipts, &ethereum_hash_0).await?;
1116		let (block1, receipts, ethereum_hash_1) = build_block(1, 1);
1117		provider.insert(&block1, &receipts, &ethereum_hash_1).await?;
1118		let (block2, receipts, ethereum_hash_2) = build_block(2, 2);
1119		provider.insert(&block2, &receipts, &ethereum_hash_2).await?;
1120		let (block3, receipts, ethereum_hash_3) = build_block(3, 3);
1121		provider.insert(&block3, &receipts, &ethereum_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		// Now build another block on height 1.
1138		let (fork_block, receipts, ethereum_hash_fork) = build_block(4, 1);
1139		provider.insert(&fork_block, &receipts, &ethereum_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		// Build two blocks at the same height with the same transaction hash
1162		let tx_hash = H256::from([42u8; 32]);
1163
1164		// Block A at height 1
1165		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, &ethereum_hash_a).await?;
1177
1178		// Verify transaction points to block A
1179		let (found_hash, _) = provider.find_transaction(&tx_hash).await.unwrap();
1180		assert_eq!(found_hash, block_a.hash);
1181
1182		// Clear the in-memory map to simulate server restart
1183		provider.block_number_to_hashes.lock().await.clear();
1184
1185		// Block B at same height 1 (re-org) with SAME transaction
1186		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, // Same tx hash!
1192				transaction_index: U256::from(0),
1193				..Default::default()
1194			},
1195		)];
1196
1197		// This should NOT fail with UNIQUE constraint violation
1198		provider.insert(&block_b, &receipts_b, &ethereum_hash_b).await?;
1199
1200		// Transaction should now point to block B
1201		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, &ethereum_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				&ethereum_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				&ethereum_hash2,
1286			)
1287			.await?;
1288
1289		let resolve_block_number = mock_resolve_block_number_with_latest(block2.number.into());
1290
1291		// Empty filter
1292		let logs = provider.logs(None, &resolve_block_number).await?;
1293		assert_eq!(logs, vec![log2.clone()]);
1294
1295		// from_block filter
1296		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		// from_block filter (using latest block)
1302		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		// to_block filter
1308		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		// block_hash filter
1314		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		// single address
1323		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		// multiple addresses
1336		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		// single topic
1348		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		// multiple topic
1361		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		// multiple topic for topic_0
1375		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		// Altogether
1387		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 mapping
1470		insert_block_mapping(&provider.db_ctx.pool, &block_map).await?;
1471
1472		// Test forward lookup
1473		let resolved = provider.get_substrate_hash(&ethereum_hash).await;
1474		assert_eq!(resolved, Some(substrate_hash));
1475
1476		// Test reverse lookup
1477		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 mappings
1494		insert_block_mapping(&provider.db_ctx.pool, &block_map1).await?;
1495		insert_block_mapping(&provider.db_ctx.pool, &block_map2).await?;
1496
1497		// Verify they exist
1498		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		// Remove one mapping
1508		provider.remove(&[block_map1]).await?;
1509
1510		// Verify removal
1511		assert_eq!(provider.get_substrate_hash(&ethereum_hash1).await, None);
1512		assert_eq!(provider.get_substrate_hash(&ethereum_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 mapping
1525		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		// Remove substrate block (this should also remove the mapping)
1532		provider.remove(&[block_map.clone()]).await?;
1533
1534		// Mapping should be gone
1535		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		// Create a log with ethereum hash
1548		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		// Insert the log
1561		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, &ethereum_hash).await?;
1572
1573		// Query logs using Ethereum block hash (should resolve to substrate hash)
1574		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		// Initially no mappings
1590		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 some mappings
1596		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		// Remove one
1602		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		// Persist first_evm_block = 42.
1613		provider
1614			.set_sync_label(ChainMetadata::FirstEvmBlock, SyncCheckpoint::from_number(42))
1615			.await?;
1616
1617		// MockBlockInfoProvider returns no blocks, so has_evm_hash is always false.
1618		// This means evm_first=42 is stale (no longer has an EVM hash).
1619		provider.restore_first_evm_block().await?;
1620
1621		// The value should have been cleared (not restored to the extractor).
1622		assert_eq!(provider.first_evm_block(), None);
1623
1624		// DB row should have been deleted.
1625		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		// First insert creates the row.
1637		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		// Higher value advances.
1644		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		// Lower and equal values are ignored (strict >).
1651		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		// First insert creates the row.
1670		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		// Lower value recedes.
1677		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		// Higher and equal values are ignored (strict <).
1684		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		// U256 > u32::MAX should never be considered "before floor"
1699		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		// Sentinel first_evm_block (u32::MAX) is permissive โ€” no queries rejected.
1710		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		// Tag-based queries are never rejected.
1715		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		// Persistent DB mode: keep_latest_n_blocks = None
1721		let provider = mock_provider()
1722			.with_db_ctx(DbContext::new(pool, DbContext::DEFAULT_MAX_VARIABLE_NUMBER));
1723
1724		// Insert more than MAX_CACHED_BLOCKS blocks.
1725		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, &ethereum_hash).await?;
1744		}
1745
1746		// The map is capped at MAX_CACHED_BLOCKS.
1747		let map = provider.block_number_to_hashes.lock().await;
1748		assert_eq!(map.len(), MAX_CACHED_BLOCKS);
1749
1750		// The oldest block (1) should have been evicted, keeping blocks 2..=MAX+1.
1751		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		// All blocks are still in the DB.
1757		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),             // exact tx chunk boundary
1822			(tx_chunk + 1, log_chunk), // crosses tx boundary; exact log chunk boundary
1823			(1000, 3),                 // multiple tx and log chunks
1824		];
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, &ethereum_hash).await?;
1833			assert_receipts_inserted(&provider, &block, &ethereum_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		// First insert.
1846		provider.insert_into_db(&block, &receipts, &ethereum_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		// Delete only the block mapping so the EXISTS guard won't short-circuit.
1852		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		// Second insert hits the actual INSERT OR REPLACE statements.
1858		provider.insert_into_db(&block, &receipts, &ethereum_hash).await?;
1859
1860		// Row counts unchanged โ€” no duplicates.
1861		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, &[], &ethereum_hash).await?;
1875
1876		// Block mapping is stored as a deduplication marker.
1877		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		// Second insert for the same block is a no-op, even with receipts.
1882		let receipts = make_receipts(0, 3, 2);
1883		provider.insert(&block, &receipts, &ethereum_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		// Use the smallest valid limit to force chunked INSERTs and DELETEs.
1894		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, &ethereum_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		// Not cached yet
1937		assert!(provider.get_processed_eth_block_hash(10, block.hash).await.is_none());
1938
1939		// Insert also populates the in-memory cache
1940		provider.insert(&block, &receipts, &ethereum_hash).await?;
1941		assert_eq!(
1942			provider.get_processed_eth_block_hash(10, block.hash).await,
1943			Some(ethereum_hash)
1944		);
1945
1946		// Wrong hash for same block number
1947		assert!(
1948			provider
1949				.get_processed_eth_block_hash(10, H256::from([0xCC; 32]))
1950				.await
1951				.is_none()
1952		);
1953
1954		// Wrong block number
1955		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		// No logs before insert
1996		let logs = provider.logs_by_block_number(42, ethereum_hash).await?;
1997		assert!(logs.is_empty());
1998
1999		provider.insert(&block, &receipts, &ethereum_hash).await?;
2000
2001		// Logs returned in log_index order
2002		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		// Different block number returns empty
2010		let logs = provider.logs_by_block_number(43, ethereum_hash).await?;
2011		assert!(logs.is_empty());
2012
2013		// Wrong ethereum hash returns empty
2014		let logs = provider.logs_by_block_number(42, H256::from([0xDD; 32])).await?;
2015		assert!(logs.is_empty());
2016
2017		Ok(())
2018	}
2019}