1use std::{
5 collections::BTreeSet,
6 marker::PhantomData,
7 ops::{Range, RangeInclusive},
8};
9
10use crate::{
11 providers::{
12 history_info, rocksdb::RocksDBBatch, HistoryInfo, StaticFileProvider,
13 StaticFileProviderRWRefMut,
14 },
15 StaticFileProviderFactory,
16};
17use alloy_primitives::{map::HashMap, Address, BlockNumber, TxHash, TxNumber, B256};
18use rayon::slice::ParallelSliceMut;
19use reth_db::{
20 cursor::{DbCursorRO, DbDupCursorRW},
21 models::{AccountBeforeTx, StorageBeforeTx},
22 static_file::TransactionSenderMask,
23 table::Value,
24 transaction::{CursorMutTy, CursorTy, DbTx, DbTxMut, DupCursorMutTy, DupCursorTy},
25};
26use reth_db_api::{
27 cursor::DbCursorRW,
28 models::{storage_sharded_key::StorageShardedKey, BlockNumberAddress, ShardedKey},
29 tables,
30 tables::BlockNumberList,
31};
32use reth_errors::ProviderError;
33use reth_node_types::NodePrimitives;
34use reth_primitives_traits::{ReceiptTy, StorageEntry};
35use reth_static_file_types::StaticFileSegment;
36use reth_storage_api::{
37 ChangeSetReader, DBProvider, DbTxProvider, NodePrimitivesProvider, StorageSettingsCache,
38};
39use reth_storage_errors::provider::ProviderResult;
40use strum::{Display, EnumIs};
41
42type EitherReaderTy<'a, 'db, P, T> = EitherReader<
44 'a,
45 'db,
46 CursorTy<<P as DbTxProvider>::Tx, T>,
47 <P as NodePrimitivesProvider>::Primitives,
48>;
49
50type DupEitherReaderTy<'a, 'db, P, T> = EitherReader<
52 'a,
53 'db,
54 DupCursorTy<<P as DbTxProvider>::Tx, T>,
55 <P as NodePrimitivesProvider>::Primitives,
56>;
57
58type DupEitherWriterTy<'a, P, T> = EitherWriter<
60 'a,
61 DupCursorMutTy<<P as DbTxProvider>::Tx, T>,
62 <P as NodePrimitivesProvider>::Primitives,
63>;
64
65type EitherWriterTy<'a, P, T> = EitherWriter<
67 'a,
68 CursorMutTy<<P as DbTxProvider>::Tx, T>,
69 <P as NodePrimitivesProvider>::Primitives,
70>;
71
72pub type RocksBatchArg<'a> = crate::providers::rocksdb::RocksDBBatch<'a>;
74
75pub type RawRocksDBBatch = rocksdb::WriteBatchWithTransaction<true>;
77
78pub type RocksDBRefArg<'a, 'db> = Option<&'a crate::providers::rocksdb::RocksReadSnapshot<'db>>;
83
84#[derive(Debug, Display)]
86pub enum EitherWriter<'a, CURSOR, N> {
87 Database(CURSOR),
89 StaticFile(StaticFileProviderRWRefMut<'a, N>),
91 RocksDB(RocksDBBatch<'a>),
93}
94
95impl<'a> EitherWriter<'a, (), ()> {
96 pub fn new_receipts<P>(
98 provider: &'a P,
99 block_number: BlockNumber,
100 ) -> ProviderResult<EitherWriterTy<'a, P, tables::Receipts<ReceiptTy<P::Primitives>>>>
101 where
102 P: DBProvider + NodePrimitivesProvider + StorageSettingsCache + StaticFileProviderFactory,
103 P::Tx: DbTxMut,
104 ReceiptTy<P::Primitives>: Value,
105 {
106 if Self::receipts_destination(provider).is_static_file() {
107 Ok(EitherWriter::StaticFile(
108 provider.get_static_file_writer(block_number, StaticFileSegment::Receipts)?,
109 ))
110 } else {
111 Ok(EitherWriter::Database(
112 provider.tx_ref().cursor_write::<tables::Receipts<ReceiptTy<P::Primitives>>>()?,
113 ))
114 }
115 }
116
117 pub fn new_senders<P>(
119 provider: &'a P,
120 block_number: BlockNumber,
121 ) -> ProviderResult<EitherWriterTy<'a, P, tables::TransactionSenders>>
122 where
123 P: DBProvider + NodePrimitivesProvider + StorageSettingsCache + StaticFileProviderFactory,
124 P::Tx: DbTxMut,
125 {
126 if EitherWriterDestination::senders(provider).is_static_file() {
127 Ok(EitherWriter::StaticFile(
128 provider
129 .get_static_file_writer(block_number, StaticFileSegment::TransactionSenders)?,
130 ))
131 } else {
132 Ok(EitherWriter::Database(
133 provider.tx_ref().cursor_write::<tables::TransactionSenders>()?,
134 ))
135 }
136 }
137
138 pub fn new_account_changesets<P>(
141 provider: &'a P,
142 block_number: BlockNumber,
143 ) -> ProviderResult<DupEitherWriterTy<'a, P, tables::AccountChangeSets>>
144 where
145 P: DBProvider + NodePrimitivesProvider + StorageSettingsCache + StaticFileProviderFactory,
146 P::Tx: DbTxMut,
147 {
148 if provider.cached_storage_settings().storage_v2 {
149 Ok(EitherWriter::StaticFile(
150 provider
151 .get_static_file_writer(block_number, StaticFileSegment::AccountChangeSets)?,
152 ))
153 } else {
154 Ok(EitherWriter::Database(
155 provider.tx_ref().cursor_dup_write::<tables::AccountChangeSets>()?,
156 ))
157 }
158 }
159
160 pub fn new_storage_changesets<P>(
162 provider: &'a P,
163 block_number: BlockNumber,
164 ) -> ProviderResult<DupEitherWriterTy<'a, P, tables::StorageChangeSets>>
165 where
166 P: DBProvider + NodePrimitivesProvider + StorageSettingsCache + StaticFileProviderFactory,
167 P::Tx: DbTxMut,
168 {
169 if provider.cached_storage_settings().storage_v2 {
170 Ok(EitherWriter::StaticFile(
171 provider
172 .get_static_file_writer(block_number, StaticFileSegment::StorageChangeSets)?,
173 ))
174 } else {
175 Ok(EitherWriter::Database(
176 provider.tx_ref().cursor_dup_write::<tables::StorageChangeSets>()?,
177 ))
178 }
179 }
180
181 pub fn receipts_destination<P: DBProvider + StorageSettingsCache>(
190 provider: &P,
191 ) -> EitherWriterDestination {
192 let receipts_in_static_files = provider.cached_storage_settings().storage_v2;
193 let prune_modes = provider.prune_modes_ref();
194
195 if !receipts_in_static_files && prune_modes.has_receipts_pruning() ||
196 receipts_in_static_files && !prune_modes.receipts_log_filter.is_empty()
198 {
199 EitherWriterDestination::Database
200 } else {
201 EitherWriterDestination::StaticFile
202 }
203 }
204
205 pub fn account_changesets_destination<P: DBProvider + StorageSettingsCache>(
209 provider: &P,
210 ) -> EitherWriterDestination {
211 if provider.cached_storage_settings().storage_v2 {
212 EitherWriterDestination::StaticFile
213 } else {
214 EitherWriterDestination::Database
215 }
216 }
217
218 pub fn storage_changesets_destination<P: DBProvider + StorageSettingsCache>(
222 provider: &P,
223 ) -> EitherWriterDestination {
224 if provider.cached_storage_settings().storage_v2 {
225 EitherWriterDestination::StaticFile
226 } else {
227 EitherWriterDestination::Database
228 }
229 }
230
231 pub fn new_storages_history<P>(
233 provider: &P,
234 _rocksdb_batch: RocksBatchArg<'a>,
235 ) -> ProviderResult<EitherWriterTy<'a, P, tables::StoragesHistory>>
236 where
237 P: DBProvider + NodePrimitivesProvider + StorageSettingsCache,
238 P::Tx: DbTxMut,
239 {
240 if provider.cached_storage_settings().storage_v2 {
241 return Ok(EitherWriter::RocksDB(_rocksdb_batch));
242 }
243
244 Ok(EitherWriter::Database(provider.tx_ref().cursor_write::<tables::StoragesHistory>()?))
245 }
246
247 pub fn new_transaction_hash_numbers<P>(
249 provider: &P,
250 _rocksdb_batch: RocksBatchArg<'a>,
251 ) -> ProviderResult<EitherWriterTy<'a, P, tables::TransactionHashNumbers>>
252 where
253 P: DBProvider + NodePrimitivesProvider + StorageSettingsCache,
254 P::Tx: DbTxMut,
255 {
256 if provider.cached_storage_settings().storage_v2 {
257 return Ok(EitherWriter::RocksDB(_rocksdb_batch));
258 }
259
260 Ok(EitherWriter::Database(
261 provider.tx_ref().cursor_write::<tables::TransactionHashNumbers>()?,
262 ))
263 }
264
265 pub fn new_accounts_history<P>(
267 provider: &P,
268 _rocksdb_batch: RocksBatchArg<'a>,
269 ) -> ProviderResult<EitherWriterTy<'a, P, tables::AccountsHistory>>
270 where
271 P: DBProvider + NodePrimitivesProvider + StorageSettingsCache,
272 P::Tx: DbTxMut,
273 {
274 if provider.cached_storage_settings().storage_v2 {
275 return Ok(EitherWriter::RocksDB(_rocksdb_batch));
276 }
277
278 Ok(EitherWriter::Database(provider.tx_ref().cursor_write::<tables::AccountsHistory>()?))
279 }
280}
281
282impl<'a, CURSOR, N: NodePrimitives> EitherWriter<'a, CURSOR, N> {
283 pub fn into_raw_rocksdb_batch(self) -> Option<rocksdb::WriteBatchWithTransaction<true>> {
291 match self {
292 Self::Database(_) | Self::StaticFile(_) => None,
293 Self::RocksDB(batch) => Some(batch.into_inner()),
294 }
295 }
296
297 pub fn increment_block(&mut self, expected_block_number: BlockNumber) -> ProviderResult<()> {
301 match self {
302 Self::Database(_) => Ok(()),
303 Self::StaticFile(writer) => writer.increment_block(expected_block_number),
304 Self::RocksDB(_) => Err(ProviderError::UnsupportedProvider),
305 }
306 }
307
308 pub fn ensure_at_block(&mut self, block_number: BlockNumber) -> ProviderResult<()> {
315 match self {
316 Self::Database(_) => Ok(()),
317 Self::StaticFile(writer) => writer.ensure_at_block(block_number),
318 Self::RocksDB(_) => Err(ProviderError::UnsupportedProvider),
319 }
320 }
321}
322
323impl<'a, CURSOR, N: NodePrimitives> EitherWriter<'a, CURSOR, N>
324where
325 N::Receipt: Value,
326 CURSOR: DbCursorRW<tables::Receipts<N::Receipt>>,
327{
328 pub fn append_receipt(&mut self, tx_num: TxNumber, receipt: &N::Receipt) -> ProviderResult<()> {
330 match self {
331 Self::Database(cursor) => Ok(cursor.append(tx_num, receipt)?),
332 Self::StaticFile(writer) => writer.append_receipt(tx_num, receipt),
333 Self::RocksDB(_) => Err(ProviderError::UnsupportedProvider),
334 }
335 }
336}
337
338impl<'a, CURSOR, N: NodePrimitives> EitherWriter<'a, CURSOR, N>
339where
340 CURSOR: DbCursorRW<tables::TransactionSenders>,
341{
342 pub fn append_sender(&mut self, tx_num: TxNumber, sender: &Address) -> ProviderResult<()> {
344 match self {
345 Self::Database(cursor) => Ok(cursor.append(tx_num, sender)?),
346 Self::StaticFile(writer) => writer.append_transaction_sender(tx_num, sender),
347 Self::RocksDB(_) => Err(ProviderError::UnsupportedProvider),
348 }
349 }
350
351 pub fn append_senders<I>(&mut self, senders: I) -> ProviderResult<()>
353 where
354 I: Iterator<Item = (TxNumber, Address)>,
355 {
356 match self {
357 Self::Database(cursor) => {
358 for (tx_num, sender) in senders {
359 cursor.append(tx_num, &sender)?;
360 }
361 Ok(())
362 }
363 Self::StaticFile(writer) => writer.append_transaction_senders(senders),
364 Self::RocksDB(_) => Err(ProviderError::UnsupportedProvider),
365 }
366 }
367
368 pub fn prune_senders(
371 &mut self,
372 unwind_tx_from: TxNumber,
373 block: BlockNumber,
374 ) -> ProviderResult<()>
375 where
376 CURSOR: DbCursorRO<tables::TransactionSenders>,
377 {
378 match self {
379 Self::Database(cursor) => {
380 let mut walker = cursor.walk_range(unwind_tx_from..)?;
381 while walker.next().transpose()?.is_some() {
382 walker.delete_current()?;
383 }
384 }
385 Self::StaticFile(writer) => {
386 let static_file_transaction_sender_num = writer
387 .reader()
388 .get_highest_static_file_tx(StaticFileSegment::TransactionSenders);
389
390 let to_delete = static_file_transaction_sender_num
391 .map(|static_num| (static_num + 1).saturating_sub(unwind_tx_from))
392 .unwrap_or_default();
393
394 writer.prune_transaction_senders(to_delete, block)?;
395 }
396 Self::RocksDB(_) => return Err(ProviderError::UnsupportedProvider),
397 }
398
399 Ok(())
400 }
401}
402
403impl<'a, CURSOR, N: NodePrimitives> EitherWriter<'a, CURSOR, N>
404where
405 CURSOR: DbCursorRW<tables::TransactionHashNumbers> + DbCursorRO<tables::TransactionHashNumbers>,
406{
407 pub fn put_transaction_hash_number(
413 &mut self,
414 hash: TxHash,
415 tx_num: TxNumber,
416 append_only: bool,
417 ) -> ProviderResult<()> {
418 match self {
419 Self::Database(cursor) => {
420 if append_only {
421 Ok(cursor.append(hash, &tx_num)?)
422 } else {
423 Ok(cursor.upsert(hash, &tx_num)?)
424 }
425 }
426 Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
427 Self::RocksDB(batch) => batch.put::<tables::TransactionHashNumbers>(hash, &tx_num),
428 }
429 }
430
431 pub fn put_transaction_hash_numbers_batch(
440 &mut self,
441 entries: Vec<(TxHash, TxNumber)>,
442 append_only: bool,
443 ) -> ProviderResult<()> {
444 match self {
445 Self::Database(cursor) => {
446 for (hash, tx_num) in entries {
447 if append_only {
448 cursor.append(hash, &tx_num)?;
449 } else {
450 cursor.upsert(hash, &tx_num)?;
451 }
452 }
453 Ok(())
454 }
455 Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
456 Self::RocksDB(batch) => {
457 for (hash, tx_num) in entries {
458 batch.put::<tables::TransactionHashNumbers>(hash, &tx_num)?;
459 }
460 Ok(())
461 }
462 }
463 }
464
465 pub fn delete_transaction_hash_number(&mut self, hash: TxHash) -> ProviderResult<()> {
467 match self {
468 Self::Database(cursor) => {
469 if cursor.seek_exact(hash)?.is_some() {
470 cursor.delete_current()?;
471 }
472 Ok(())
473 }
474 Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
475 Self::RocksDB(batch) => batch.delete::<tables::TransactionHashNumbers>(hash),
476 }
477 }
478}
479
480impl<'a, CURSOR, N: NodePrimitives> EitherWriter<'a, CURSOR, N>
481where
482 CURSOR: DbCursorRW<tables::StoragesHistory> + DbCursorRO<tables::StoragesHistory>,
483{
484 pub fn put_storage_history(
486 &mut self,
487 key: StorageShardedKey,
488 value: &BlockNumberList,
489 ) -> ProviderResult<()> {
490 match self {
491 Self::Database(cursor) => Ok(cursor.upsert(key, value)?),
492 Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
493 Self::RocksDB(batch) => batch.put::<tables::StoragesHistory>(key, value),
494 }
495 }
496
497 pub fn delete_storage_history(&mut self, key: StorageShardedKey) -> ProviderResult<()> {
499 match self {
500 Self::Database(cursor) => {
501 if cursor.seek_exact(key)?.is_some() {
502 cursor.delete_current()?;
503 }
504 Ok(())
505 }
506 Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
507 Self::RocksDB(batch) => batch.delete::<tables::StoragesHistory>(key),
508 }
509 }
510
511 pub fn append_storage_history(
513 &mut self,
514 key: StorageShardedKey,
515 value: &BlockNumberList,
516 ) -> ProviderResult<()> {
517 match self {
518 Self::Database(cursor) => Ok(cursor.append(key, value)?),
519 Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
520 Self::RocksDB(batch) => batch.put::<tables::StoragesHistory>(key, value),
521 }
522 }
523
524 pub fn upsert_storage_history(
526 &mut self,
527 key: StorageShardedKey,
528 value: &BlockNumberList,
529 ) -> ProviderResult<()> {
530 match self {
531 Self::Database(cursor) => Ok(cursor.upsert(key, value)?),
532 Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
533 Self::RocksDB(batch) => batch.put::<tables::StoragesHistory>(key, value),
534 }
535 }
536
537 pub fn get_last_storage_history_shard(
539 &mut self,
540 address: Address,
541 storage_key: B256,
542 ) -> ProviderResult<Option<BlockNumberList>> {
543 let key = StorageShardedKey::last(address, storage_key);
544 match self {
545 Self::Database(cursor) => Ok(cursor.seek_exact(key)?.map(|(_, v)| v)),
546 Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
547 Self::RocksDB(batch) => batch.get::<tables::StoragesHistory>(key),
548 }
549 }
550}
551
552impl<'a, CURSOR, N: NodePrimitives> EitherWriter<'a, CURSOR, N>
553where
554 CURSOR: DbCursorRW<tables::AccountsHistory> + DbCursorRO<tables::AccountsHistory>,
555{
556 pub fn append_account_history(
558 &mut self,
559 key: ShardedKey<Address>,
560 value: &BlockNumberList,
561 ) -> ProviderResult<()> {
562 match self {
563 Self::Database(cursor) => Ok(cursor.append(key, value)?),
564 Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
565 Self::RocksDB(batch) => batch.put::<tables::AccountsHistory>(key, value),
566 }
567 }
568
569 pub fn upsert_account_history(
571 &mut self,
572 key: ShardedKey<Address>,
573 value: &BlockNumberList,
574 ) -> ProviderResult<()> {
575 match self {
576 Self::Database(cursor) => Ok(cursor.upsert(key, value)?),
577 Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
578 Self::RocksDB(batch) => batch.put::<tables::AccountsHistory>(key, value),
579 }
580 }
581
582 pub fn get_last_account_history_shard(
584 &mut self,
585 address: Address,
586 ) -> ProviderResult<Option<BlockNumberList>> {
587 match self {
588 Self::Database(cursor) => {
589 Ok(cursor.seek_exact(ShardedKey::last(address))?.map(|(_, v)| v))
590 }
591 Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
592 Self::RocksDB(batch) => batch.get::<tables::AccountsHistory>(ShardedKey::last(address)),
593 }
594 }
595
596 pub fn delete_account_history(&mut self, key: ShardedKey<Address>) -> ProviderResult<()> {
598 match self {
599 Self::Database(cursor) => {
600 if cursor.seek_exact(key)?.is_some() {
601 cursor.delete_current()?;
602 }
603 Ok(())
604 }
605 Self::StaticFile(_) => Err(ProviderError::UnsupportedProvider),
606 Self::RocksDB(batch) => batch.delete::<tables::AccountsHistory>(key),
607 }
608 }
609}
610
611impl<'a, CURSOR, N: NodePrimitives> EitherWriter<'a, CURSOR, N>
612where
613 CURSOR: DbDupCursorRW<tables::AccountChangeSets>,
614{
615 pub fn append_account_changeset(
619 &mut self,
620 block_number: BlockNumber,
621 mut changeset: Vec<AccountBeforeTx>,
622 ) -> ProviderResult<()> {
623 changeset.par_sort_by_key(|a| a.address);
625 match self {
626 Self::Database(cursor) => {
627 for change in changeset {
628 cursor.append_dup(block_number, change)?;
629 }
630 }
631 Self::StaticFile(writer) => {
632 writer.append_account_changeset(changeset, block_number)?;
633 }
634 Self::RocksDB(_) => return Err(ProviderError::UnsupportedProvider),
635 }
636
637 Ok(())
638 }
639}
640
641impl<'a, CURSOR, N: NodePrimitives> EitherWriter<'a, CURSOR, N>
642where
643 CURSOR: DbDupCursorRW<tables::StorageChangeSets>,
644{
645 pub fn append_storage_changeset(
649 &mut self,
650 block_number: BlockNumber,
651 mut changeset: Vec<StorageBeforeTx>,
652 ) -> ProviderResult<()> {
653 changeset.par_sort_by_key(|change| (change.address, change.key));
654
655 match self {
656 Self::Database(cursor) => {
657 for change in changeset {
658 let storage_id = BlockNumberAddress((block_number, change.address));
659 cursor.append_dup(
660 storage_id,
661 StorageEntry { key: change.key, value: change.value },
662 )?;
663 }
664 }
665 Self::StaticFile(writer) => {
666 writer.append_storage_changeset(changeset, block_number)?;
667 }
668 Self::RocksDB(_) => return Err(ProviderError::UnsupportedProvider),
669 }
670
671 Ok(())
672 }
673}
674
675#[derive(Debug, Display)]
677pub enum EitherReader<'a, 'db, CURSOR, N> {
678 Database(CURSOR, PhantomData<&'a ()>),
680 StaticFile(StaticFileProvider<N>, PhantomData<&'a ()>),
682 RocksDB(&'a crate::providers::rocksdb::RocksReadSnapshot<'db>),
684}
685
686impl<'a, 'db> EitherReader<'a, 'db, (), ()> {
687 pub fn new_senders<P>(
689 provider: &P,
690 ) -> ProviderResult<EitherReaderTy<'a, 'db, P, tables::TransactionSenders>>
691 where
692 P: DBProvider + NodePrimitivesProvider + StorageSettingsCache + StaticFileProviderFactory,
693 P::Tx: DbTx,
694 {
695 if EitherWriterDestination::senders(provider).is_static_file() {
696 Ok(EitherReader::StaticFile(provider.static_file_provider(), PhantomData))
697 } else {
698 Ok(EitherReader::Database(
699 provider.tx_ref().cursor_read::<tables::TransactionSenders>()?,
700 PhantomData,
701 ))
702 }
703 }
704
705 pub fn new_storages_history<P>(
707 provider: &P,
708 rocksdb: RocksDBRefArg<'a, 'db>,
709 ) -> ProviderResult<EitherReaderTy<'a, 'db, P, tables::StoragesHistory>>
710 where
711 P: DBProvider + NodePrimitivesProvider + StorageSettingsCache,
712 P::Tx: DbTx,
713 {
714 if provider.cached_storage_settings().storage_v2 {
715 return Ok(EitherReader::RocksDB(
716 rocksdb.expect("storages_history_in_rocksdb requires rocksdb snapshot"),
717 ));
718 }
719
720 Ok(EitherReader::Database(
721 provider.tx_ref().cursor_read::<tables::StoragesHistory>()?,
722 PhantomData,
723 ))
724 }
725
726 pub fn new_transaction_hash_numbers<P>(
728 provider: &P,
729 rocksdb: RocksDBRefArg<'a, 'db>,
730 ) -> ProviderResult<EitherReaderTy<'a, 'db, P, tables::TransactionHashNumbers>>
731 where
732 P: DBProvider + NodePrimitivesProvider + StorageSettingsCache,
733 P::Tx: DbTx,
734 {
735 if provider.cached_storage_settings().storage_v2 {
736 return Ok(EitherReader::RocksDB(
737 rocksdb.expect("transaction_hash_numbers_in_rocksdb requires rocksdb snapshot"),
738 ));
739 }
740
741 Ok(EitherReader::Database(
742 provider.tx_ref().cursor_read::<tables::TransactionHashNumbers>()?,
743 PhantomData,
744 ))
745 }
746
747 pub fn new_accounts_history<P>(
749 provider: &P,
750 rocksdb: RocksDBRefArg<'a, 'db>,
751 ) -> ProviderResult<EitherReaderTy<'a, 'db, P, tables::AccountsHistory>>
752 where
753 P: DBProvider + NodePrimitivesProvider + StorageSettingsCache,
754 P::Tx: DbTx,
755 {
756 if provider.cached_storage_settings().storage_v2 {
757 return Ok(EitherReader::RocksDB(
758 rocksdb.expect("account_history_in_rocksdb requires rocksdb snapshot"),
759 ));
760 }
761
762 Ok(EitherReader::Database(
763 provider.tx_ref().cursor_read::<tables::AccountsHistory>()?,
764 PhantomData,
765 ))
766 }
767
768 pub fn new_account_changesets<P>(
770 provider: &P,
771 ) -> ProviderResult<DupEitherReaderTy<'a, 'db, P, tables::AccountChangeSets>>
772 where
773 P: DBProvider + NodePrimitivesProvider + StorageSettingsCache + StaticFileProviderFactory,
774 P::Tx: DbTx,
775 {
776 if EitherWriterDestination::account_changesets(provider).is_static_file() {
777 Ok(EitherReader::StaticFile(provider.static_file_provider(), PhantomData))
778 } else {
779 Ok(EitherReader::Database(
780 provider.tx_ref().cursor_dup_read::<tables::AccountChangeSets>()?,
781 PhantomData,
782 ))
783 }
784 }
785}
786
787impl<CURSOR, N: NodePrimitives> EitherReader<'_, '_, CURSOR, N>
788where
789 CURSOR: DbCursorRO<tables::TransactionSenders>,
790{
791 pub fn senders_by_tx_range(
793 &mut self,
794 range: Range<TxNumber>,
795 ) -> ProviderResult<HashMap<TxNumber, Address>> {
796 match self {
797 Self::Database(cursor, _) => cursor
798 .walk_range(range)?
799 .map(|result| result.map_err(ProviderError::from))
800 .collect::<ProviderResult<HashMap<_, _>>>(),
801 Self::StaticFile(provider, _) => range
802 .clone()
803 .zip(provider.fetch_range_iter(
804 StaticFileSegment::TransactionSenders,
805 range,
806 |cursor, number| cursor.get_one::<TransactionSenderMask>(number.into()),
807 )?)
808 .filter_map(|(tx_num, sender)| {
809 let result = sender.transpose()?;
810 Some(result.map(|sender| (tx_num, sender)))
811 })
812 .collect::<ProviderResult<HashMap<_, _>>>(),
813 Self::RocksDB(_) => Err(ProviderError::UnsupportedProvider),
814 }
815 }
816}
817
818impl<CURSOR, N: NodePrimitives> EitherReader<'_, '_, CURSOR, N>
819where
820 CURSOR: DbCursorRO<tables::TransactionHashNumbers>,
821{
822 pub fn get_transaction_hash_number(
824 &mut self,
825 hash: TxHash,
826 ) -> ProviderResult<Option<TxNumber>> {
827 match self {
828 Self::Database(cursor, _) => Ok(cursor.seek_exact(hash)?.map(|(_, v)| v)),
829 Self::StaticFile(_, _) => Err(ProviderError::UnsupportedProvider),
830 Self::RocksDB(snapshot) => snapshot.get::<tables::TransactionHashNumbers>(hash),
831 }
832 }
833}
834
835impl<CURSOR, N: NodePrimitives> EitherReader<'_, '_, CURSOR, N>
836where
837 CURSOR: DbCursorRO<tables::StoragesHistory>,
838{
839 pub fn get_storage_history(
841 &mut self,
842 key: StorageShardedKey,
843 ) -> ProviderResult<Option<BlockNumberList>> {
844 match self {
845 Self::Database(cursor, _) => Ok(cursor.seek_exact(key)?.map(|(_, v)| v)),
846 Self::StaticFile(_, _) => Err(ProviderError::UnsupportedProvider),
847 Self::RocksDB(snapshot) => snapshot.get::<tables::StoragesHistory>(key),
848 }
849 }
850
851 pub fn storage_history_info(
853 &mut self,
854 address: Address,
855 storage_key: alloy_primitives::B256,
856 block_number: BlockNumber,
857 lowest_available_block_number: Option<BlockNumber>,
858 visible_tip: BlockNumber,
859 ) -> ProviderResult<HistoryInfo> {
860 match self {
861 Self::Database(cursor, _) => {
862 let key = StorageShardedKey::new(address, storage_key, block_number);
863 history_info::<tables::StoragesHistory, _, _>(
864 cursor,
865 key,
866 block_number,
867 |k| k.address == address && k.sharded_key.key == storage_key,
868 lowest_available_block_number,
869 )
870 }
871 Self::StaticFile(_, _) => Err(ProviderError::UnsupportedProvider),
872 Self::RocksDB(snapshot) => snapshot.storage_history_info(
873 address,
874 storage_key,
875 block_number,
876 lowest_available_block_number,
877 visible_tip,
878 ),
879 }
880 }
881}
882
883impl<CURSOR, N: NodePrimitives> EitherReader<'_, '_, CURSOR, N>
884where
885 CURSOR: DbCursorRO<tables::AccountsHistory>,
886{
887 pub fn get_account_history(
889 &mut self,
890 key: ShardedKey<Address>,
891 ) -> ProviderResult<Option<BlockNumberList>> {
892 match self {
893 Self::Database(cursor, _) => Ok(cursor.seek_exact(key)?.map(|(_, v)| v)),
894 Self::StaticFile(_, _) => Err(ProviderError::UnsupportedProvider),
895 Self::RocksDB(snapshot) => snapshot.get::<tables::AccountsHistory>(key),
896 }
897 }
898
899 pub fn account_history_info(
901 &mut self,
902 address: Address,
903 block_number: BlockNumber,
904 lowest_available_block_number: Option<BlockNumber>,
905 visible_tip: BlockNumber,
906 ) -> ProviderResult<HistoryInfo> {
907 match self {
908 Self::Database(cursor, _) => {
909 let key = ShardedKey::new(address, block_number);
910 history_info::<tables::AccountsHistory, _, _>(
911 cursor,
912 key,
913 block_number,
914 |k| k.key == address,
915 lowest_available_block_number,
916 )
917 }
918 Self::StaticFile(_, _) => Err(ProviderError::UnsupportedProvider),
919 Self::RocksDB(snapshot) => snapshot.account_history_info(
920 address,
921 block_number,
922 lowest_available_block_number,
923 visible_tip,
924 ),
925 }
926 }
927}
928
929impl<CURSOR, N: NodePrimitives> EitherReader<'_, '_, CURSOR, N>
930where
931 CURSOR: DbCursorRO<tables::AccountChangeSets>,
932{
933 pub fn changed_accounts_with_range(
935 &mut self,
936 range: RangeInclusive<BlockNumber>,
937 ) -> ProviderResult<BTreeSet<Address>> {
938 match self {
939 Self::StaticFile(provider, _) => {
940 let highest_static_block =
941 provider.get_highest_static_file_block(StaticFileSegment::AccountChangeSets);
942
943 let Some(highest) = highest_static_block else {
944 return Err(ProviderError::MissingHighestStaticFileBlock(
945 StaticFileSegment::AccountChangeSets,
946 ))
947 };
948
949 let start = *range.start();
950 let static_end = (*range.end()).min(highest);
951
952 let mut changed_accounts = BTreeSet::default();
953 if start <= static_end {
954 for block in start..=static_end {
955 let block_changesets = provider.account_block_changeset(block)?;
956 for changeset in block_changesets {
957 changed_accounts.insert(changeset.address);
958 }
959 }
960 }
961
962 Ok(changed_accounts)
963 }
964 Self::Database(provider, _) => provider
965 .walk_range(range)?
966 .map(|entry| {
967 entry.map(|(_, account_before)| account_before.address).map_err(Into::into)
968 })
969 .collect(),
970 Self::RocksDB(_) => Err(ProviderError::UnsupportedProvider),
971 }
972 }
973}
974
975#[derive(Debug, EnumIs)]
977pub enum EitherWriterDestination {
978 Database,
980 StaticFile,
982 RocksDB,
984}
985
986impl EitherWriterDestination {
987 pub fn senders<P>(provider: &P) -> Self
989 where
990 P: StorageSettingsCache,
991 {
992 if provider.cached_storage_settings().storage_v2 {
994 Self::StaticFile
995 } else {
996 Self::Database
997 }
998 }
999
1000 pub fn account_changesets<P>(provider: &P) -> Self
1002 where
1003 P: StorageSettingsCache,
1004 {
1005 if provider.cached_storage_settings().storage_v2 {
1007 Self::StaticFile
1008 } else {
1009 Self::Database
1010 }
1011 }
1012
1013 pub fn storage_changesets<P>(provider: &P) -> Self
1015 where
1016 P: StorageSettingsCache,
1017 {
1018 if provider.cached_storage_settings().storage_v2 {
1020 Self::StaticFile
1021 } else {
1022 Self::Database
1023 }
1024 }
1025}
1026
1027#[cfg(test)]
1028mod tests {
1029 use crate::{test_utils::create_test_provider_factory, StaticFileWriter};
1030
1031 use super::*;
1032 use alloy_primitives::Address;
1033 use reth_db::models::AccountBeforeTx;
1034 use reth_static_file_types::StaticFileSegment;
1035 use reth_storage_api::{DatabaseProviderFactory, StorageSettings};
1036
1037 #[test]
1049 fn test_changed_accounts_with_range_caps_at_static_file_tip() {
1050 let factory = create_test_provider_factory();
1051 let highest_block = 5u64;
1052
1053 let addresses: Vec<Address> = (0..=highest_block)
1054 .map(|i| {
1055 let mut addr = Address::ZERO;
1056 addr.0[0] = i as u8;
1057 addr
1058 })
1059 .collect();
1060
1061 {
1062 let sf_provider = factory.static_file_provider();
1063 let mut writer =
1064 sf_provider.latest_writer(StaticFileSegment::AccountChangeSets).unwrap();
1065
1066 for block_num in 0..=highest_block {
1067 let changeset =
1068 vec![AccountBeforeTx { address: addresses[block_num as usize], info: None }];
1069 writer.append_account_changeset(changeset, block_num).unwrap();
1070 }
1071 writer.commit().unwrap();
1072 }
1073
1074 factory.set_storage_settings_cache(StorageSettings::v2());
1075
1076 let provider = factory.database_provider_ro().unwrap();
1077
1078 let sf_tip = provider
1079 .static_file_provider()
1080 .get_highest_static_file_block(StaticFileSegment::AccountChangeSets);
1081 assert_eq!(sf_tip, Some(highest_block));
1082
1083 let mut reader = EitherReader::new_account_changesets(&provider).unwrap();
1084 assert!(matches!(reader, EitherReader::StaticFile(_, _)));
1085
1086 let result = reader.changed_accounts_with_range(0..=10).unwrap();
1088
1089 let expected: BTreeSet<Address> = addresses.into_iter().collect();
1090 assert_eq!(result, expected);
1091 }
1092
1093 #[test]
1094 fn test_reader_senders_by_tx_range() {
1095 let factory = create_test_provider_factory();
1096
1097 let senders = [
1099 (1, Address::random()),
1100 (2, Address::random()),
1101 (3, Address::random()),
1102 (4, Address::random()),
1103 ];
1104
1105 for transaction_senders_in_static_files in [false, true] {
1106 factory.set_storage_settings_cache(if transaction_senders_in_static_files {
1107 StorageSettings::v2()
1108 } else {
1109 StorageSettings::v1()
1110 });
1111
1112 let provider = factory.database_provider_rw().unwrap();
1113 let mut writer = EitherWriter::new_senders(&provider, 0).unwrap();
1114 if transaction_senders_in_static_files {
1115 assert!(matches!(writer, EitherWriter::StaticFile(_)));
1116 } else {
1117 assert!(matches!(writer, EitherWriter::Database(_)));
1118 }
1119
1120 writer.increment_block(0).unwrap();
1121 writer.append_senders(senders.iter().copied()).unwrap();
1122 drop(writer);
1123 provider.commit().unwrap();
1124
1125 let provider = factory.database_provider_ro().unwrap();
1126 let mut reader = EitherReader::new_senders(&provider).unwrap();
1127 if transaction_senders_in_static_files {
1128 assert!(matches!(reader, EitherReader::StaticFile(_, _)));
1129 } else {
1130 assert!(matches!(reader, EitherReader::Database(_, _)));
1131 }
1132
1133 assert_eq!(
1134 reader.senders_by_tx_range(0..6).unwrap(),
1135 senders.iter().copied().collect::<HashMap<_, _>>(),
1136 "{reader}"
1137 );
1138 }
1139 }
1140}
1141
1142#[cfg(test)]
1143mod rocksdb_tests {
1144 use super::*;
1145 use crate::{
1146 providers::rocksdb::{RocksDBBuilder, RocksDBProvider},
1147 test_utils::create_test_provider_factory,
1148 RocksDBProviderFactory,
1149 };
1150 use alloy_primitives::{Address, B256};
1151 use reth_db_api::{
1152 models::{storage_sharded_key::StorageShardedKey, IntegerList, ShardedKey},
1153 tables,
1154 transaction::DbTxMut,
1155 };
1156 use reth_ethereum_primitives::EthPrimitives;
1157 use reth_storage_api::{DatabaseProviderFactory, StorageSettings};
1158 use std::marker::PhantomData;
1159 use tempfile::TempDir;
1160
1161 fn create_rocksdb_provider() -> (TempDir, RocksDBProvider) {
1162 let temp_dir = TempDir::new().unwrap();
1163 let provider = RocksDBBuilder::new(temp_dir.path())
1164 .with_table::<tables::TransactionHashNumbers>()
1165 .with_table::<tables::StoragesHistory>()
1166 .with_table::<tables::AccountsHistory>()
1167 .build()
1168 .unwrap();
1169 (temp_dir, provider)
1170 }
1171
1172 #[test]
1176 fn test_either_writer_transaction_hash_numbers_with_rocksdb() {
1177 let factory = create_test_provider_factory();
1178
1179 factory.set_storage_settings_cache(StorageSettings::v2());
1181
1182 let hash1 = B256::from([1u8; 32]);
1183 let hash2 = B256::from([2u8; 32]);
1184 let tx_num1 = 100u64;
1185 let tx_num2 = 200u64;
1186
1187 let rocksdb = factory.rocksdb_provider();
1189 let batch = rocksdb.batch();
1190
1191 let provider = factory.database_provider_rw().unwrap();
1193 let mut writer = EitherWriter::new_transaction_hash_numbers(&provider, batch).unwrap();
1194
1195 assert!(matches!(writer, EitherWriter::RocksDB(_)));
1197
1198 writer.put_transaction_hash_number(hash1, tx_num1, false).unwrap();
1200 writer.put_transaction_hash_number(hash2, tx_num2, false).unwrap();
1201
1202 if let Some(batch) = writer.into_raw_rocksdb_batch() {
1204 provider.set_pending_rocksdb_batch(batch);
1205 }
1206
1207 provider.commit().unwrap();
1209
1210 let rocksdb = factory.rocksdb_provider();
1212 assert_eq!(rocksdb.get::<tables::TransactionHashNumbers>(hash1).unwrap(), Some(tx_num1));
1213 assert_eq!(rocksdb.get::<tables::TransactionHashNumbers>(hash2).unwrap(), Some(tx_num2));
1214 }
1215
1216 #[test]
1218 fn test_either_writer_delete_transaction_hash_number_with_rocksdb() {
1219 let factory = create_test_provider_factory();
1220
1221 factory.set_storage_settings_cache(StorageSettings::v2());
1223
1224 let hash = B256::from([1u8; 32]);
1225 let tx_num = 100u64;
1226
1227 let rocksdb = factory.rocksdb_provider();
1229 rocksdb.put::<tables::TransactionHashNumbers>(hash, &tx_num).unwrap();
1230 assert_eq!(rocksdb.get::<tables::TransactionHashNumbers>(hash).unwrap(), Some(tx_num));
1231
1232 let batch = rocksdb.batch();
1234 let provider = factory.database_provider_rw().unwrap();
1235 let mut writer = EitherWriter::new_transaction_hash_numbers(&provider, batch).unwrap();
1236 writer.delete_transaction_hash_number(hash).unwrap();
1237
1238 if let Some(batch) = writer.into_raw_rocksdb_batch() {
1240 provider.set_pending_rocksdb_batch(batch);
1241 }
1242 provider.commit().unwrap();
1243
1244 let rocksdb = factory.rocksdb_provider();
1246 assert_eq!(rocksdb.get::<tables::TransactionHashNumbers>(hash).unwrap(), None);
1247 }
1248
1249 #[test]
1250 fn test_rocksdb_batch_transaction_hash_numbers() {
1251 let (_temp_dir, provider) = create_rocksdb_provider();
1252
1253 let hash1 = B256::from([1u8; 32]);
1254 let hash2 = B256::from([2u8; 32]);
1255 let tx_num1 = 100u64;
1256 let tx_num2 = 200u64;
1257
1258 let mut batch = provider.batch();
1260 batch.put::<tables::TransactionHashNumbers>(hash1, &tx_num1).unwrap();
1261 batch.put::<tables::TransactionHashNumbers>(hash2, &tx_num2).unwrap();
1262 batch.commit().unwrap();
1263
1264 let tx = provider.tx();
1266 assert_eq!(tx.get::<tables::TransactionHashNumbers>(hash1).unwrap(), Some(tx_num1));
1267 assert_eq!(tx.get::<tables::TransactionHashNumbers>(hash2).unwrap(), Some(tx_num2));
1268
1269 let missing_hash = B256::from([99u8; 32]);
1271 assert_eq!(tx.get::<tables::TransactionHashNumbers>(missing_hash).unwrap(), None);
1272 }
1273
1274 #[test]
1275 fn test_rocksdb_batch_storage_history() {
1276 let (_temp_dir, provider) = create_rocksdb_provider();
1277
1278 let address = Address::random();
1279 let storage_key = B256::from([1u8; 32]);
1280 let key = StorageShardedKey::new(address, storage_key, 1000);
1281 let value = IntegerList::new([1, 5, 10, 50]).unwrap();
1282
1283 let mut batch = provider.batch();
1285 batch.put::<tables::StoragesHistory>(key.clone(), &value).unwrap();
1286 batch.commit().unwrap();
1287
1288 let tx = provider.tx();
1290 let result = tx.get::<tables::StoragesHistory>(key).unwrap();
1291 assert_eq!(result, Some(value));
1292
1293 let missing_key = StorageShardedKey::new(Address::random(), B256::ZERO, 0);
1295 assert_eq!(tx.get::<tables::StoragesHistory>(missing_key).unwrap(), None);
1296 }
1297
1298 #[test]
1299 fn test_rocksdb_batch_account_history() {
1300 let (_temp_dir, provider) = create_rocksdb_provider();
1301
1302 let address = Address::random();
1303 let key = ShardedKey::new(address, 1000);
1304 let value = IntegerList::new([1, 10, 100, 500]).unwrap();
1305
1306 let mut batch = provider.batch();
1308 batch.put::<tables::AccountsHistory>(key.clone(), &value).unwrap();
1309 batch.commit().unwrap();
1310
1311 let tx = provider.tx();
1313 let result = tx.get::<tables::AccountsHistory>(key).unwrap();
1314 assert_eq!(result, Some(value));
1315
1316 let missing_key = ShardedKey::new(Address::random(), 0);
1318 assert_eq!(tx.get::<tables::AccountsHistory>(missing_key).unwrap(), None);
1319 }
1320
1321 #[test]
1322 fn test_rocksdb_batch_delete_transaction_hash_number() {
1323 let (_temp_dir, provider) = create_rocksdb_provider();
1324
1325 let hash = B256::from([1u8; 32]);
1326 let tx_num = 100u64;
1327
1328 provider.put::<tables::TransactionHashNumbers>(hash, &tx_num).unwrap();
1330 assert_eq!(provider.get::<tables::TransactionHashNumbers>(hash).unwrap(), Some(tx_num));
1331
1332 let mut batch = provider.batch();
1334 batch.delete::<tables::TransactionHashNumbers>(hash).unwrap();
1335 batch.commit().unwrap();
1336
1337 assert_eq!(provider.get::<tables::TransactionHashNumbers>(hash).unwrap(), None);
1339 }
1340
1341 #[test]
1342 fn test_rocksdb_batch_delete_storage_history() {
1343 let (_temp_dir, provider) = create_rocksdb_provider();
1344
1345 let address = Address::random();
1346 let storage_key = B256::from([1u8; 32]);
1347 let key = StorageShardedKey::new(address, storage_key, 1000);
1348 let value = IntegerList::new([1, 5, 10]).unwrap();
1349
1350 provider.put::<tables::StoragesHistory>(key.clone(), &value).unwrap();
1352 assert!(provider.get::<tables::StoragesHistory>(key.clone()).unwrap().is_some());
1353
1354 let mut batch = provider.batch();
1356 batch.delete::<tables::StoragesHistory>(key.clone()).unwrap();
1357 batch.commit().unwrap();
1358
1359 assert_eq!(provider.get::<tables::StoragesHistory>(key).unwrap(), None);
1361 }
1362
1363 #[test]
1364 fn test_rocksdb_batch_delete_account_history() {
1365 let (_temp_dir, provider) = create_rocksdb_provider();
1366
1367 let address = Address::random();
1368 let key = ShardedKey::new(address, 1000);
1369 let value = IntegerList::new([1, 10, 100]).unwrap();
1370
1371 provider.put::<tables::AccountsHistory>(key.clone(), &value).unwrap();
1373 assert!(provider.get::<tables::AccountsHistory>(key.clone()).unwrap().is_some());
1374
1375 let mut batch = provider.batch();
1377 batch.delete::<tables::AccountsHistory>(key.clone()).unwrap();
1378 batch.commit().unwrap();
1379
1380 assert_eq!(provider.get::<tables::AccountsHistory>(key).unwrap(), None);
1382 }
1383
1384 struct HistoryQuery {
1391 block_number: BlockNumber,
1392 lowest_available: Option<BlockNumber>,
1393 expected: HistoryInfo,
1394 }
1395
1396 type AccountsHistoryWriteCursor =
1398 reth_db::mdbx::cursor::Cursor<reth_db::mdbx::RW, tables::AccountsHistory>;
1399 type StoragesHistoryWriteCursor =
1400 reth_db::mdbx::cursor::Cursor<reth_db::mdbx::RW, tables::StoragesHistory>;
1401 type AccountsHistoryReadCursor =
1402 reth_db::mdbx::cursor::Cursor<reth_db::mdbx::RO, tables::AccountsHistory>;
1403 type StoragesHistoryReadCursor =
1404 reth_db::mdbx::cursor::Cursor<reth_db::mdbx::RO, tables::StoragesHistory>;
1405
1406 fn run_account_history_scenario(
1409 scenario_name: &str,
1410 address: Address,
1411 shards: &[(BlockNumber, Vec<BlockNumber>)], queries: &[HistoryQuery],
1413 ) {
1414 let factory = create_test_provider_factory();
1416 let mdbx_provider = factory.database_provider_rw().unwrap();
1417 let (temp_dir, rocks_provider) = create_rocksdb_provider();
1418
1419 let mut mdbx_writer: EitherWriter<'_, AccountsHistoryWriteCursor, EthPrimitives> =
1421 EitherWriter::Database(
1422 mdbx_provider.tx_ref().cursor_write::<tables::AccountsHistory>().unwrap(),
1423 );
1424 let mut rocks_writer: EitherWriter<'_, AccountsHistoryWriteCursor, EthPrimitives> =
1425 EitherWriter::RocksDB(rocks_provider.batch());
1426
1427 for (highest_block, blocks) in shards {
1429 let key = ShardedKey::new(address, *highest_block);
1430 let value = IntegerList::new(blocks.clone()).unwrap();
1431 mdbx_writer.upsert_account_history(key.clone(), &value).unwrap();
1432 rocks_writer.upsert_account_history(key, &value).unwrap();
1433 }
1434
1435 drop(mdbx_writer);
1437 mdbx_provider.commit().unwrap();
1438 if let EitherWriter::RocksDB(batch) = rocks_writer {
1439 batch.commit().unwrap();
1440 }
1441
1442 let mdbx_ro = factory.database_provider_ro().unwrap();
1444 let rocks_snapshot = rocks_provider.snapshot();
1445
1446 for (i, query) in queries.iter().enumerate() {
1447 let mut mdbx_reader: EitherReader<'_, '_, AccountsHistoryReadCursor, EthPrimitives> =
1449 EitherReader::Database(
1450 mdbx_ro.tx_ref().cursor_read::<tables::AccountsHistory>().unwrap(),
1451 PhantomData,
1452 );
1453 let mdbx_result = mdbx_reader
1454 .account_history_info(address, query.block_number, query.lowest_available, u64::MAX)
1455 .unwrap();
1456
1457 let rocks_result = rocks_snapshot
1459 .account_history_info(address, query.block_number, query.lowest_available, u64::MAX)
1460 .unwrap();
1461
1462 assert_eq!(
1464 mdbx_result,
1465 rocks_result,
1466 "Backend mismatch in scenario '{}' query {}: block={}, lowest={:?}\n\
1467 MDBX: {:?}, RocksDB: {:?}",
1468 scenario_name,
1469 i,
1470 query.block_number,
1471 query.lowest_available,
1472 mdbx_result,
1473 rocks_result
1474 );
1475
1476 assert_eq!(
1478 mdbx_result,
1479 query.expected,
1480 "Unexpected result in scenario '{}' query {}: block={}, lowest={:?}\n\
1481 Got: {:?}, Expected: {:?}",
1482 scenario_name,
1483 i,
1484 query.block_number,
1485 query.lowest_available,
1486 mdbx_result,
1487 query.expected
1488 );
1489 }
1490
1491 drop(temp_dir);
1492 }
1493
1494 fn run_storage_history_scenario(
1497 scenario_name: &str,
1498 address: Address,
1499 storage_key: B256,
1500 shards: &[(BlockNumber, Vec<BlockNumber>)], queries: &[HistoryQuery],
1502 ) {
1503 let factory = create_test_provider_factory();
1505 let mdbx_provider = factory.database_provider_rw().unwrap();
1506 let (temp_dir, rocks_provider) = create_rocksdb_provider();
1507
1508 let mut mdbx_writer: EitherWriter<'_, StoragesHistoryWriteCursor, EthPrimitives> =
1510 EitherWriter::Database(
1511 mdbx_provider.tx_ref().cursor_write::<tables::StoragesHistory>().unwrap(),
1512 );
1513 let mut rocks_writer: EitherWriter<'_, StoragesHistoryWriteCursor, EthPrimitives> =
1514 EitherWriter::RocksDB(rocks_provider.batch());
1515
1516 for (highest_block, blocks) in shards {
1518 let key = StorageShardedKey::new(address, storage_key, *highest_block);
1519 let value = IntegerList::new(blocks.clone()).unwrap();
1520 mdbx_writer.put_storage_history(key.clone(), &value).unwrap();
1521 rocks_writer.put_storage_history(key, &value).unwrap();
1522 }
1523
1524 drop(mdbx_writer);
1526 mdbx_provider.commit().unwrap();
1527 if let EitherWriter::RocksDB(batch) = rocks_writer {
1528 batch.commit().unwrap();
1529 }
1530
1531 let mdbx_ro = factory.database_provider_ro().unwrap();
1533 let rocks_snapshot = rocks_provider.snapshot();
1534
1535 for (i, query) in queries.iter().enumerate() {
1536 let mut mdbx_reader: EitherReader<'_, '_, StoragesHistoryReadCursor, EthPrimitives> =
1538 EitherReader::Database(
1539 mdbx_ro.tx_ref().cursor_read::<tables::StoragesHistory>().unwrap(),
1540 PhantomData,
1541 );
1542 let mdbx_result = mdbx_reader
1543 .storage_history_info(
1544 address,
1545 storage_key,
1546 query.block_number,
1547 query.lowest_available,
1548 u64::MAX,
1549 )
1550 .unwrap();
1551
1552 let rocks_result = rocks_snapshot
1554 .storage_history_info(
1555 address,
1556 storage_key,
1557 query.block_number,
1558 query.lowest_available,
1559 u64::MAX,
1560 )
1561 .unwrap();
1562
1563 assert_eq!(
1565 mdbx_result,
1566 rocks_result,
1567 "Backend mismatch in scenario '{}' query {}: block={}, lowest={:?}\n\
1568 MDBX: {:?}, RocksDB: {:?}",
1569 scenario_name,
1570 i,
1571 query.block_number,
1572 query.lowest_available,
1573 mdbx_result,
1574 rocks_result
1575 );
1576
1577 assert_eq!(
1579 mdbx_result,
1580 query.expected,
1581 "Unexpected result in scenario '{}' query {}: block={}, lowest={:?}\n\
1582 Got: {:?}, Expected: {:?}",
1583 scenario_name,
1584 i,
1585 query.block_number,
1586 query.lowest_available,
1587 mdbx_result,
1588 query.expected
1589 );
1590 }
1591
1592 drop(temp_dir);
1593 }
1594
1595 #[test]
1603 fn test_account_history_info_both_backends() {
1604 let address = Address::from([0x42; 20]);
1605
1606 run_account_history_scenario(
1608 "single_shard",
1609 address,
1610 &[(u64::MAX, vec![100, 200, 300])],
1611 &[
1612 HistoryQuery {
1614 block_number: 50,
1615 lowest_available: None,
1616 expected: HistoryInfo::NotYetWritten,
1617 },
1618 HistoryQuery {
1620 block_number: 150,
1621 lowest_available: None,
1622 expected: HistoryInfo::InChangeset(200),
1623 },
1624 HistoryQuery {
1626 block_number: 300,
1627 lowest_available: None,
1628 expected: HistoryInfo::InChangeset(300),
1629 },
1630 HistoryQuery {
1632 block_number: 500,
1633 lowest_available: None,
1634 expected: HistoryInfo::InPlainState,
1635 },
1636 ],
1637 );
1638
1639 run_account_history_scenario(
1641 "multiple_shards",
1642 address,
1643 &[
1644 (500, vec![100, 200, 300, 400, 500]), (u64::MAX, vec![600, 700, 800]), ],
1647 &[
1648 HistoryQuery {
1650 block_number: 50,
1651 lowest_available: None,
1652 expected: HistoryInfo::NotYetWritten,
1653 },
1654 HistoryQuery {
1656 block_number: 150,
1657 lowest_available: None,
1658 expected: HistoryInfo::InChangeset(200),
1659 },
1660 HistoryQuery {
1662 block_number: 550,
1663 lowest_available: None,
1664 expected: HistoryInfo::InChangeset(600),
1665 },
1666 HistoryQuery {
1668 block_number: 900,
1669 lowest_available: None,
1670 expected: HistoryInfo::InPlainState,
1671 },
1672 ],
1673 );
1674
1675 let address_without_history = Address::from([0x43; 20]);
1677 run_account_history_scenario(
1678 "no_history",
1679 address_without_history,
1680 &[], &[HistoryQuery {
1682 block_number: 150,
1683 lowest_available: None,
1684 expected: HistoryInfo::NotYetWritten,
1685 }],
1686 );
1687
1688 run_account_history_scenario(
1694 "with_pruning_boundary",
1695 address,
1696 &[(u64::MAX, vec![100, 200, 300])],
1697 &[
1698 HistoryQuery {
1700 block_number: 100,
1701 lowest_available: Some(100),
1702 expected: HistoryInfo::InChangeset(100),
1703 },
1704 HistoryQuery {
1706 block_number: 150,
1707 lowest_available: Some(100),
1708 expected: HistoryInfo::InChangeset(200),
1709 },
1710 ],
1711 );
1712 }
1713
1714 #[test]
1716 fn test_storage_history_info_both_backends() {
1717 let address = Address::from([0x42; 20]);
1718 let storage_key = B256::from([0x01; 32]);
1719 let other_storage_key = B256::from([0x02; 32]);
1720
1721 run_storage_history_scenario(
1723 "storage_single_shard",
1724 address,
1725 storage_key,
1726 &[(u64::MAX, vec![100, 200, 300])],
1727 &[
1728 HistoryQuery {
1730 block_number: 50,
1731 lowest_available: None,
1732 expected: HistoryInfo::NotYetWritten,
1733 },
1734 HistoryQuery {
1736 block_number: 150,
1737 lowest_available: None,
1738 expected: HistoryInfo::InChangeset(200),
1739 },
1740 HistoryQuery {
1742 block_number: 500,
1743 lowest_available: None,
1744 expected: HistoryInfo::InPlainState,
1745 },
1746 ],
1747 );
1748
1749 run_storage_history_scenario(
1751 "storage_no_history",
1752 address,
1753 other_storage_key,
1754 &[], &[HistoryQuery {
1756 block_number: 150,
1757 lowest_available: None,
1758 expected: HistoryInfo::NotYetWritten,
1759 }],
1760 );
1761 }
1762
1763 #[test]
1766 fn test_rocksdb_commits_at_provider_level() {
1767 let factory = create_test_provider_factory();
1768
1769 factory.set_storage_settings_cache(StorageSettings::v2());
1771
1772 let hash1 = B256::from([1u8; 32]);
1773 let hash2 = B256::from([2u8; 32]);
1774 let tx_num1 = 100u64;
1775 let tx_num2 = 200u64;
1776
1777 let rocksdb = factory.rocksdb_provider();
1779 let batch = rocksdb.batch();
1780
1781 let provider = factory.database_provider_rw().unwrap();
1783 let mut writer = EitherWriter::new_transaction_hash_numbers(&provider, batch).unwrap();
1784
1785 writer.put_transaction_hash_number(hash1, tx_num1, false).unwrap();
1787 writer.put_transaction_hash_number(hash2, tx_num2, false).unwrap();
1788
1789 let raw_batch = writer.into_raw_rocksdb_batch();
1791 if let Some(batch) = raw_batch {
1792 provider.set_pending_rocksdb_batch(batch);
1793 }
1794
1795 let rocksdb = factory.rocksdb_provider();
1797 assert_eq!(
1798 rocksdb.get::<tables::TransactionHashNumbers>(hash1).unwrap(),
1799 None,
1800 "Data should not be visible before provider.commit()"
1801 );
1802
1803 provider.commit().unwrap();
1805
1806 let rocksdb = factory.rocksdb_provider();
1808 assert_eq!(
1809 rocksdb.get::<tables::TransactionHashNumbers>(hash1).unwrap(),
1810 Some(tx_num1),
1811 "Data should be visible after provider.commit()"
1812 );
1813 assert_eq!(
1814 rocksdb.get::<tables::TransactionHashNumbers>(hash2).unwrap(),
1815 Some(tx_num2),
1816 "Data should be visible after provider.commit()"
1817 );
1818 }
1819
1820 #[test]
1824 #[should_panic(expected = "account_history_in_rocksdb requires rocksdb snapshot")]
1825 fn test_settings_mismatch_panics() {
1826 let factory = create_test_provider_factory();
1827
1828 factory.set_storage_settings_cache(StorageSettings::v2());
1829
1830 let provider = factory.database_provider_ro().unwrap();
1831 let _ = EitherReader::<(), ()>::new_accounts_history(&provider, None);
1832 }
1833}