1use 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
27pub struct ChainImportResult {
29 pub nodes: Vec<NodeHelperType<EthereumNode>>,
31 pub wallet: Wallet,
33 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
47pub 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 let mut nodes = Vec::with_capacity(num_nodes);
81 let mut temp_dirs = Vec::with_capacity(num_nodes);
83
84 for idx in 0..num_nodes {
85 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 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 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
132pub 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 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
155async 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 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 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 reth_db_common::init::init_genesis_with_settings(&provider_factory, storage_settings)?;
185
186 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 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 drop(provider_factory);
233 drop(db);
234
235 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 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 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 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 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 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 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 let fcu_state =
365 load_forkchoice_state(&fcu_path).expect("Failed to load forkchoice state");
366
367 let node = &result.nodes[0];
368
369 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 node.update_forkchoice(fcu_state.finalized_block_hash, fcu_state.head_block_hash)
381 .await
382 .expect("Failed to update forkchoice");
383
384 node.sync_to(fcu_state.head_block_hash).await.expect("Failed to sync to head");
386
387 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}