client/common/
mod.rs

1// Copyright 2020 Conflux Foundation. All rights reserved.
2// Conflux is free software and distributed under GNU General Public License.
3// See http://www.gnu.org/licenses/
4
5use log::{debug, info};
6use std::{
7    collections::HashMap,
8    fs::create_dir_all,
9    path::Path,
10    sync::{Arc, Weak},
11    thread,
12    time::{Duration, Instant},
13};
14
15use cfx_rpc_builder::RpcServerHandle;
16use cfx_util_macros::bail;
17use parking_lot::{Condvar, Mutex};
18use rand_08::{prelude::StdRng, rngs::OsRng, SeedableRng};
19use threadpool::ThreadPool;
20
21use blockgen::BlockGenerator;
22use cfx_executor::machine::{Machine, VmFactory};
23use cfx_parameters::genesis::{
24    DEV_GENESIS_KEY_PAIR_2, GENESIS_ACCOUNT_ADDRESS,
25};
26use cfx_rpc_cfx_types::apis::ApiSet;
27use cfx_storage::StorageManager;
28use cfx_tasks::TaskManager;
29use cfx_types::{address_util::AddressUtil, Address, Space, U256};
30pub use cfxcore::pos::pos::PosDropHandle;
31use cfxcore::{
32    block_data_manager::BlockDataManager,
33    consensus::{
34        pivot_hint::PivotHint,
35        pos_handler::{PosConfiguration, PosVerifier},
36    },
37    genesis_block::{self as genesis, genesis_block},
38    pow::PowComputer,
39    statistics::Statistics,
40    sync::SyncPhaseType,
41    ConsensusGraph, LightProvider, NodeType, Notifications, Stopable,
42    SynchronizationGraph, SynchronizationService, TransactionPool,
43    WORKER_COMPUTATION_PARALLELISM,
44};
45use cfxcore_accounts::AccountProvider;
46use cfxkey::public_to_address;
47use diem_config::keys::ConfigKey;
48use diem_crypto::{
49    key_file::{load_pri_key, save_pri_key},
50    PrivateKey, Uniform,
51};
52use diem_types::validator_config::{
53    ConsensusPrivateKey, ConsensusVRFPrivateKey,
54};
55use malloc_size_of::{new_malloc_size_ops, MallocSizeOf, MallocSizeOfOps};
56use network::NetworkService;
57use secret_store::{SecretStore, SharedSecretStore};
58use tokio::runtime::Runtime as TokioRuntime;
59use txgen::{DirectTransactionGenerator, TransactionGenerator};
60
61use cfx_config::{parse_config_address_string, Configuration};
62
63use crate::{
64    accounts::{account_provider, keys_path},
65    keylib::KeyPair,
66    rpc_starter::{launch_async_rpc_servers, launch_cfx_async_rpc_servers},
67};
68#[cfg(all(unix, feature = "jemalloc-prof"))]
69use cfx_mallocator_utils::start_pprf_server;
70use cfxcore::consensus::pos_handler::read_initial_nodes_from_file;
71
72pub mod panic_handler;
73pub mod shutdown_handler;
74
75/// Hold all top-level components for a type of client.
76/// This struct implement ClientShutdownTrait.
77pub struct ClientComponents<BlockGenT, Rest> {
78    pub data_manager_weak_ptr: Weak<BlockDataManager>,
79    pub blockgen: Option<Arc<BlockGenT>>,
80    pub pos_handler: Option<Arc<PosVerifier>>,
81    pub other_components: Rest,
82}
83
84impl<BlockGenT, Rest: MallocSizeOf> MallocSizeOf
85    for ClientComponents<BlockGenT, Rest>
86{
87    fn size_of(&self, ops: &mut MallocSizeOfOps) -> usize {
88        if let Some(data_man) = self.data_manager_weak_ptr.upgrade() {
89            let data_manager_size = data_man.size_of(ops);
90            data_manager_size + self.other_components.size_of(ops)
91        } else {
92            // If data_man is `None`, we will be just shutting down (dropping
93            // components) so we don't care about the size
94            0
95        }
96    }
97}
98
99impl<BlockGenT: 'static + Stopable, Rest> ClientTrait
100    for ClientComponents<BlockGenT, Rest>
101{
102    fn take_out_components_for_shutdown(
103        &self,
104    ) -> (
105        Weak<BlockDataManager>,
106        Option<Arc<PosVerifier>>,
107        Option<Arc<dyn Stopable>>,
108    ) {
109        debug!("take_out_components_for_shutdown");
110        let data_manager_weak_ptr = self.data_manager_weak_ptr.clone();
111        let blockgen: Option<Arc<dyn Stopable>> = match self.blockgen.clone() {
112            Some(blockgen) => Some(blockgen),
113            None => None,
114        };
115
116        (data_manager_weak_ptr, self.pos_handler.clone(), blockgen)
117    }
118}
119
120pub trait ClientTrait {
121    fn take_out_components_for_shutdown(
122        &self,
123    ) -> (
124        Weak<BlockDataManager>,
125        Option<Arc<PosVerifier>>,
126        Option<Arc<dyn Stopable>>,
127    );
128}
129
130pub fn initialize_common_modules(
131    conf: &mut Configuration, exit: Arc<(Mutex<bool>, Condvar)>,
132    node_type: NodeType,
133) -> Result<
134    (
135        Arc<Machine>,
136        Arc<SecretStore>,
137        HashMap<Address, U256>,
138        Arc<BlockDataManager>,
139        Arc<PowComputer>,
140        Arc<PosVerifier>,
141        Arc<TransactionPool>,
142        Arc<ConsensusGraph>,
143        Arc<SynchronizationGraph>,
144        Arc<NetworkService>,
145        Arc<AccountProvider>,
146        Arc<Notifications>,
147        Arc<TokioRuntime>,
148    ),
149    String,
150> {
151    info!("Working directory: {:?}", std::env::current_dir());
152
153    // TODO(lpl): Keep it properly and allow not running pos.
154    let (self_pos_private_key, self_vrf_private_key) = {
155        let key_path = Path::new(&conf.raw_conf.pos_private_key_path);
156
157        let read_pos_password = |prompt: &str| -> Result<Vec<u8>, String> {
158            match rpassword::prompt_password(prompt) {
159                Ok(password) => Ok(password.into_bytes()),
160                Err(e) => {
161                    let mut msg = format!("{:?}", e);
162                    // On macOS, attempting to open `/dev/tty` without a
163                    // controlling TTY can return ENXIO
164                    // ("Device not configured"). This commonly happens when
165                    // running under a debugger/IDE that doesn't allocate a real
166                    // terminal.
167                    if e.raw_os_error() == Some(6) {
168                        msg.push_str(" Hint: maybe no controlling TTY detected (macOS ENXIO: \"Device not configured\"). If you are running under VS Code debugger or with redirected stdio, set `CFX_POS_KEY_ENCRYPTION_PASSWORD` env var to avoid interactive prompting.");
169                    }
170                    Err(msg)
171                }
172            }
173        };
174
175        let default_passwd = if conf.is_test_or_dev_mode() {
176            Some(vec![])
177        } else {
178            conf.raw_conf
179                .dev_pos_private_key_encryption_password
180                .clone()
181                // If the password is not set in the config file, read it from
182                // the environment variable.
183                .or(std::env::var("CFX_POS_KEY_ENCRYPTION_PASSWORD").ok())
184                .map(|s| s.into_bytes())
185        };
186        if key_path.exists() {
187            let passwd = match default_passwd {
188                Some(p) => p,
189                None => read_pos_password(
190                    "PoS key detected, please input your encryption password.\nPassword:",
191                )?,
192            };
193            match load_pri_key(key_path, &passwd) {
194                Ok((sk, vrf_sk)) => {
195                    (ConfigKey::new(sk), ConfigKey::new(vrf_sk))
196                }
197                Err(e) => {
198                    bail!("Load pos_key failed: {}", e);
199                }
200            }
201        } else {
202            create_dir_all(key_path.parent().unwrap()).unwrap();
203            let passwd = match default_passwd {
204                Some(p) => p,
205                None => {
206                    let p = read_pos_password("PoS key is not detected and will be generated instead, please input your encryption password. This password is needed when you restart the node\nPassword:")?;
207                    let p2 = read_pos_password("Repeat Password:")?;
208                    if p != p2 {
209                        bail!("Passwords do not match!");
210                    }
211                    p
212                }
213            };
214            let mut rng = StdRng::from_rng(OsRng).unwrap();
215            let private_key = ConsensusPrivateKey::generate(&mut rng);
216            let vrf_private_key = ConsensusVRFPrivateKey::generate(&mut rng);
217            save_pri_key(key_path, &passwd, &(&private_key, &vrf_private_key))
218                .expect("error saving private key");
219            (ConfigKey::new(private_key), ConfigKey::new(vrf_private_key))
220        }
221    };
222
223    let worker_thread_pool = Arc::new(Mutex::new(ThreadPool::with_name(
224        "Tx Recover".into(),
225        WORKER_COMPUTATION_PARALLELISM,
226    )));
227
228    let network_config = conf.net_config()?;
229    let cache_config = conf.cache_config();
230
231    let (db_path, db_config) = conf.db_config();
232    let ledger_db = db::open_database(db_path.to_str().unwrap(), &db_config)
233        .map_err(|e| format!("Failed to open database {:?}", e))?;
234
235    let secret_store = Arc::new(SecretStore::new());
236    let storage_manager = Arc::new(
237        StorageManager::new(conf.storage_config(&node_type))
238            .expect("Failed to initialize storage."),
239    );
240    {
241        let storage_manager_log_weak_ptr = Arc::downgrade(&storage_manager);
242        let exit_clone = exit.clone();
243        thread::spawn(move || loop {
244            let mut exit_lock = exit_clone.0.lock();
245            if exit_clone
246                .1
247                .wait_for(&mut exit_lock, Duration::from_millis(5000))
248                .timed_out()
249            {
250                let manager = storage_manager_log_weak_ptr.upgrade();
251                match manager {
252                    None => return,
253                    Some(manager) => manager.log_usage(),
254                };
255            } else {
256                return;
257            }
258        });
259    }
260
261    let genesis_accounts = if conf.is_test_or_dev_mode() {
262        match (
263            &conf.raw_conf.genesis_secrets,
264            &conf.raw_conf.genesis_evm_secrets,
265        ) {
266            (Some(file), evm_file) => {
267                // Load core space accounts
268                let mut accounts = genesis::load_secrets_file(
269                    file,
270                    secret_store.as_ref(),
271                    Space::Native,
272                )?;
273
274                // Load EVM space accounts if specified
275                if let Some(evm_file) = evm_file {
276                    let evm_accounts = genesis::load_secrets_file(
277                        evm_file,
278                        secret_store.as_ref(),
279                        Space::Ethereum,
280                    )?;
281                    accounts.extend(evm_accounts);
282                }
283                accounts
284            }
285            (None, Some(evm_file)) => {
286                // Only load EVM space accounts
287                genesis::load_secrets_file(
288                    evm_file,
289                    secret_store.as_ref(),
290                    Space::Ethereum,
291                )?
292            }
293            (None, None) => genesis::default(conf.is_test_or_dev_mode()),
294        }
295    } else {
296        match conf.raw_conf.genesis_accounts {
297            Some(ref file) => genesis::load_file(file, |addr_str| {
298                parse_config_address_string(
299                    addr_str,
300                    network_config.get_network_type(),
301                )
302            })?,
303            None => genesis::default(conf.is_test_or_dev_mode()),
304        }
305    };
306
307    // Only try to setup PoW genesis block if pos is enabled from genesis.
308    let initial_nodes = if conf.raw_conf.pos_reference_enable_height == 0 {
309        Some(
310            read_initial_nodes_from_file(
311                conf.raw_conf.pos_initial_nodes_path.as_str(),
312            )
313            .expect("Genesis must have been initialized with pos"),
314        )
315    } else {
316        None
317    };
318
319    let consensus_conf = conf.consensus_config();
320    let vm = VmFactory::new(1024 * 32);
321    let machine = Arc::new(Machine::new_with_builtin(conf.common_params(), vm));
322
323    let genesis_block = genesis_block(
324        &storage_manager,
325        genesis_accounts.clone(),
326        GENESIS_ACCOUNT_ADDRESS,
327        U256::zero(),
328        machine.clone(),
329        conf.raw_conf.execute_genesis, /* need_to_execute */
330        conf.raw_conf.chain_id,
331        &initial_nodes,
332    );
333    storage_manager.notify_genesis_hash(genesis_block.hash());
334    let mut genesis_accounts = genesis_accounts;
335    let genesis_accounts = genesis_accounts
336        .drain()
337        .filter(|(addr, _)| addr.space == Space::Native)
338        .map(|(addr, x)| (addr.address, x))
339        .collect();
340    debug!("Initialize genesis_block={:?}", genesis_block);
341    if conf.raw_conf.pos_genesis_pivot_decision.is_none() {
342        conf.raw_conf.pos_genesis_pivot_decision = Some(genesis_block.hash());
343    }
344
345    let pow_config = conf.pow_config();
346    let pow = Arc::new(PowComputer::new(pow_config.use_octopus()));
347
348    let data_man = Arc::new(BlockDataManager::new(
349        cache_config,
350        Arc::new(genesis_block),
351        ledger_db.clone(),
352        storage_manager,
353        worker_thread_pool,
354        conf.data_mananger_config(),
355        pow.clone(),
356    ));
357
358    let network = {
359        let mut rng = StdRng::from_rng(OsRng).unwrap();
360        let private_key = ConsensusPrivateKey::generate(&mut rng);
361        let vrf_private_key = ConsensusVRFPrivateKey::generate(&mut rng);
362        let mut network = NetworkService::new(network_config.clone());
363        network
364            .initialize((
365                private_key.public_key(),
366                vrf_private_key.public_key(),
367            ))
368            .unwrap();
369        Arc::new(network)
370    };
371
372    let pos_verifier = Arc::new(PosVerifier::new(
373        Some(network.clone()),
374        PosConfiguration {
375            bls_key: self_pos_private_key,
376            vrf_key: self_vrf_private_key,
377            diem_conf_path: conf.raw_conf.pos_config_path.clone(),
378            protocol_conf: conf.protocol_config(),
379            pos_initial_nodes_path: conf
380                .raw_conf
381                .pos_initial_nodes_path
382                .clone(),
383            vrf_proposal_threshold: conf.raw_conf.vrf_proposal_threshold,
384            pos_state_config: conf.pos_state_config(),
385        },
386        conf.raw_conf.pos_reference_enable_height,
387    ));
388    let verification_config = conf.verification_config(machine.clone());
389    let txpool = Arc::new(TransactionPool::new(
390        conf.txpool_config(),
391        verification_config.clone(),
392        data_man.clone(),
393        machine.clone(),
394    ));
395
396    let statistics = Arc::new(Statistics::new());
397    let notifications = Notifications::init();
398    let pivot_hint = if let Some(conf) = &consensus_conf.pivot_hint_conf {
399        Some(Arc::new(PivotHint::new(conf)?))
400    } else {
401        None
402    };
403
404    let consensus = Arc::new(ConsensusGraph::new(
405        consensus_conf,
406        txpool.clone(),
407        statistics.clone(),
408        data_man.clone(),
409        pow_config.clone(),
410        pow.clone(),
411        notifications.clone(),
412        conf.execution_config(),
413        verification_config.clone(),
414        node_type,
415        pos_verifier.clone(),
416        pivot_hint,
417        conf.common_params(),
418    ));
419
420    for terminal in data_man
421        .terminals_from_db()
422        .unwrap_or(vec![data_man.get_cur_consensus_era_genesis_hash()])
423    {
424        if data_man.block_height_by_hash(&terminal).unwrap()
425            >= conf.raw_conf.pos_reference_enable_height
426        {
427            pos_verifier.initialize(consensus.clone())?;
428            break;
429        }
430    }
431
432    let sync_config = conf.sync_graph_config();
433
434    let sync_graph = Arc::new(SynchronizationGraph::new(
435        consensus.clone(),
436        data_man.clone(),
437        statistics.clone(),
438        verification_config,
439        pow_config,
440        pow.clone(),
441        sync_config,
442        notifications.clone(),
443        machine.clone(),
444        pos_verifier.clone(),
445    ));
446    let refresh_time =
447        Duration::from_millis(conf.raw_conf.account_provider_refresh_time_ms);
448
449    let accounts = Arc::new(
450        account_provider(
451            Some(keys_path()),
452            None, /* sstore_iterations */
453            Some(refresh_time),
454        )
455        .expect("failed to initialize account provider"),
456    );
457
458    let tokio_runtime =
459        Arc::new(TokioRuntime::new().map_err(|e| e.to_string())?);
460
461    Ok((
462        machine,
463        secret_store,
464        genesis_accounts,
465        data_man,
466        pow,
467        pos_verifier,
468        txpool,
469        consensus,
470        sync_graph,
471        network,
472        accounts,
473        notifications,
474        tokio_runtime,
475    ))
476}
477
478pub fn initialize_not_light_node_modules(
479    conf: &mut Configuration, exit: Arc<(Mutex<bool>, Condvar)>,
480    node_type: NodeType,
481) -> Result<
482    (
483        Arc<BlockDataManager>,
484        Arc<PowComputer>,
485        Arc<TransactionPool>,
486        Arc<ConsensusGraph>,
487        Arc<SynchronizationService>,
488        Arc<BlockGenerator>,
489        Arc<PosVerifier>,
490        Arc<TokioRuntime>,
491        Option<RpcServerHandle>,
492        Option<RpcServerHandle>,
493        Option<RpcServerHandle>,
494        TaskManager,
495    ),
496    String,
497> {
498    let (
499        _machine,
500        secret_store,
501        genesis_accounts,
502        data_man,
503        pow,
504        pos_verifier,
505        txpool,
506        consensus,
507        sync_graph,
508        network,
509        accounts,
510        notifications,
511        tokio_runtime,
512    ) = initialize_common_modules(conf, exit.clone(), node_type)?;
513
514    let light_provider = Arc::new(LightProvider::new(
515        consensus.clone(),
516        sync_graph.clone(),
517        Arc::downgrade(&network),
518        txpool.clone(),
519        conf.raw_conf.throttling_conf.clone(),
520        node_type,
521    ));
522    light_provider.register(network.clone()).unwrap();
523
524    let sync = Arc::new(SynchronizationService::new(
525        node_type,
526        network.clone(),
527        sync_graph.clone(),
528        conf.protocol_config(),
529        conf.state_sync_config(),
530        SyncPhaseType::CatchUpRecoverBlockHeaderFromDB,
531        light_provider,
532        consensus.clone(),
533    ));
534    sync.register().unwrap();
535
536    if let Some(print_memory_usage_period_s) =
537        conf.raw_conf.print_memory_usage_period_s
538    {
539        let secret_store = secret_store.clone();
540        let data_man = data_man.clone();
541        let txpool = txpool.clone();
542        let consensus = consensus.clone();
543        let sync = sync.clone();
544        thread::Builder::new().name("MallocSizeOf".into()).spawn(
545            move || loop {
546                let start = Instant::now();
547                let mb = 1_000_000;
548                let mut ops = new_malloc_size_ops();
549                let secret_store_size = secret_store.size_of(&mut ops) / mb;
550                // Note `db_manager` is not wrapped in Arc, so it will still be included
551                // in `data_man_size`.
552                let data_manager_db_cache_size = data_man.db_manager.size_of(&mut ops) / mb;
553                let storage_manager_size = data_man.storage_manager.size_of(&mut ops) / mb;
554                let data_man_size = data_man.size_of(&mut ops) / mb;
555                let tx_pool_size = txpool.size_of(&mut ops) / mb;
556                let consensus_graph_size = consensus.size_of(&mut ops) / mb;
557                let sync_graph_size =
558                    sync.get_synchronization_graph().size_of(&mut ops) / mb;
559                let sync_service_size = sync.size_of(&mut ops) / mb;
560                info!(
561                    "Malloc Size(MB): secret_store={} data_manager_db_cache_size={} \
562                    storage_manager_size={} data_man={} txpool={} consensus={} sync_graph={}\
563                    sync_service={}, \
564                    time elapsed={:?}",
565                    secret_store_size,data_manager_db_cache_size,storage_manager_size,
566                    data_man_size, tx_pool_size, consensus_graph_size, sync_graph_size,
567                    sync_service_size, start.elapsed(),
568                );
569                thread::sleep(Duration::from_secs(
570                    print_memory_usage_period_s,
571                ));
572            },
573        ).expect("Memory usage thread start fails");
574    }
575
576    let (maybe_txgen, maybe_direct_txgen) = initialize_txgens(
577        consensus.clone(),
578        txpool.clone(),
579        sync.clone(),
580        secret_store.clone(),
581        genesis_accounts,
582        conf,
583        network.net_key_pair().unwrap(),
584    );
585
586    let maybe_author: Option<Address> =
587        conf.raw_conf.mining_author.as_ref().map(|addr_str| {
588            parse_config_address_string(addr_str, network.get_network_type())
589                .unwrap_or_else(|err| {
590                    panic!("Error parsing mining-author {}", err)
591                })
592        });
593    let pow_config = conf.pow_config();
594    let blockgen = Arc::new(BlockGenerator::new(
595        sync_graph,
596        txpool.clone(),
597        sync.clone(),
598        maybe_txgen.clone(),
599        pow_config.clone(),
600        pow.clone(),
601        maybe_author.unwrap_or_default(),
602        pos_verifier.clone(),
603    ));
604    if conf.is_dev_mode() {
605        // If `dev_block_interval_ms` is None, blocks are generated after
606        // receiving RPC `cfx_sendRawTransaction`.
607        if let Some(interval_ms) = conf.raw_conf.dev_block_interval_ms {
608            // Automatic block generation with fixed interval.
609            let bg = blockgen.test_api();
610            info!("Start auto block generation");
611            thread::Builder::new()
612                .name("auto_mining".into())
613                .spawn(move || {
614                    bg.auto_block_generation(interval_ms);
615                })
616                .expect("Mining thread spawn error");
617        }
618    } else if let Some(author) = maybe_author {
619        if !author.is_genesis_valid_address() || author.is_builtin_address() {
620            panic!("mining-author must be user address or contract address, otherwise you will not get mining rewards!!!");
621        }
622        if pow_config.enable_mining() {
623            let bg = blockgen.clone();
624            thread::Builder::new()
625                .name("mining".into())
626                .spawn(move || {
627                    bg.mine();
628                })
629                .expect("Mining thread spawn error");
630        }
631    }
632
633    info!("Using new async core space RPC implementation");
634
635    let task_manager = TaskManager::new(tokio_runtime.handle().clone());
636    let task_executor = task_manager.executor();
637
638    let eth_rpc_server_handle =
639        tokio_runtime.block_on(launch_async_rpc_servers(
640            consensus.clone(),
641            sync.clone(),
642            txpool.clone(),
643            notifications.clone(),
644            task_executor.clone(),
645            conf,
646        ))?;
647
648    // Start V2 async core space RPC servers.
649    let cfx_rpc_server_handle =
650        tokio_runtime.block_on(launch_cfx_async_rpc_servers(
651            consensus.clone(),
652            sync.clone(),
653            txpool.clone(),
654            data_man.clone(),
655            network.clone(),
656            pos_verifier.clone(),
657            notifications.clone(),
658            task_executor.clone(),
659            accounts.clone(),
660            exit.clone(),
661            blockgen.test_api(),
662            maybe_txgen.clone(),
663            maybe_direct_txgen.clone(),
664            conf,
665            conf.raw_conf.public_rpc_apis.clone(),
666            false,
667        ))?;
668
669    let debug_cfx_rpc_server_handle =
670        tokio_runtime.block_on(launch_cfx_async_rpc_servers(
671            consensus.clone(),
672            sync.clone(),
673            txpool.clone(),
674            data_man.clone(),
675            network.clone(),
676            pos_verifier.clone(),
677            notifications.clone(),
678            task_executor.clone(),
679            accounts.clone(),
680            exit.clone(),
681            blockgen.test_api(),
682            maybe_txgen.clone(),
683            maybe_direct_txgen.clone(),
684            conf,
685            ApiSet::All,
686            true,
687        ))?;
688
689    // start pprf server, which is used to serve the pprof data for heap
690    // profiling
691    #[cfg(all(unix, feature = "jemalloc-prof"))]
692    if let Some(pprf_addr) = conf.raw_conf.profiling_listen_addr.as_ref() {
693        let pprf_addr = pprf_addr.clone();
694        let _pprf_server_handle = tokio_runtime.spawn(async move {
695            if let Err(e) = start_pprf_server(&pprf_addr).await {
696                eprintln!("Error starting pprof server: {}", e);
697            }
698        });
699    }
700
701    metrics::initialize(conf.metrics_config(), task_executor.clone());
702
703    network.start();
704
705    Ok((
706        data_man,
707        pow,
708        txpool,
709        consensus,
710        sync,
711        blockgen,
712        pos_verifier,
713        tokio_runtime,
714        eth_rpc_server_handle,
715        cfx_rpc_server_handle,
716        debug_cfx_rpc_server_handle,
717        task_manager,
718    ))
719}
720
721pub fn initialize_txgens(
722    consensus: Arc<ConsensusGraph>, txpool: Arc<TransactionPool>,
723    sync: Arc<SynchronizationService>, secret_store: SharedSecretStore,
724    genesis_accounts: HashMap<Address, U256>, conf: &Configuration,
725    network_key_pair: KeyPair,
726) -> (
727    Option<Arc<TransactionGenerator>>,
728    Option<Arc<Mutex<DirectTransactionGenerator>>>,
729) {
730    // This tx generator directly push simple transactions and erc20
731    // transactions into blocks.
732    let maybe_direct_txgen_with_contract = if conf.is_test_or_dev_mode() {
733        Some(Arc::new(Mutex::new(DirectTransactionGenerator::new(
734            network_key_pair,
735            &public_to_address(DEV_GENESIS_KEY_PAIR_2.public(), true),
736            U256::from_dec_str("10000000000000000").unwrap(),
737            U256::from_dec_str("10000000000000000").unwrap(),
738        ))))
739    } else {
740        None
741    };
742
743    // This tx generator generates transactions from preconfigured multiple
744    // genesis accounts and it pushes transactions into transaction pool.
745    let maybe_multi_genesis_txgen = if let Some(txgen_conf) =
746        conf.tx_gen_config()
747    {
748        let multi_genesis_txgen = Arc::new(TransactionGenerator::new(
749            consensus.clone(),
750            txpool.clone(),
751            sync.clone(),
752            secret_store.clone(),
753        ));
754        if txgen_conf.generate_tx {
755            let txgen_clone = multi_genesis_txgen.clone();
756            let join_handle =
757                thread::Builder::new()
758                    .name("txgen".into())
759                    .spawn(move || {
760                        TransactionGenerator::generate_transactions_with_multiple_genesis_accounts(
761                            txgen_clone,
762                            txgen_conf,
763                            genesis_accounts,
764                        );
765                    })
766                    .expect("should succeed");
767            multi_genesis_txgen.set_join_handle(join_handle);
768        }
769        Some(multi_genesis_txgen)
770    } else {
771        None
772    };
773
774    (maybe_multi_genesis_txgen, maybe_direct_txgen_with_contract)
775}