Skip to main content

reth_db_common/
init.rs

1//! Reth genesis initialization utility functions.
2
3use alloy_consensus::BlockHeader;
4use alloy_genesis::GenesisAccount;
5use alloy_primitives::{
6    keccak256,
7    map::{AddressMap, B256Map, B256Set, HashMap},
8    Address, B256, U256,
9};
10use reth_chainspec::EthChainSpec;
11use reth_codecs::Compact;
12use reth_config::config::EtlConfig;
13use reth_db_api::{
14    cursor::{DbCursorRW, DbDupCursorRW},
15    models::{
16        storage_sharded_key::StorageShardedKey, AccountBeforeTx, BlockNumberAddress, IntegerList,
17        ShardedKey,
18    },
19    tables,
20    transaction::DbTxMut,
21    DatabaseError,
22};
23use reth_etl::Collector;
24use reth_execution_errors::StateRootError;
25use reth_primitives_traits::{
26    Account, Bytecode, GotExpected, NodePrimitives, SealedHeader, StorageEntry,
27};
28use reth_provider::{
29    errors::provider::ProviderResult, providers::StaticFileWriter, BlockHashReader, BlockNumReader,
30    BundleStateInit, ChainSpecProvider, DBProvider, DatabaseProviderFactory, ExecutionOutcome,
31    HashingWriter, HeaderProvider, HistoryWriter, MetadataProvider, MetadataWriter,
32    NodePrimitivesProvider, OriginalValuesKnown, ProviderError, RevertsInit,
33    RocksDBProviderFactory, StageCheckpointReader, StageCheckpointWriter, StateWriteConfig,
34    StateWriter, StaticFileProviderFactory, StorageSettings, StorageSettingsCache, TrieWriter,
35};
36use reth_stages_types::{StageCheckpoint, StageId};
37use reth_static_file_types::StaticFileSegment;
38use reth_trie::{
39    prefix_set::TriePrefixSets, IntermediateStateRootState, StateRoot as StateRootComputer,
40    StateRootProgress,
41};
42use reth_trie_db::DatabaseStateRoot;
43
44type DbStateRoot<'a, TX, A> = StateRootComputer<
45    reth_trie_db::DatabaseTrieCursorFactory<&'a TX, A>,
46    reth_trie_db::DatabaseHashedCursorFactory<&'a TX>,
47>;
48
49use serde::{Deserialize, Serialize};
50use std::io::BufRead;
51use tracing::{debug, error, info, trace, warn};
52
53pub use reth_provider::init::{
54    insert_account_history, insert_genesis_account_history, insert_genesis_history,
55    insert_genesis_storage_history, insert_history, insert_storage_history,
56};
57
58/// Default soft limit for number of bytes to read from state dump file, before inserting into
59/// database.
60///
61/// Default is 1 GB.
62pub const DEFAULT_SOFT_LIMIT_BYTE_LEN_ACCOUNTS_CHUNK: usize = 1_000_000_000;
63
64/// Soft limit for the number of flushed updates after which to log progress summary.
65const SOFT_LIMIT_COUNT_FLUSHED_UPDATES: usize = 1_000_000;
66
67/// Max number of storage "units" (1 per account + 1 per storage slot) before committing
68/// the current MDBX transaction and opening a new one. This bounds dirty page accumulation
69/// and prevents OOM on large state imports.
70const STORAGE_COMMIT_THRESHOLD: usize = 100_000;
71
72/// Max number of trie updates retained before init-state state root computation commits progress.
73const STATE_ROOT_COMMIT_THRESHOLD: u64 = 25_000;
74
75/// Storage initialization error type.
76#[derive(Debug, thiserror::Error, Clone)]
77pub enum InitStorageError {
78    /// Genesis header found on static files but the database is empty.
79    #[error(
80        "static files found, but the database is uninitialized. If attempting to re-syncing, delete both."
81    )]
82    UninitializedDatabase,
83    /// An existing genesis block was found in the database, and its hash did not match the hash of
84    /// the chainspec.
85    #[error(
86        "genesis hash in the storage does not match the specified chainspec: chainspec is {chainspec_hash}, database is {storage_hash}"
87    )]
88    GenesisHashMismatch {
89        /// Expected genesis hash.
90        chainspec_hash: B256,
91        /// Actual genesis hash.
92        storage_hash: B256,
93    },
94    /// Provider error.
95    #[error(transparent)]
96    Provider(#[from] ProviderError),
97    /// State root error while computing the state root
98    #[error(transparent)]
99    StateRootError(#[from] StateRootError),
100    /// State root doesn't match the expected one.
101    #[error("state root mismatch: {_0}")]
102    StateRootMismatch(GotExpected<B256>),
103}
104
105impl From<DatabaseError> for InitStorageError {
106    fn from(error: DatabaseError) -> Self {
107        Self::Provider(ProviderError::Database(error))
108    }
109}
110
111/// Write the genesis block if it has not already been written
112pub fn init_genesis<PF>(factory: &PF) -> Result<B256, InitStorageError>
113where
114    PF: DatabaseProviderFactory
115        + StaticFileProviderFactory<Primitives: NodePrimitives<BlockHeader: Compact>>
116        + ChainSpecProvider
117        + StageCheckpointReader
118        + BlockNumReader
119        + MetadataProvider
120        + StorageSettingsCache,
121    PF::ProviderRW: StaticFileProviderFactory<Primitives = PF::Primitives>
122        + StageCheckpointWriter
123        + HistoryWriter
124        + HeaderProvider
125        + HashingWriter
126        + StateWriter
127        + TrieWriter
128        + MetadataWriter
129        + ChainSpecProvider
130        + StorageSettingsCache
131        + RocksDBProviderFactory
132        + NodePrimitivesProvider
133        + AsRef<PF::ProviderRW>,
134    PF::ChainSpec: EthChainSpec<Header = <PF::Primitives as NodePrimitives>::BlockHeader>,
135{
136    init_genesis_with_settings(factory, StorageSettings::base())
137}
138
139/// Write the genesis block if it has not already been written with [`StorageSettings`].
140pub fn init_genesis_with_settings<PF>(
141    factory: &PF,
142    genesis_storage_settings: StorageSettings,
143) -> Result<B256, InitStorageError>
144where
145    PF: DatabaseProviderFactory
146        + StaticFileProviderFactory<Primitives: NodePrimitives<BlockHeader: Compact>>
147        + ChainSpecProvider
148        + StageCheckpointReader
149        + BlockNumReader
150        + MetadataProvider
151        + StorageSettingsCache,
152    PF::ProviderRW: StaticFileProviderFactory<Primitives = PF::Primitives>
153        + StageCheckpointWriter
154        + HistoryWriter
155        + HeaderProvider
156        + HashingWriter
157        + StateWriter
158        + TrieWriter
159        + MetadataWriter
160        + ChainSpecProvider
161        + StorageSettingsCache
162        + RocksDBProviderFactory
163        + NodePrimitivesProvider
164        + AsRef<PF::ProviderRW>,
165    PF::ChainSpec: EthChainSpec<Header = <PF::Primitives as NodePrimitives>::BlockHeader>,
166{
167    init_genesis_with_settings_and_validate(factory, genesis_storage_settings, true)
168}
169
170/// Write the genesis block if it has not already been written with [`StorageSettings`],
171/// optionally validating the DB-resident genesis hash against the chainspec hash.
172pub fn init_genesis_with_settings_and_validate<PF>(
173    factory: &PF,
174    genesis_storage_settings: StorageSettings,
175    validate_genesis_hash: bool,
176) -> Result<B256, InitStorageError>
177where
178    PF: DatabaseProviderFactory
179        + StaticFileProviderFactory<Primitives: NodePrimitives<BlockHeader: Compact>>
180        + ChainSpecProvider
181        + StageCheckpointReader
182        + BlockNumReader
183        + MetadataProvider
184        + StorageSettingsCache,
185    PF::ProviderRW: StaticFileProviderFactory<Primitives = PF::Primitives>
186        + StageCheckpointWriter
187        + HistoryWriter
188        + HeaderProvider
189        + HashingWriter
190        + StateWriter
191        + TrieWriter
192        + MetadataWriter
193        + ChainSpecProvider
194        + StorageSettingsCache
195        + RocksDBProviderFactory
196        + NodePrimitivesProvider
197        + AsRef<PF::ProviderRW>,
198    PF::ChainSpec: EthChainSpec<Header = <PF::Primitives as NodePrimitives>::BlockHeader>,
199{
200    let chain = factory.chain_spec();
201
202    let genesis = chain.genesis();
203    let hash = chain.genesis_hash();
204
205    // Get the genesis block number from the chain spec
206    let genesis_block_number = chain.genesis_header().number();
207
208    // Check if we already have the genesis header or if we have the wrong one.
209    match factory.block_hash(genesis_block_number) {
210        Ok(None) | Err(ProviderError::MissingStaticFileBlock(StaticFileSegment::Headers, _)) => {}
211        Ok(Some(block_hash)) => {
212            if block_hash == hash {
213                // Some users will at times attempt to re-sync from scratch by just deleting the
214                // database. Since `factory.block_hash` will only query the static files, we need to
215                // make sure that our database has been written to, and throw error if it's empty.
216                if factory.get_stage_checkpoint(StageId::Headers)?.is_none() {
217                    error!(target: "reth::storage", "Genesis header found on static files, but database is uninitialized.");
218                    return Err(InitStorageError::UninitializedDatabase)
219                }
220
221                let stored = factory.storage_settings()?.unwrap_or_else(StorageSettings::v1);
222                if stored != genesis_storage_settings {
223                    warn!(
224                        target: "reth::storage",
225                        ?stored,
226                        requested = ?genesis_storage_settings,
227                        "Storage settings mismatch detected. Using the stored settings from the existing database."
228                    );
229                }
230
231                debug!("Genesis already written, skipping.");
232                return Ok(hash)
233            }
234
235            if !validate_genesis_hash {
236                warn!(
237                    target: "reth::storage",
238                    chainspec_hash = %hash,
239                    storage_hash = %block_hash,
240                    "Genesis hash mismatch with chainspec; trusting DB per --debug.skip-genesis-validation"
241                );
242                return Ok(block_hash)
243            }
244            return Err(InitStorageError::GenesisHashMismatch {
245                chainspec_hash: hash,
246                storage_hash: block_hash,
247            })
248        }
249        Err(e) => {
250            debug!(?e);
251            return Err(e.into());
252        }
253    }
254
255    debug!("Writing genesis block.");
256
257    // Make sure to set storage settings before anything writes
258    factory.set_storage_settings_cache(genesis_storage_settings);
259
260    let alloc = &genesis.alloc;
261
262    // use transaction to insert genesis header
263    let provider_rw = factory.database_provider_rw()?;
264
265    // Behaviour reserved only for new nodes should be set in the storage settings.
266    provider_rw.write_storage_settings(genesis_storage_settings)?;
267
268    // For non-zero genesis blocks, set expected_block_start BEFORE insert_genesis_state.
269    // When block_range is None, next_block_number() uses expected_block_start. By default,
270    // expected_block_start comes from find_fixed_range which returns the file range start (0),
271    // not the genesis block number. This would cause increment_block(N) to fail.
272    let static_file_provider = provider_rw.static_file_provider();
273    if genesis_block_number > 0 {
274        if genesis_storage_settings.storage_v2 {
275            static_file_provider
276                .get_writer(genesis_block_number, StaticFileSegment::AccountChangeSets)?
277                .user_header_mut()
278                .set_expected_block_start(genesis_block_number);
279        }
280        if genesis_storage_settings.storage_v2 {
281            static_file_provider
282                .get_writer(genesis_block_number, StaticFileSegment::StorageChangeSets)?
283                .user_header_mut()
284                .set_expected_block_start(genesis_block_number);
285        }
286    }
287
288    insert_genesis_hashes(&provider_rw, alloc.iter())?;
289    insert_genesis_history(&provider_rw, alloc.iter())?;
290
291    // Insert header
292    insert_genesis_header(&provider_rw, &chain)?;
293
294    insert_genesis_state(&provider_rw, alloc.iter())?;
295
296    // compute state root to populate trie tables
297    compute_state_root(&provider_rw, None)?;
298
299    // set stage checkpoint to genesis block number for all stages
300    let checkpoint = StageCheckpoint::new(genesis_block_number);
301    for stage in StageId::ALL {
302        provider_rw.save_stage_checkpoint(stage, checkpoint)?;
303    }
304
305    // Static file segments start empty, so we need to initialize the block range.
306    // For genesis blocks with non-zero block numbers, we use get_writer() instead of
307    // latest_writer() and set_block_range() to ensure static files start at the correct block.
308    let static_file_provider = provider_rw.static_file_provider();
309
310    static_file_provider
311        .get_writer(genesis_block_number, StaticFileSegment::Receipts)?
312        .user_header_mut()
313        .set_block_range(genesis_block_number, genesis_block_number);
314    static_file_provider
315        .get_writer(genesis_block_number, StaticFileSegment::Transactions)?
316        .user_header_mut()
317        .set_block_range(genesis_block_number, genesis_block_number);
318
319    if genesis_storage_settings.storage_v2 {
320        static_file_provider
321            .get_writer(genesis_block_number, StaticFileSegment::TransactionSenders)?
322            .user_header_mut()
323            .set_block_range(genesis_block_number, genesis_block_number);
324    }
325
326    // `commit_unwind`` will first commit the DB and then the static file provider, which is
327    // necessary on `init_genesis`.
328    provider_rw.commit()?;
329
330    Ok(hash)
331}
332
333/// Inserts the genesis state into the database.
334pub fn insert_genesis_state<'a, 'b, Provider>(
335    provider: &Provider,
336    alloc: impl Iterator<Item = (&'a Address, &'b GenesisAccount)>,
337) -> ProviderResult<()>
338where
339    Provider: StaticFileProviderFactory
340        + DBProvider<Tx: DbTxMut>
341        + HeaderProvider
342        + StateWriter
343        + ChainSpecProvider
344        + AsRef<Provider>,
345{
346    let genesis_block_number = provider.chain_spec().genesis_header().number();
347    insert_state(provider, alloc, genesis_block_number)
348}
349
350/// Inserts state at given block into database.
351pub fn insert_state<'a, 'b, Provider>(
352    provider: &Provider,
353    alloc: impl Iterator<Item = (&'a Address, &'b GenesisAccount)>,
354    block: u64,
355) -> ProviderResult<()>
356where
357    Provider: StaticFileProviderFactory
358        + DBProvider<Tx: DbTxMut>
359        + HeaderProvider
360        + StateWriter
361        + AsRef<Provider>,
362{
363    let capacity = alloc.size_hint().1.unwrap_or(0);
364    let mut state_init: BundleStateInit =
365        AddressMap::with_capacity_and_hasher(capacity, Default::default());
366    let mut reverts_init: AddressMap<_> =
367        AddressMap::with_capacity_and_hasher(capacity, Default::default());
368    let mut contracts: B256Map<Bytecode> =
369        B256Map::with_capacity_and_hasher(capacity, Default::default());
370
371    for (address, account) in alloc {
372        let bytecode_hash = if let Some(code) = &account.code {
373            match Bytecode::new_raw_checked(code.clone()) {
374                Ok(bytecode) => {
375                    let hash = bytecode.hash_slow();
376                    contracts.insert(hash, bytecode);
377                    Some(hash)
378                }
379                Err(err) => {
380                    error!(%address, %err, "Failed to decode genesis bytecode.");
381                    return Err(DatabaseError::Other(err.to_string()).into());
382                }
383            }
384        } else {
385            None
386        };
387
388        // get state
389        let storage = account
390            .storage
391            .as_ref()
392            .map(|m| {
393                m.iter()
394                    .map(|(key, value)| {
395                        let value = U256::from_be_bytes(value.0);
396                        (*key, (U256::ZERO, value))
397                    })
398                    .collect::<B256Map<_>>()
399            })
400            .unwrap_or_default();
401
402        reverts_init.insert(
403            *address,
404            (Some(None), storage.keys().map(|k| StorageEntry::new(*k, U256::ZERO)).collect()),
405        );
406
407        state_init.insert(
408            *address,
409            (None, Some(Account { bytecode_hash, ..Account::from(account) }), storage),
410        );
411    }
412    let all_reverts_init: RevertsInit = HashMap::from_iter([(block, reverts_init)]);
413
414    let execution_outcome = ExecutionOutcome::new_init(
415        state_init,
416        all_reverts_init,
417        contracts,
418        Vec::default(),
419        block,
420        Vec::new(),
421    );
422
423    provider.write_state(
424        &execution_outcome,
425        OriginalValuesKnown::Yes,
426        StateWriteConfig::default(),
427    )?;
428
429    trace!(target: "reth::cli", "Inserted state");
430
431    Ok(())
432}
433
434/// Inserts hashes for the genesis state.
435pub fn insert_genesis_hashes<'a, 'b, Provider>(
436    provider: &Provider,
437    alloc: impl Iterator<Item = (&'a Address, &'b GenesisAccount)> + Clone,
438) -> ProviderResult<()>
439where
440    Provider: DBProvider<Tx: DbTxMut> + HashingWriter,
441{
442    // insert and hash accounts to hashing table
443    let alloc_accounts = alloc.clone().map(|(addr, account)| (*addr, Some(Account::from(account))));
444    provider.insert_account_for_hashing(alloc_accounts)?;
445
446    trace!(target: "reth::cli", "Inserted account hashes");
447
448    let alloc_storage = alloc.filter_map(|(addr, account)| {
449        // only return Some if there is storage
450        account.storage.as_ref().map(|storage| {
451            (*addr, storage.iter().map(|(&key, &value)| StorageEntry { key, value: value.into() }))
452        })
453    });
454    provider.insert_storage_for_hashing(alloc_storage)?;
455
456    trace!(target: "reth::cli", "Inserted storage hashes");
457
458    Ok(())
459}
460
461/// Inserts header for the genesis state.
462pub fn insert_genesis_header<Provider, Spec>(
463    provider: &Provider,
464    chain: &Spec,
465) -> ProviderResult<()>
466where
467    Provider: StaticFileProviderFactory<Primitives: NodePrimitives<BlockHeader: Compact>>
468        + DBProvider<Tx: DbTxMut>,
469    Spec: EthChainSpec<Header = <Provider::Primitives as NodePrimitives>::BlockHeader>,
470{
471    let (header, block_hash) = (chain.genesis_header(), chain.genesis_hash());
472    let static_file_provider = provider.static_file_provider();
473
474    // Get the actual genesis block number from the header
475    let genesis_block_number = header.number();
476
477    match static_file_provider.block_hash(genesis_block_number) {
478        Ok(None) | Err(ProviderError::MissingStaticFileBlock(StaticFileSegment::Headers, _)) => {
479            let difficulty = header.difficulty();
480
481            // For genesis blocks with non-zero block numbers, we need to ensure they are stored
482            // in the correct static file range. We use get_writer() with the genesis block number
483            // to ensure the genesis block is stored in the correct static file range.
484            let mut writer = static_file_provider
485                .get_writer(genesis_block_number, StaticFileSegment::Headers)?;
486
487            // For non-zero genesis blocks, we need to set block range to genesis_block_number and
488            // append header without increment block
489            if genesis_block_number > 0 {
490                writer
491                    .user_header_mut()
492                    .set_block_range(genesis_block_number, genesis_block_number);
493                writer.append_header_direct(header, difficulty, &block_hash)?;
494            } else {
495                // For zero genesis blocks, use normal append_header
496                writer.append_header(header, &block_hash)?;
497            }
498        }
499        Ok(Some(_)) => {}
500        Err(e) => return Err(e),
501    }
502
503    provider.tx_ref().put::<tables::HeaderNumbers>(block_hash, genesis_block_number)?;
504    provider.tx_ref().put::<tables::BlockBodyIndices>(genesis_block_number, Default::default())?;
505
506    Ok(())
507}
508
509/// Reads account state from a [`BufRead`] reader and initializes it at the highest block that can
510/// be found on database.
511///
512/// It's similar to [`init_genesis`] but supports importing state too big to fit in memory, and can
513/// be set to the highest block present. One practical usecase is to import OP mainnet state at
514/// bedrock transition block.
515pub fn init_from_state_dump<PF>(
516    mut reader: impl BufRead,
517    provider_factory: &PF,
518    etl_config: EtlConfig,
519) -> eyre::Result<B256>
520where
521    PF: DatabaseProviderFactory<
522        ProviderRW: StaticFileProviderFactory
523                        + DBProvider<Tx: DbTxMut>
524                        + BlockNumReader
525                        + BlockHashReader
526                        + ChainSpecProvider
527                        + StageCheckpointWriter
528                        + HistoryWriter
529                        + HeaderProvider
530                        + HashingWriter
531                        + TrieWriter
532                        + StateWriter
533                        + StorageSettingsCache
534                        + RocksDBProviderFactory
535                        + NodePrimitivesProvider
536                        + AsRef<PF::ProviderRW>,
537    >,
538{
539    if etl_config.file_size == 0 {
540        return Err(eyre::eyre!("ETL file size cannot be zero"))
541    }
542
543    let (block, hash, expected_state_root) = {
544        let provider_rw = provider_factory.database_provider_rw()?;
545        let block = provider_rw.last_block_number()?;
546        let hash = provider_rw
547            .block_hash(block)?
548            .ok_or_else(|| eyre::eyre!("Block hash not found for block {}", block))?;
549        let header = provider_rw
550            .header_by_number(block)?
551            .map(|h| SealedHeader::new(h, hash))
552            .ok_or_else(|| ProviderError::HeaderNotFound(block.into()))?;
553        let state_root = header.state_root();
554
555        debug!(target: "reth::cli",
556            block,
557            chain=%provider_rw.chain_spec().chain(),
558            "Initializing state at block"
559        );
560
561        (block, hash, state_root)
562    };
563
564    // first line can be state root
565    let dump_state_root = parse_state_root(&mut reader)?;
566    if expected_state_root != dump_state_root {
567        error!(target: "reth::cli",
568            ?dump_state_root,
569            ?expected_state_root,
570            "State root from state dump does not match state root in current header."
571        );
572        return Err(InitStorageError::StateRootMismatch(GotExpected {
573            got: dump_state_root,
574            expected: expected_state_root,
575        })
576        .into())
577    }
578
579    // remaining lines are accounts
580    let collector = parse_accounts(&mut reader, etl_config)?;
581
582    // write state to db with chunked commits to avoid OOM
583    dump_state(collector, provider_factory, block)?;
584
585    info!(target: "reth::cli", "All accounts written to database, starting state root computation (may take some time)");
586
587    // clear trie tables so state root is computed from scratch
588    {
589        let provider_rw = provider_factory.database_provider_rw()?;
590        provider_rw.tx_ref().clear::<tables::AccountsTrie>()?;
591        provider_rw.tx_ref().clear::<tables::StoragesTrie>()?;
592        provider_rw.commit()?;
593    }
594
595    // compute and compare state root
596    let computed_state_root = compute_state_root_chunked(provider_factory)?;
597    if computed_state_root == expected_state_root {
598        info!(target: "reth::cli",
599            ?computed_state_root,
600            "Computed state root matches state root in state dump"
601        );
602    } else {
603        error!(target: "reth::cli",
604            ?computed_state_root,
605            ?expected_state_root,
606            "Computed state root does not match state root in state dump"
607        );
608
609        return Err(InitStorageError::StateRootMismatch(GotExpected {
610            got: computed_state_root,
611            expected: expected_state_root,
612        })
613        .into())
614    }
615
616    // insert sync stages for stages that require state
617    {
618        let provider_rw = provider_factory.database_provider_rw()?;
619        for stage in StageId::STATE_REQUIRED {
620            provider_rw.save_stage_checkpoint(stage, StageCheckpoint::new(block))?;
621        }
622        provider_rw.commit()?;
623    }
624
625    Ok(hash)
626}
627
628/// Parses and returns expected state root.
629fn parse_state_root(reader: &mut impl BufRead) -> eyre::Result<B256> {
630    let mut line = String::new();
631    reader.read_line(&mut line)?;
632
633    let expected_state_root = serde_json::from_str::<StateRoot>(&line)?.root;
634    trace!(target: "reth::cli",
635        root=%expected_state_root,
636        "Read state root from file"
637    );
638    Ok(expected_state_root)
639}
640
641/// Parses accounts and pushes them to a [`Collector`].
642fn parse_accounts(
643    reader: impl BufRead,
644    etl_config: EtlConfig,
645) -> Result<Collector<Address, GenesisAccount>, eyre::Error> {
646    let mut collector = Collector::new(etl_config.file_size, etl_config.dir);
647    let mut parsed_accounts = 0usize;
648
649    let stream =
650        serde_json::Deserializer::from_reader(reader).into_iter::<GenesisAccountWithAddress>();
651    for account in stream {
652        let GenesisAccountWithAddress { genesis_account, address } = account?;
653        collector.insert(address, genesis_account)?;
654
655        parsed_accounts += 1;
656        if parsed_accounts.is_multiple_of(100_000) {
657            info!(target: "reth::cli", parsed_accounts, "Parsed accounts");
658        }
659    }
660
661    Ok(collector)
662}
663
664/// Takes a [`Collector`] and writes all accounts directly to database tables.
665///
666/// This bypasses the higher-level `insert_state`/`insert_genesis_hashes`/`insert_history`
667/// functions which build intermediate structures (`BundleStateInit`, `RevertsInit`,
668/// `ExecutionOutcome`) that duplicate all storage data 2-3x in memory. For accounts with
669/// millions of storage entries this causes OOM.
670///
671/// Instead, each account is written directly to all required tables using cursor operations,
672/// using `append`/`append_dup` for sorted tables where possible (MDBX fast path that skips
673/// B-tree traversal). Commits happen every [`STORAGE_COMMIT_THRESHOLD`] storage units to
674/// bound MDBX dirty page accumulation.
675///
676/// NOTE: This function is not idempotent. If the process crashes mid-import, the database
677/// must be wiped before retrying.
678fn dump_state<PF>(
679    mut collector: Collector<Address, GenesisAccount>,
680    provider_factory: &PF,
681    block: u64,
682) -> Result<(), eyre::Error>
683where
684    PF: DatabaseProviderFactory<ProviderRW: DBProvider<Tx: DbTxMut>>,
685    PF::ProviderRW: StaticFileProviderFactory
686        + StorageSettingsCache
687        + RocksDBProviderFactory
688        + NodePrimitivesProvider,
689{
690    let storage_settings = provider_factory.database_provider_rw()?.cached_storage_settings();
691    if storage_settings.storage_v2 {
692        return dump_state_v2(collector, provider_factory, block)
693    }
694
695    let accounts_len = collector.len();
696    let mut total_accounts: usize = 0;
697    let mut storage_units: usize = 0;
698
699    // pre-allocate the history list once — every entry uses the same single-block bitmap
700    let history_list = IntegerList::new([block])?;
701
702    // track seen bytecode hashes to avoid re-hashing and re-writing duplicates
703    let mut seen_bytecodes: B256Set = B256Set::default();
704
705    let mut provider_rw = provider_factory.database_provider_rw()?;
706
707    for entry in collector.iter()? {
708        let (address_raw, account_raw) = entry?;
709        let (address, _) = Address::from_compact(address_raw.as_slice(), address_raw.len());
710        let (account, _) = GenesisAccount::from_compact(account_raw.as_slice(), account_raw.len());
711
712        let account_storage_len = account.storage.as_ref().map_or(0, |s| s.len());
713        let account_units = 1 + account_storage_len;
714
715        // commit before this account would push us over the threshold
716        if storage_units > 0 && storage_units + account_units > STORAGE_COMMIT_THRESHOLD {
717            provider_rw.commit()?;
718            provider_rw = provider_factory.database_provider_rw()?;
719            info!(target: "reth::cli",
720                total_accounts,
721                accounts_len,
722                storage_units,
723                "Committed chunk"
724            );
725            storage_units = 0;
726            seen_bytecodes = B256Set::default();
727        }
728
729        write_account_to_db(
730            provider_rw.tx_ref(),
731            &address,
732            &account,
733            block,
734            &history_list,
735            &mut seen_bytecodes,
736        )?;
737
738        total_accounts += 1;
739        storage_units += account_units;
740
741        if total_accounts.is_multiple_of(100_000) {
742            info!(target: "reth::cli", total_accounts, accounts_len, "Writing accounts...");
743        }
744    }
745
746    // commit final batch
747    provider_rw.commit()?;
748
749    info!(target: "reth::cli", total_accounts, "All accounts written to database");
750
751    Ok(())
752}
753
754fn dump_state_v2<PF>(
755    mut collector: Collector<Address, GenesisAccount>,
756    provider_factory: &PF,
757    block: u64,
758) -> Result<(), eyre::Error>
759where
760    PF: DatabaseProviderFactory<
761        ProviderRW: StaticFileProviderFactory
762                        + DBProvider<Tx: DbTxMut>
763                        + StorageSettingsCache
764                        + RocksDBProviderFactory
765                        + NodePrimitivesProvider,
766    >,
767{
768    let accounts_len = collector.len();
769    let mut total_accounts: usize = 0;
770    let mut storage_units: usize = 0;
771
772    // pre-allocate the history list once — every entry uses the same single-block bitmap
773    let history_list = IntegerList::new([block])?;
774
775    // track seen bytecode hashes to avoid re-hashing and re-writing duplicates
776    let mut seen_bytecodes: B256Set = B256Set::default();
777
778    let mut provider_rw = provider_factory.database_provider_rw()?;
779    let static_file_provider = provider_rw.static_file_provider();
780    let rocksdb_provider = provider_rw.rocksdb_provider();
781    let mut history_batch = rocksdb_provider.batch_with_auto_commit();
782    if snapshot_state_tables_empty(provider_rw.tx_ref())? {
783        reset_pre_snapshot_changeset_segment(
784            &static_file_provider,
785            StaticFileSegment::AccountChangeSets,
786            block,
787        )?;
788        reset_pre_snapshot_changeset_segment(
789            &static_file_provider,
790            StaticFileSegment::StorageChangeSets,
791            block,
792        )?;
793    }
794
795    {
796        let mut account_changeset_writer =
797            static_file_provider.get_writer(block, StaticFileSegment::AccountChangeSets)?;
798        let mut storage_changeset_writer =
799            static_file_provider.get_writer(block, StaticFileSegment::StorageChangeSets)?;
800        prepare_account_changeset_writer(&mut account_changeset_writer, block)?;
801        prepare_storage_changeset_writer(&mut storage_changeset_writer, block)?;
802
803        for entry in collector.iter()? {
804            let (address_raw, account_raw) = entry?;
805            let (address, _) = Address::from_compact(address_raw.as_slice(), address_raw.len());
806            let (account, _) =
807                GenesisAccount::from_compact(account_raw.as_slice(), account_raw.len());
808
809            let account_storage_len = account.storage.as_ref().map_or(0, |s| s.len());
810            let account_units = 1 + account_storage_len;
811
812            // commit before this account would push us over the threshold
813            if storage_units > 0 && storage_units + account_units > STORAGE_COMMIT_THRESHOLD {
814                history_batch.commit()?;
815                commit_mdbx_only(provider_rw)?;
816                provider_rw = provider_factory.database_provider_rw()?;
817                history_batch = rocksdb_provider.batch_with_auto_commit();
818                info!(target: "reth::cli",
819                    total_accounts,
820                    accounts_len,
821                    storage_units,
822                    "Committed chunk"
823                );
824                storage_units = 0;
825                seen_bytecodes = B256Set::default();
826            }
827
828            write_account_to_db_v2(
829                provider_rw.tx_ref(),
830                (&mut account_changeset_writer, &mut storage_changeset_writer),
831                &mut history_batch,
832                &address,
833                &account,
834                &history_list,
835                &mut seen_bytecodes,
836            )?;
837
838            total_accounts += 1;
839            storage_units += account_units;
840
841            if total_accounts.is_multiple_of(100_000) {
842                info!(target: "reth::cli", total_accounts, accounts_len, "Writing accounts...");
843            }
844        }
845    }
846
847    history_batch.commit()?;
848    commit_mdbx_only(provider_rw)?;
849    static_file_provider.finalize()?;
850
851    info!(target: "reth::cli", total_accounts, "All accounts written to database");
852
853    Ok(())
854}
855
856fn prepare_account_changeset_writer<N: NodePrimitives>(
857    writer: &mut reth_provider::providers::StaticFileProviderRWRefMut<'_, N>,
858    block: u64,
859) -> ProviderResult<()> {
860    let next_block = writer.next_block_number();
861    if next_block < block {
862        info!(
863            target: "reth::cli",
864            from_block = next_block,
865            to_block = block - 1,
866            "Padding empty account changesets before state import"
867        );
868        for empty_block in next_block..block {
869            writer.append_account_changeset(Vec::new(), empty_block)?;
870            if empty_block > next_block && empty_block.is_multiple_of(1_000_000) {
871                info!(
872                    target: "reth::cli",
873                    padded_to_block = empty_block,
874                    "Padded empty account changesets"
875                );
876            }
877        }
878    }
879
880    writer.begin_account_changeset(block)
881}
882
883fn prepare_storage_changeset_writer<N: NodePrimitives>(
884    writer: &mut reth_provider::providers::StaticFileProviderRWRefMut<'_, N>,
885    block: u64,
886) -> ProviderResult<()> {
887    let next_block = writer.next_block_number();
888    if next_block < block {
889        info!(
890            target: "reth::cli",
891            from_block = next_block,
892            to_block = block - 1,
893            "Padding empty storage changesets before state import"
894        );
895        for empty_block in next_block..block {
896            writer.append_storage_changeset(Vec::new(), empty_block)?;
897            if empty_block > next_block && empty_block.is_multiple_of(1_000_000) {
898                info!(
899                    target: "reth::cli",
900                    padded_to_block = empty_block,
901                    "Padded empty storage changesets"
902                );
903            }
904        }
905    }
906
907    writer.begin_storage_changeset(block)
908}
909
910fn snapshot_state_tables_empty<TX: reth_db_api::transaction::DbTx>(
911    tx: &TX,
912) -> ProviderResult<bool> {
913    Ok(tx.entries::<tables::PlainAccountState>()? == 0 &&
914        tx.entries::<tables::PlainStorageState>()? == 0 &&
915        tx.entries::<tables::HashedAccounts>()? == 0 &&
916        tx.entries::<tables::HashedStorages>()? == 0 &&
917        tx.entries::<tables::AccountChangeSets>()? == 0 &&
918        tx.entries::<tables::StorageChangeSets>()? == 0 &&
919        tx.entries::<tables::Bytecodes>()? == 0)
920}
921
922fn reset_pre_snapshot_changeset_segment<N: NodePrimitives>(
923    static_file_provider: &reth_provider::providers::StaticFileProvider<N>,
924    segment: StaticFileSegment,
925    block: u64,
926) -> ProviderResult<()> {
927    if block == 0 {
928        return Ok(())
929    }
930
931    let Some(highest_block) = static_file_provider.get_highest_static_file_block(segment) else {
932        return Ok(())
933    };
934
935    if highest_block >= block {
936        return Ok(())
937    }
938
939    let file_start = static_file_provider.find_fixed_range(segment, block).start();
940    info!(
941        target: "reth::cli",
942        ?segment,
943        highest_block,
944        import_block = block,
945        file_start,
946        "Resetting pre-snapshot changeset static files before state import"
947    );
948    static_file_provider.delete_segment(segment)?;
949
950    Ok(())
951}
952
953fn commit_mdbx_only<Provider>(provider: Provider) -> ProviderResult<()>
954where
955    Provider: DBProvider<Tx: DbTxMut>,
956{
957    reth_db_api::transaction::DbTx::commit(provider.into_tx()).map_err(ProviderError::from)
958}
959
960/// Writes a single account and all its storage to every required DB table directly,
961/// without building intermediary structures.
962///
963/// Uses `append_dup` for `DupSort` tables where insertion order matches key order (the ETL
964/// collector sorts by address, so `AccountChangeSets`, `PlainStorageState`, and
965/// `StorageChangeSets` receive data in sorted order within each account). For `HashedAccounts`
966/// and `HashedStorages`, insertion order is unsorted (keccak scrambles address order), so we
967/// use `put`/`upsert` which do a full B-tree lookup.
968#[allow(clippy::clone_on_copy)]
969fn write_account_to_db<TX: DbTxMut>(
970    tx: &TX,
971    address: &Address,
972    genesis_account: &GenesisAccount,
973    block: u64,
974    history_list: &IntegerList,
975    seen_bytecodes: &mut B256Set,
976) -> Result<(), eyre::Error> {
977    let bytecode_hash = if let Some(code) = &genesis_account.code {
978        let bytecode = Bytecode::new_raw_checked(code.clone())
979            .map_err(|e| eyre::eyre!("Invalid bytecode for {address}: {e}"))?;
980        let hash = bytecode.hash_slow();
981        if seen_bytecodes.insert(hash) {
982            tx.put::<tables::Bytecodes>(hash, bytecode)?;
983        }
984        Some(hash)
985    } else {
986        None
987    };
988
989    let account = Account { bytecode_hash, ..Account::from(genesis_account) };
990
991    let hashed_address = keccak256(address);
992
993    // plain state — sorted by address (ETL order), use append
994    tx.put::<tables::PlainAccountState>(*address, account.clone())?;
995
996    // hashed state — unsorted (keccak scrambles order), must use put
997    tx.put::<tables::HashedAccounts>(hashed_address, account)?;
998
999    // account changeset — DupSort keyed by block, subkey sorted by address (ETL order)
1000    let mut acct_cs_cursor = tx.cursor_dup_write::<tables::AccountChangeSets>()?;
1001    acct_cs_cursor.append_dup(block, AccountBeforeTx { address: *address, info: None })?;
1002
1003    // account history
1004    tx.put::<tables::AccountsHistory>(ShardedKey::new(*address, u64::MAX), history_list.clone())?;
1005
1006    // storage entries
1007    if let Some(storage) = &genesis_account.storage {
1008        let mut hashed_storage_cursor = tx.cursor_dup_write::<tables::HashedStorages>()?;
1009        let mut plain_storage_cursor = tx.cursor_dup_write::<tables::PlainStorageState>()?;
1010        let mut storage_cs_cursor = tx.cursor_dup_write::<tables::StorageChangeSets>()?;
1011
1012        for (&key, &value) in storage {
1013            let value_u256 = U256::from_be_bytes(value.0);
1014
1015            // plain storage — sorted by (address, key), use append_dup
1016            plain_storage_cursor.append_dup(*address, StorageEntry { key, value: value_u256 })?;
1017
1018            // hashed storage — unsorted keccak order, use upsert
1019            let hashed_key = keccak256(key);
1020            hashed_storage_cursor
1021                .upsert(hashed_address, &StorageEntry { key: hashed_key, value: value_u256 })?;
1022
1023            // storage changeset — sorted by (block, address), then by key via append_dup
1024            storage_cs_cursor.append_dup(
1025                BlockNumberAddress((block, *address)),
1026                StorageEntry { key, value: U256::ZERO },
1027            )?;
1028
1029            // storage history
1030            tx.put::<tables::StoragesHistory>(
1031                StorageShardedKey::new(*address, key, u64::MAX),
1032                history_list.clone(),
1033            )?;
1034        }
1035    }
1036
1037    Ok(())
1038}
1039
1040/// Writes a single account to the v2 storage destinations.
1041///
1042/// Storage v2 uses hashed state as the canonical state, static-file change sets, and `RocksDB`
1043/// history indices. The ETL collector yields accounts sorted by address and genesis storage is a
1044/// `BTreeMap`, so the streaming static-file writes preserve the required order.
1045fn write_account_to_db_v2<TX, N>(
1046    tx: &TX,
1047    changeset_writers: (
1048        &mut reth_provider::providers::StaticFileProviderRWRefMut<'_, N>,
1049        &mut reth_provider::providers::StaticFileProviderRWRefMut<'_, N>,
1050    ),
1051    history_batch: &mut reth_provider::providers::RocksDBBatch<'_>,
1052    address: &Address,
1053    genesis_account: &GenesisAccount,
1054    history_list: &IntegerList,
1055    seen_bytecodes: &mut B256Set,
1056) -> Result<(), eyre::Error>
1057where
1058    TX: DbTxMut,
1059    N: NodePrimitives,
1060{
1061    let bytecode_hash = if let Some(code) = &genesis_account.code {
1062        let bytecode = Bytecode::new_raw_checked(code.clone())
1063            .map_err(|e| eyre::eyre!("Invalid bytecode for {address}: {e}"))?;
1064        let hash = bytecode.hash_slow();
1065        if seen_bytecodes.insert(hash) {
1066            tx.put::<tables::Bytecodes>(hash, bytecode)?;
1067        }
1068        Some(hash)
1069    } else {
1070        None
1071    };
1072
1073    let account = Account { bytecode_hash, ..Account::from(genesis_account) };
1074
1075    let hashed_address = keccak256(address);
1076    let (account_changeset_writer, storage_changeset_writer) = changeset_writers;
1077
1078    tx.put::<tables::HashedAccounts>(hashed_address, account)?;
1079    account_changeset_writer
1080        .append_account_changeset_entry(AccountBeforeTx { address: *address, info: None })?;
1081    history_batch
1082        .put::<tables::AccountsHistory>(ShardedKey::new(*address, u64::MAX), history_list)?;
1083
1084    if let Some(storage) = &genesis_account.storage {
1085        let mut hashed_storage_cursor = tx.cursor_dup_write::<tables::HashedStorages>()?;
1086
1087        for (&key, &value) in storage {
1088            let value_u256 = U256::from_be_bytes(value.0);
1089
1090            let hashed_key = keccak256(key);
1091            hashed_storage_cursor
1092                .upsert(hashed_address, &StorageEntry { key: hashed_key, value: value_u256 })?;
1093
1094            storage_changeset_writer.append_storage_changeset_entry(
1095                reth_db_api::models::StorageBeforeTx { address: *address, key, value: U256::ZERO },
1096            )?;
1097
1098            history_batch.put::<tables::StoragesHistory>(
1099                StorageShardedKey::new(*address, key, u64::MAX),
1100                history_list,
1101            )?;
1102        }
1103    }
1104
1105    Ok(())
1106}
1107
1108/// Computes the state root (from scratch) based on the accounts and storages present in the
1109/// database.
1110fn compute_state_root<Provider>(
1111    provider: &Provider,
1112    prefix_sets: Option<TriePrefixSets>,
1113) -> Result<B256, InitStorageError>
1114where
1115    Provider: DBProvider<Tx: DbTxMut> + TrieWriter + StorageSettingsCache,
1116{
1117    reth_trie_db::with_adapter!(provider, |A| {
1118        compute_state_root_inner::<_, A>(provider, prefix_sets)
1119    })
1120}
1121
1122fn compute_state_root_inner<Provider, A>(
1123    provider: &Provider,
1124    prefix_sets: Option<TriePrefixSets>,
1125) -> Result<B256, InitStorageError>
1126where
1127    Provider: DBProvider<Tx: DbTxMut> + TrieWriter + StorageSettingsCache,
1128    A: reth_trie_db::TrieTableAdapter,
1129{
1130    trace!(target: "reth::cli", "Computing state root");
1131
1132    let tx = provider.tx_ref();
1133    let mut intermediate_state: Option<IntermediateStateRootState> = None;
1134    let mut total_flushed_updates = 0;
1135
1136    loop {
1137        let mut state_root =
1138            DbStateRoot::<_, A>::from_tx(tx).with_intermediate_state(intermediate_state);
1139
1140        if let Some(sets) = prefix_sets.clone() {
1141            state_root = state_root.with_prefix_sets(sets);
1142        }
1143
1144        match state_root.root_with_progress()? {
1145            StateRootProgress::Progress(state, _, updates) => {
1146                let updated_len = provider.write_trie_updates(updates)?;
1147                total_flushed_updates += updated_len;
1148
1149                trace!(target: "reth::cli",
1150                    last_account_key = %state.account_root_state.last_hashed_key,
1151                    updated_len,
1152                    total_flushed_updates,
1153                    "Flushing trie updates"
1154                );
1155
1156                intermediate_state = Some(*state);
1157
1158                if total_flushed_updates.is_multiple_of(SOFT_LIMIT_COUNT_FLUSHED_UPDATES) {
1159                    info!(target: "reth::cli",
1160                        total_flushed_updates,
1161                        "Flushing trie updates"
1162                    );
1163                }
1164            }
1165            StateRootProgress::Complete(root, _, updates) => {
1166                let updated_len = provider.write_trie_updates(updates)?;
1167                total_flushed_updates += updated_len;
1168
1169                trace!(target: "reth::cli",
1170                    %root,
1171                    updated_len,
1172                    total_flushed_updates,
1173                    "State root has been computed"
1174                );
1175
1176                return Ok(root)
1177            }
1178        }
1179    }
1180}
1181
1182/// Computes the state root (from scratch) with periodic commits to free MDBX dirty pages.
1183///
1184/// Opens a fresh transaction each iteration to release dirty pages, preventing OOM on large
1185/// states where trie updates accumulate gigabytes of MDBX dirty pages.
1186fn compute_state_root_chunked<PF>(provider_factory: &PF) -> Result<B256, InitStorageError>
1187where
1188    PF: DatabaseProviderFactory<
1189        ProviderRW: DBProvider<Tx: DbTxMut> + TrieWriter + StorageSettingsCache,
1190    >,
1191{
1192    let provider_rw = provider_factory.database_provider_rw().map_err(provider_db_err)?;
1193
1194    reth_trie_db::with_adapter!(&provider_rw, |A| {
1195        drop(provider_rw);
1196        compute_state_root_chunked_inner::<PF, A>(provider_factory)
1197    })
1198}
1199
1200fn compute_state_root_chunked_inner<PF, A>(provider_factory: &PF) -> Result<B256, InitStorageError>
1201where
1202    PF: DatabaseProviderFactory<
1203        ProviderRW: DBProvider<Tx: DbTxMut> + TrieWriter + StorageSettingsCache,
1204    >,
1205    A: reth_trie_db::TrieTableAdapter,
1206{
1207    trace!(target: "reth::cli", "Computing state root");
1208
1209    let mut intermediate_state: Option<IntermediateStateRootState> = None;
1210    let mut total_flushed_updates = 0;
1211
1212    loop {
1213        let provider_rw = provider_factory.database_provider_rw().map_err(provider_db_err)?;
1214        let tx = provider_rw.tx_ref();
1215
1216        let state_root = DbStateRoot::<_, A>::from_tx(tx)
1217            .with_intermediate_state(intermediate_state.take())
1218            .with_threshold(STATE_ROOT_COMMIT_THRESHOLD);
1219
1220        match state_root.root_with_progress()? {
1221            StateRootProgress::Progress(state, _, updates) => {
1222                let updated_len = provider_rw.write_trie_updates(updates)?;
1223                total_flushed_updates += updated_len;
1224
1225                info!(target: "reth::cli",
1226                    last_account_key = %state.account_root_state.last_hashed_key,
1227                    updated_len,
1228                    total_flushed_updates,
1229                    "Flushing trie updates (committing to free memory)"
1230                );
1231
1232                intermediate_state = Some(*state);
1233                provider_rw.commit().map_err(provider_db_err)?;
1234            }
1235            StateRootProgress::Complete(root, _, updates) => {
1236                let updated_len = provider_rw.write_trie_updates(updates)?;
1237                total_flushed_updates += updated_len;
1238
1239                info!(target: "reth::cli",
1240                    %root,
1241                    updated_len,
1242                    total_flushed_updates,
1243                    "State root computation complete"
1244                );
1245
1246                provider_rw.commit().map_err(provider_db_err)?;
1247                return Ok(root)
1248            }
1249        }
1250    }
1251}
1252
1253/// Converts a provider error into an [`InitStorageError`].
1254fn provider_db_err(e: impl std::fmt::Display) -> InitStorageError {
1255    InitStorageError::from(StateRootError::Database(DatabaseError::Other(e.to_string())))
1256}
1257
1258/// Type to deserialize state root from state dump file.
1259#[derive(Debug, Serialize, Deserialize, PartialEq, Eq)]
1260struct StateRoot {
1261    root: B256,
1262}
1263
1264/// An account as in the state dump file. This contains a [`GenesisAccount`] and the account's
1265/// address.
1266#[derive(Debug, Serialize, Deserialize)]
1267struct GenesisAccountWithAddress {
1268    /// The account's balance, nonce, code, and storage.
1269    #[serde(flatten)]
1270    genesis_account: GenesisAccount,
1271    /// The account's address.
1272    address: Address,
1273}
1274
1275#[cfg(test)]
1276mod tests {
1277    use super::*;
1278    use alloy_consensus::constants::{
1279        HOLESKY_GENESIS_HASH, MAINNET_GENESIS_HASH, SEPOLIA_GENESIS_HASH,
1280    };
1281    use alloy_genesis::Genesis;
1282    use reth_chainspec::{Chain, ChainSpec, HOLESKY, MAINNET, SEPOLIA};
1283    use reth_db::DatabaseEnv;
1284    use reth_db_api::{
1285        cursor::DbCursorRO,
1286        models::{storage_sharded_key::StorageShardedKey, IntegerList, ShardedKey},
1287        table::{Table, TableRow},
1288        transaction::DbTx,
1289        Database,
1290    };
1291    use reth_provider::{
1292        test_utils::{create_test_provider_factory_with_chain_spec, MockNodeTypesWithDB},
1293        ProviderFactory, RocksDBProviderFactory,
1294    };
1295    use std::{collections::BTreeMap, sync::Arc};
1296
1297    fn collect_table_entries<DB, T>(
1298        tx: &<DB as Database>::TX,
1299    ) -> Result<Vec<TableRow<T>>, InitStorageError>
1300    where
1301        DB: Database,
1302        T: Table,
1303    {
1304        Ok(tx.cursor_read::<T>()?.walk_range(..)?.collect::<Result<Vec<_>, _>>()?)
1305    }
1306
1307    #[test]
1308    fn parse_accounts_streams_jsonl_accounts() {
1309        let input = br#"{"address":"0x0000000000000000000000000000000000000002","balance":"0x2"}
1310{"address":"0x0000000000000000000000000000000000000001","balance":"0x1"}
1311"#;
1312
1313        let mut collector =
1314            parse_accounts(&input[..], EtlConfig::new(None, 128)).expect("parse succeeds");
1315
1316        let accounts = collector
1317            .iter()
1318            .unwrap()
1319            .map(|entry| {
1320                let (address_raw, account_raw) = entry.unwrap();
1321                let (address, _) = Address::from_compact(address_raw.as_slice(), address_raw.len());
1322                let (account, _) =
1323                    GenesisAccount::from_compact(account_raw.as_slice(), account_raw.len());
1324                (address, account.balance)
1325            })
1326            .collect::<Vec<_>>();
1327
1328        assert_eq!(
1329            accounts,
1330            vec![
1331                (Address::with_last_byte(1), U256::from(1)),
1332                (Address::with_last_byte(2), U256::from(2))
1333            ]
1334        );
1335    }
1336
1337    #[test]
1338    fn dump_state_uses_storage_v2_destinations() {
1339        let storage_key = B256::with_last_byte(3);
1340        let input = br#"{"address":"0x0000000000000000000000000000000000000001","balance":"0x1"}
1341{"address":"0x0000000000000000000000000000000000000002","balance":"0x0","storage":{"0x0000000000000000000000000000000000000000000000000000000000000003":"0x0000000000000000000000000000000000000000000000000000000000000004"}}
1342"#;
1343
1344        let collector = parse_accounts(&input[..], EtlConfig::new(None, 128)).unwrap();
1345        let factory = create_test_provider_factory_with_chain_spec(MAINNET.clone());
1346        factory.set_storage_settings_cache(StorageSettings::v2());
1347        let block = 10;
1348
1349        dump_state(collector, &factory, block).unwrap();
1350
1351        let provider = factory.provider().unwrap();
1352        let tx = provider.tx_ref();
1353        assert_eq!(tx.entries::<tables::PlainAccountState>().unwrap(), 0);
1354        assert_eq!(tx.entries::<tables::PlainStorageState>().unwrap(), 0);
1355        assert_eq!(tx.entries::<tables::AccountChangeSets>().unwrap(), 0);
1356        assert_eq!(tx.entries::<tables::StorageChangeSets>().unwrap(), 0);
1357        assert_eq!(tx.entries::<tables::HashedAccounts>().unwrap(), 2);
1358        assert_eq!(tx.entries::<tables::HashedStorages>().unwrap(), 1);
1359
1360        let address_with_balance = Address::with_last_byte(1);
1361        let address_with_storage = Address::with_last_byte(2);
1362        assert_eq!(
1363            reth_provider::ChangeSetReader::account_block_changeset(&provider, block).unwrap(),
1364            vec![
1365                AccountBeforeTx { address: address_with_balance, info: None },
1366                AccountBeforeTx { address: address_with_storage, info: None }
1367            ]
1368        );
1369        assert_eq!(
1370            reth_provider::StorageChangeSetReader::storage_changeset(&provider, block).unwrap(),
1371            vec![(
1372                BlockNumberAddress((block, address_with_storage)),
1373                StorageEntry { key: storage_key, value: U256::ZERO }
1374            )]
1375        );
1376
1377        let rocksdb = factory.rocksdb_provider();
1378        let accounts = rocksdb
1379            .iter::<tables::AccountsHistory>()
1380            .unwrap()
1381            .collect::<Result<Vec<_>, _>>()
1382            .unwrap();
1383        let storages = rocksdb
1384            .iter::<tables::StoragesHistory>()
1385            .unwrap()
1386            .collect::<Result<Vec<_>, _>>()
1387            .unwrap();
1388
1389        assert_eq!(
1390            accounts,
1391            vec![
1392                (
1393                    ShardedKey::new(address_with_balance, u64::MAX),
1394                    IntegerList::new([block]).unwrap()
1395                ),
1396                (
1397                    ShardedKey::new(address_with_storage, u64::MAX),
1398                    IntegerList::new([block]).unwrap()
1399                )
1400            ]
1401        );
1402        assert_eq!(
1403            storages,
1404            vec![(
1405                StorageShardedKey::new(address_with_storage, storage_key, u64::MAX),
1406                IntegerList::new([block]).unwrap()
1407            )]
1408        );
1409    }
1410
1411    #[test]
1412    fn dump_state_v2_resets_presnapshot_changeset_static_files() {
1413        let storage_key = B256::with_last_byte(3);
1414        let input = br#"{"address":"0x0000000000000000000000000000000000000002","balance":"0x0","storage":{"0x0000000000000000000000000000000000000000000000000000000000000003":"0x0000000000000000000000000000000000000000000000000000000000000004"}}
1415"#;
1416
1417        let collector = parse_accounts(&input[..], EtlConfig::new(None, 128)).unwrap();
1418        let factory = create_test_provider_factory_with_chain_spec(MAINNET.clone());
1419        factory.set_storage_settings_cache(StorageSettings::v2());
1420        let static_files = factory.static_file_provider();
1421
1422        {
1423            let mut writer =
1424                static_files.get_writer(0, StaticFileSegment::AccountChangeSets).unwrap();
1425            writer.append_account_changeset(Vec::new(), 0).unwrap();
1426        }
1427        {
1428            let mut writer =
1429                static_files.get_writer(0, StaticFileSegment::StorageChangeSets).unwrap();
1430            writer.append_storage_changeset(Vec::new(), 0).unwrap();
1431        }
1432        static_files.commit().unwrap();
1433
1434        let block = 500_010;
1435        dump_state(collector, &factory, block).unwrap();
1436        static_files.initialize_index().unwrap();
1437
1438        let provider = factory.provider().unwrap();
1439        let address = Address::with_last_byte(2);
1440        assert!(reth_provider::ChangeSetReader::account_block_changeset(&provider, 5)
1441            .unwrap()
1442            .is_empty());
1443        assert!(reth_provider::StorageChangeSetReader::storage_changeset(&provider, 5)
1444            .unwrap()
1445            .is_empty());
1446
1447        let account_file_start =
1448            static_files.find_fixed_range(StaticFileSegment::AccountChangeSets, block).start();
1449        let storage_file_start =
1450            static_files.find_fixed_range(StaticFileSegment::StorageChangeSets, block).start();
1451        assert_eq!(account_file_start, 500_000);
1452        assert_eq!(storage_file_start, 500_000);
1453        assert!(reth_provider::ChangeSetReader::account_block_changeset(
1454            &provider,
1455            account_file_start
1456        )
1457        .unwrap()
1458        .is_empty());
1459        assert!(reth_provider::StorageChangeSetReader::storage_changeset(
1460            &provider,
1461            storage_file_start
1462        )
1463        .unwrap()
1464        .is_empty());
1465
1466        assert_eq!(
1467            reth_provider::ChangeSetReader::account_block_changeset(&provider, block).unwrap(),
1468            vec![AccountBeforeTx { address, info: None }]
1469        );
1470        assert_eq!(
1471            reth_provider::StorageChangeSetReader::storage_changeset(&provider, block).unwrap(),
1472            vec![(
1473                BlockNumberAddress((block, address)),
1474                StorageEntry { key: storage_key, value: U256::ZERO }
1475            )]
1476        );
1477
1478        let account_offsets = static_files
1479            .get_segment_provider_for_block(StaticFileSegment::AccountChangeSets, block, None)
1480            .unwrap()
1481            .read_changeset_offsets()
1482            .unwrap()
1483            .unwrap();
1484        let storage_offsets = static_files
1485            .get_segment_provider_for_block(StaticFileSegment::StorageChangeSets, block, None)
1486            .unwrap()
1487            .read_changeset_offsets()
1488            .unwrap()
1489            .unwrap();
1490        assert_eq!(account_offsets.len() as u64, block - account_file_start + 1);
1491        assert_eq!(storage_offsets.len() as u64, block - storage_file_start + 1);
1492    }
1493
1494    #[test]
1495    fn success_init_genesis_mainnet() {
1496        let genesis_hash =
1497            init_genesis(&create_test_provider_factory_with_chain_spec(MAINNET.clone())).unwrap();
1498
1499        // actual, expected
1500        assert_eq!(genesis_hash, MAINNET_GENESIS_HASH);
1501    }
1502
1503    #[test]
1504    fn success_init_genesis_sepolia() {
1505        let genesis_hash =
1506            init_genesis(&create_test_provider_factory_with_chain_spec(SEPOLIA.clone())).unwrap();
1507
1508        // actual, expected
1509        assert_eq!(genesis_hash, SEPOLIA_GENESIS_HASH);
1510    }
1511
1512    #[test]
1513    fn success_init_genesis_holesky() {
1514        let genesis_hash =
1515            init_genesis(&create_test_provider_factory_with_chain_spec(HOLESKY.clone())).unwrap();
1516
1517        // actual, expected
1518        assert_eq!(genesis_hash, HOLESKY_GENESIS_HASH);
1519    }
1520
1521    #[test]
1522    fn fail_init_inconsistent_db() {
1523        let factory = create_test_provider_factory_with_chain_spec(SEPOLIA.clone());
1524        let static_file_provider = factory.static_file_provider();
1525        let rocksdb_provider = factory.rocksdb_provider();
1526        init_genesis(&factory).unwrap();
1527
1528        // Try to init db with a different genesis block
1529        let genesis_hash = init_genesis(
1530            &ProviderFactory::<MockNodeTypesWithDB>::new(
1531                factory.into_db(),
1532                MAINNET.clone(),
1533                static_file_provider,
1534                rocksdb_provider,
1535                reth_tasks::Runtime::test(),
1536            )
1537            .unwrap(),
1538        );
1539
1540        assert!(matches!(
1541            genesis_hash.unwrap_err(),
1542            InitStorageError::GenesisHashMismatch {
1543                chainspec_hash: MAINNET_GENESIS_HASH,
1544                storage_hash: SEPOLIA_GENESIS_HASH
1545            }
1546        ))
1547    }
1548
1549    #[test]
1550    fn skip_genesis_hash_validation_accepts_mismatched_db() {
1551        let factory = create_test_provider_factory_with_chain_spec(SEPOLIA.clone());
1552        let static_file_provider = factory.static_file_provider();
1553        let rocksdb_provider = factory.rocksdb_provider();
1554        init_genesis(&factory).unwrap();
1555
1556        let result = init_genesis_with_settings_and_validate(
1557            &ProviderFactory::<MockNodeTypesWithDB>::new(
1558                factory.into_db(),
1559                MAINNET.clone(),
1560                static_file_provider,
1561                rocksdb_provider,
1562                reth_tasks::Runtime::test(),
1563            )
1564            .unwrap(),
1565            StorageSettings::base(),
1566            false,
1567        );
1568
1569        let returned = result.expect("skip_genesis_validation should suppress mismatch error");
1570        assert_eq!(
1571            returned, SEPOLIA_GENESIS_HASH,
1572            "bypass returns the DB-resident hash, not the chainspec hash",
1573        );
1574    }
1575
1576    #[test]
1577    fn init_genesis_history() {
1578        let address_with_balance = Address::with_last_byte(1);
1579        let address_with_storage = Address::with_last_byte(2);
1580        let storage_key = B256::with_last_byte(1);
1581        let chain_spec = Arc::new(ChainSpec {
1582            chain: Chain::from_id(1),
1583            genesis: Genesis {
1584                alloc: BTreeMap::from([
1585                    (
1586                        address_with_balance,
1587                        GenesisAccount { balance: U256::from(1), ..Default::default() },
1588                    ),
1589                    (
1590                        address_with_storage,
1591                        GenesisAccount {
1592                            storage: Some(BTreeMap::from([(storage_key, B256::random())])),
1593                            ..Default::default()
1594                        },
1595                    ),
1596                ]),
1597                ..Default::default()
1598            },
1599            hardforks: Default::default(),
1600            paris_block_and_final_difficulty: None,
1601            deposit_contract: None,
1602            ..Default::default()
1603        });
1604
1605        let factory = create_test_provider_factory_with_chain_spec(chain_spec);
1606        init_genesis(&factory).unwrap();
1607
1608        let expected_accounts = vec![
1609            (ShardedKey::new(address_with_balance, u64::MAX), IntegerList::new([0]).unwrap()),
1610            (ShardedKey::new(address_with_storage, u64::MAX), IntegerList::new([0]).unwrap()),
1611        ];
1612        let expected_storages = vec![(
1613            StorageShardedKey::new(address_with_storage, storage_key, u64::MAX),
1614            IntegerList::new([0]).unwrap(),
1615        )];
1616
1617        let collect_from_mdbx = |factory: &ProviderFactory<MockNodeTypesWithDB>| {
1618            let provider = factory.provider().unwrap();
1619            let tx = provider.tx_ref();
1620            (
1621                collect_table_entries::<DatabaseEnv, tables::AccountsHistory>(tx).unwrap(),
1622                collect_table_entries::<DatabaseEnv, tables::StoragesHistory>(tx).unwrap(),
1623            )
1624        };
1625
1626        {
1627            let settings = factory.cached_storage_settings();
1628            let rocksdb = factory.rocksdb_provider();
1629
1630            let collect_rocksdb = |rocksdb: &reth_provider::providers::RocksDBProvider| {
1631                (
1632                    rocksdb
1633                        .iter::<tables::AccountsHistory>()
1634                        .unwrap()
1635                        .collect::<Result<Vec<_>, _>>()
1636                        .unwrap(),
1637                    rocksdb
1638                        .iter::<tables::StoragesHistory>()
1639                        .unwrap()
1640                        .collect::<Result<Vec<_>, _>>()
1641                        .unwrap(),
1642                )
1643            };
1644
1645            let (accounts, storages) = if settings.storage_v2 {
1646                collect_rocksdb(&rocksdb)
1647            } else {
1648                collect_from_mdbx(&factory)
1649            };
1650            assert_eq!(accounts, expected_accounts);
1651            assert_eq!(storages, expected_storages);
1652        }
1653    }
1654
1655    #[test]
1656    fn warn_storage_settings_mismatch() {
1657        let factory = create_test_provider_factory_with_chain_spec(MAINNET.clone());
1658        init_genesis_with_settings(&factory, StorageSettings::v1()).unwrap();
1659
1660        // Request different settings - should warn but succeed
1661        let result = init_genesis_with_settings(&factory, StorageSettings::v2());
1662
1663        // Should succeed (warning is logged, not an error)
1664        assert!(result.is_ok());
1665    }
1666
1667    #[test]
1668    fn allow_same_storage_settings() {
1669        let factory = create_test_provider_factory_with_chain_spec(MAINNET.clone());
1670        let settings = StorageSettings::v2();
1671        init_genesis_with_settings(&factory, settings).unwrap();
1672
1673        let result = init_genesis_with_settings(&factory, settings);
1674
1675        assert!(result.is_ok());
1676    }
1677}