Skip to main content

reth_e2e_test_utils/
setup_import.rs

1//! Setup utilities for importing RLP chain data before starting nodes.
2
3use crate::{
4    eth_payload_attributes,
5    setup_builder::{launch_test_node, test_node_config, LaunchArgs},
6    wallet::Wallet,
7    NodeHelperType,
8};
9use eyre::WrapErr;
10use reth_chainspec::ChainSpec;
11use reth_cli_commands::import_core::{import_blocks_from_file, ImportConfig, ImportResult};
12use reth_config::Config;
13use reth_db::DatabaseEnv;
14use reth_node_api::{NodeTypesWithDBAdapter, TreeConfig};
15use reth_node_core::args::StorageArgs;
16use reth_node_ethereum::EthereumNode;
17use reth_provider::{
18    providers::{RocksDBProvider, StaticFileProvider},
19    DatabaseProviderFactory, ProviderFactory, StageCheckpointReader, StorageSettings,
20};
21use reth_stages_types::StageId;
22use reth_tasks::Runtime;
23use std::{path::Path, sync::Arc};
24use tempfile::TempDir;
25use tracing::{debug, info, span, Instrument, Level};
26
27/// Setup result containing nodes and temporary directories that must be kept alive
28pub struct ChainImportResult {
29    /// The nodes that were created
30    pub nodes: Vec<NodeHelperType<EthereumNode>>,
31    /// The wallet for testing
32    pub wallet: Wallet,
33    /// Temporary directories that must be kept alive for the duration of the test
34    pub _temp_dirs: Vec<TempDir>,
35}
36
37impl std::fmt::Debug for ChainImportResult {
38    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
39        f.debug_struct("ChainImportResult")
40            .field("nodes", &self.nodes.len())
41            .field("wallet", &self.wallet)
42            .field("temp_dirs", &self._temp_dirs.len())
43            .finish()
44    }
45}
46
47/// Creates a test setup with Ethereum nodes that have pre-imported chain data from RLP files.
48///
49/// This function:
50/// 1. Creates a temporary datadir for each node
51/// 2. Imports the specified RLP chain data into the datadir
52/// 3. Starts the nodes with the pre-populated database
53/// 4. Returns the running nodes ready for testing
54///
55/// The import and all nodes share a single [`Runtime::test`] runtime, and the nodes build payloads
56/// with [`eth_payload_attributes`] for the given chain spec. `storage_v2` selects the storage
57/// layout (`--storage.v2`) for both the imported database and the launched nodes.
58///
59/// Note: This function is currently specific to `EthereumNode` because the import process
60/// uses Ethereum-specific consensus and block format. It can be made generic in the future
61/// by abstracting the import process.
62/// It uses `NoopConsensus` during import to bypass validation checks like gas limit constraints,
63/// which allows importing test chains that may not strictly conform to mainnet consensus rules. The
64/// nodes themselves still run with proper consensus when started.
65pub async fn setup_engine_with_chain_import(
66    num_nodes: usize,
67    chain_spec: Arc<ChainSpec>,
68    is_dev: bool,
69    storage_v2: bool,
70    tree_config: TreeConfig,
71    rlp_path: &Path,
72) -> eyre::Result<ChainImportResult> {
73    let runtime = Runtime::test();
74    let attributes_generator = {
75        let chain_spec = chain_spec.clone();
76        Arc::new(move |timestamp| eth_payload_attributes(&chain_spec, timestamp))
77    };
78
79    // Create nodes with imported data.
80    let mut nodes = Vec::with_capacity(num_nodes);
81    // Keep the temp dirs alive for the lifetime of the nodes.
82    let mut temp_dirs = Vec::with_capacity(num_nodes);
83
84    for idx in 0..num_nodes {
85        // Create a temporary datadir for this node.
86        let temp_dir = TempDir::new()?;
87        let datadir = temp_dir.path().to_path_buf();
88        debug!(target: "e2e::import", "Node {idx} datadir: {datadir:?}");
89
90        let span = span!(Level::INFO, "node", idx);
91        let node_config = test_node_config(chain_spec.clone())
92            .set_dev(is_dev)
93            .with_storage(StorageArgs { v2: storage_v2 });
94
95        // First, import the chain data into this datadir.
96        import_chain(
97            &datadir,
98            chain_spec.clone(),
99            node_config.storage_settings(),
100            rlp_path,
101            runtime.clone(),
102        )
103        .instrument(span.clone())
104        .await
105        .wrap_err_with(|| format!("chain import failed for node {idx}"))?;
106
107        // Now launch the node with the pre-populated datadir.
108        debug!(target: "e2e::import", "Launching node with datadir: {:?}", datadir);
109
110        let node = launch_test_node::<EthereumNode>(LaunchArgs {
111            node_config,
112            runtime: runtime.clone(),
113            tree_config: tree_config.clone(),
114            datadir,
115            attributes_generator: attributes_generator.clone(),
116            dev_payload_attributes: None,
117        })
118        .instrument(span)
119        .await?;
120
121        nodes.push(node);
122        temp_dirs.push(temp_dir);
123    }
124
125    Ok(ChainImportResult {
126        nodes,
127        wallet: Wallet::default().with_chain_id(chain_spec.chain.id()),
128        _temp_dirs: temp_dirs,
129    })
130}
131
132/// Helper to load forkchoice state from a JSON file
133pub fn load_forkchoice_state(path: &Path) -> eyre::Result<alloy_rpc_types_engine::ForkchoiceState> {
134    let json_str = std::fs::read_to_string(path)?;
135    let fcu_data: serde_json::Value = serde_json::from_str(&json_str)?;
136
137    // The headfcu.json file contains a JSON-RPC request with the forkchoice state in params[0]
138    let state = &fcu_data["params"][0];
139    Ok(alloy_rpc_types_engine::ForkchoiceState {
140        head_block_hash: state["headBlockHash"]
141            .as_str()
142            .ok_or_else(|| eyre::eyre!("missing headBlockHash"))?
143            .parse()?,
144        safe_block_hash: state["safeBlockHash"]
145            .as_str()
146            .ok_or_else(|| eyre::eyre!("missing safeBlockHash"))?
147            .parse()?,
148        finalized_block_hash: state["finalizedBlockHash"]
149            .as_str()
150            .ok_or_else(|| eyre::eyre!("missing finalizedBlockHash"))?
151            .parse()?,
152    })
153}
154
155/// Initializes the database in `datadir` with the given storage settings and imports the RLP
156/// encoded chain at `rlp_path`.
157///
158/// All database handles are released on return so the datadir can be reopened, e.g. by a node,
159/// which reads the storage settings back from the database.
160async fn import_chain(
161    datadir: &Path,
162    chain_spec: Arc<ChainSpec>,
163    storage_settings: StorageSettings,
164    rlp_path: &Path,
165    runtime: Runtime,
166) -> eyre::Result<ImportResult> {
167    info!(target: "test", "Importing chain data from {:?} into {:?}", rlp_path, datadir);
168
169    // Initialize the database using init_db, same as the CLI import command.
170    let db_args = reth_node_core::args::DatabaseArgs::default().database_args();
171    let db = reth_db::init_db(datadir.join("db"), db_args)?;
172
173    // Create a provider factory with the initialized database, not a `TempDatabase`.
174    let provider_factory =
175        ProviderFactory::<NodeTypesWithDBAdapter<EthereumNode, DatabaseEnv>>::new(
176            db.clone(),
177            chain_spec.clone(),
178            StaticFileProvider::read_write(datadir.join("static_files"))?,
179            RocksDBProvider::builder(datadir.join("rocksdb")).with_default_tables().build()?,
180            runtime.clone(),
181        )?;
182
183    // Initialize genesis with the storage settings of the node that later opens the database.
184    reth_db_common::init::init_genesis_with_settings(&provider_factory, storage_settings)?;
185
186    // Use NoopConsensus to skip gas limit validation for test imports.
187    let result = import_blocks_from_file(
188        rlp_path,
189        ImportConfig::default(),
190        provider_factory.clone(),
191        &Config::default(),
192        reth_node_ethereum::EthEvmConfig::new(chain_spec),
193        reth_consensus::noop::NoopConsensus::arc(),
194        runtime,
195    )
196    .await?;
197
198    info!(
199        target: "test",
200        "Imported {} blocks and {} transactions",
201        result.total_imported_blocks,
202        result.total_imported_txns,
203    );
204
205    debug!(target: "e2e::import",
206        "Import result: decoded {} blocks, imported {} blocks, complete: {}",
207        result.total_decoded_blocks,
208        result.total_imported_blocks,
209        result.is_complete()
210    );
211
212    eyre::ensure!(
213        result.total_decoded_blocks == result.total_imported_blocks,
214        "block count mismatch: decoded {} != imported {}",
215        result.total_decoded_blocks,
216        result.total_imported_blocks
217    );
218    eyre::ensure!(
219        result.total_decoded_txns == result.total_imported_txns,
220        "transaction count mismatch: decoded {} != imported {}",
221        result.total_decoded_txns,
222        result.total_imported_txns
223    );
224
225    // Verify the database was properly initialized by checking stage checkpoints.
226    let headers_checkpoint =
227        provider_factory.database_provider_ro()?.get_stage_checkpoint(StageId::Headers)?;
228    eyre::ensure!(headers_checkpoint.is_some(), "Headers stage checkpoint is missing after import");
229    debug!(target: "e2e::import", "Headers stage checkpoint after import: {headers_checkpoint:?}");
230
231    // Close all database handles to release locks before launching the node.
232    drop(provider_factory);
233    drop(db);
234
235    // The header and body downloader tasks spawned on the runtime hold provider factory clones
236    // and only exit once they are polled after the pipeline was dropped, so yield long enough for
237    // them to release the database.
238    tokio::time::sleep(std::time::Duration::from_millis(100)).await;
239
240    Ok(result)
241}
242
243#[cfg(test)]
244mod tests {
245    use super::*;
246    use crate::{
247        test_chain_spec,
248        test_rlp_utils::{create_fcu_json, generate_test_blocks, write_blocks_to_rlp},
249    };
250    use reth_chainspec::EthereumHardfork;
251    use reth_db::mdbx::DatabaseArguments;
252    use reth_ethereum_primitives::Block;
253    use reth_primitives_traits::SealedBlock;
254    use reth_provider::{BlockHashReader, BlockNumReader, BlockReaderIdExt, MetadataProvider};
255    use std::path::PathBuf;
256
257    /// Helper to setup test blocks and write to RLP.
258    fn setup_test_blocks_and_rlp(
259        chain_spec: &ChainSpec,
260        block_count: u64,
261        temp_dir: &Path,
262    ) -> (Vec<SealedBlock<Block>>, PathBuf) {
263        let test_blocks = generate_test_blocks(chain_spec, block_count);
264        assert_eq!(
265            test_blocks.len(),
266            block_count as usize,
267            "Should have generated expected blocks"
268        );
269
270        let rlp_path = temp_dir.join("test_chain.rlp");
271        write_blocks_to_rlp(&test_blocks, &rlp_path).expect("Failed to write RLP data");
272
273        let rlp_size = std::fs::metadata(&rlp_path).expect("RLP file should exist").len();
274        debug!(target: "e2e::import", "Wrote RLP file with size: {rlp_size} bytes");
275
276        (test_blocks, rlp_path)
277    }
278
279    /// Reopens the database in `datadir` that was written by [`import_chain`].
280    fn reopen_provider_factory(
281        datadir: &Path,
282        chain_spec: Arc<ChainSpec>,
283        runtime: Runtime,
284    ) -> ProviderFactory<NodeTypesWithDBAdapter<EthereumNode, DatabaseEnv>> {
285        let db = reth_db::init_db(datadir.join("db"), DatabaseArguments::default()).unwrap();
286        ProviderFactory::new(
287            db,
288            chain_spec,
289            StaticFileProvider::read_only(datadir.join("static_files")).unwrap(),
290            RocksDBProvider::builder(datadir.join("rocksdb"))
291                .with_default_tables()
292                .build()
293                .unwrap(),
294            runtime,
295        )
296        .expect("failed to create provider factory")
297    }
298
299    #[tokio::test]
300    async fn test_import_blocks_and_reopen() {
301        // Tests the block import without a node, and that the imported chain and stage
302        // checkpoints persist when reopening the database.
303        reth_tracing::init_test_tracing();
304
305        let runtime = Runtime::test();
306        let chain_spec = test_chain_spec(EthereumHardfork::Shanghai);
307        let temp_dir = tempfile::tempdir().expect("Failed to create temp dir");
308        let (test_blocks, rlp_path) = setup_test_blocks_and_rlp(&chain_spec, 10, temp_dir.path());
309
310        let datadir = temp_dir.path().join("datadir");
311        std::fs::create_dir_all(&datadir).unwrap();
312        let storage_settings = StorageSettings { storage_v2: StorageArgs::default().v2 };
313        let result = import_chain(
314            &datadir,
315            chain_spec.clone(),
316            storage_settings,
317            &rlp_path,
318            runtime.clone(),
319        )
320        .await
321        .unwrap();
322        assert_eq!(result.total_decoded_blocks, 10);
323        assert_eq!(result.total_imported_blocks, 10);
324        assert_eq!(result.total_decoded_txns, 0);
325        assert_eq!(result.total_imported_txns, 0);
326
327        let provider_factory = reopen_provider_factory(&datadir, chain_spec, runtime);
328        let provider = provider_factory.database_provider_ro().unwrap();
329        assert_eq!(provider.last_block_number().unwrap(), 10);
330        assert_eq!(provider.block_hash(10).unwrap(), Some(test_blocks[9].hash()));
331        let headers_checkpoint = provider.get_stage_checkpoint(StageId::Headers).unwrap();
332        assert_eq!(headers_checkpoint.map(|checkpoint| checkpoint.block_number), Some(10));
333    }
334
335    #[tokio::test]
336    async fn test_import_with_node_integration() {
337        // Tests the full integration with node setup, forkchoice updates, and syncing.
338        reth_tracing::init_test_tracing();
339
340        let chain_spec = test_chain_spec(EthereumHardfork::Shanghai);
341        let temp_dir = tempfile::tempdir().expect("Failed to create temp dir");
342        let (test_blocks, rlp_path) = setup_test_blocks_and_rlp(&chain_spec, 10, temp_dir.path());
343
344        // Create FCU data for the tip.
345        let tip = test_blocks.last().expect("Should have generated blocks");
346        let fcu_path = temp_dir.path().join("test_fcu.json");
347        std::fs::write(&fcu_path, create_fcu_json(tip).to_string())
348            .expect("Failed to write FCU data");
349
350        for storage_v2 in [false, true] {
351            // Setup nodes with imported chain.
352            let result = setup_engine_with_chain_import(
353                1,
354                chain_spec.clone(),
355                false,
356                storage_v2,
357                TreeConfig::default(),
358                &rlp_path,
359            )
360            .await
361            .expect("Failed to setup nodes with chain import");
362
363            // Load and apply forkchoice state.
364            let fcu_state =
365                load_forkchoice_state(&fcu_path).expect("Failed to load forkchoice state");
366
367            let node = &result.nodes[0];
368
369            // The imported database and the node use the requested storage layout.
370            let settings = node
371                .inner
372                .provider
373                .database_provider_ro()
374                .expect("Failed to open database provider")
375                .storage_settings()
376                .expect("Failed to read storage settings");
377            assert_eq!(settings, Some(StorageSettings { storage_v2 }));
378
379            // Send forkchoice update to make the imported chain canonical.
380            node.update_forkchoice(fcu_state.finalized_block_hash, fcu_state.head_block_hash)
381                .await
382                .expect("Failed to update forkchoice");
383
384            // Wait for the node to sync to the head.
385            node.sync_to(fcu_state.head_block_hash).await.expect("Failed to sync to head");
386
387            // Verify the chain tip.
388            let latest = node
389                .inner
390                .provider
391                .sealed_header_by_id(alloy_eips::BlockId::latest())
392                .expect("Failed to get latest header")
393                .expect("No latest header found");
394
395            assert_eq!(
396                latest.hash(),
397                fcu_state.head_block_hash,
398                "Chain tip does not match expected head"
399            );
400        }
401    }
402}