massa_node/
main.rs

1// Copyright (c) 2022 MASSA LABS <info@massa.net>
2
3#![doc = include_str!("../../README.md")]
4#![warn(missing_docs)]
5#![warn(unused_crate_dependencies)]
6
7#[cfg(unix)]
8mod jemalloc_init;
9
10extern crate massa_logging;
11
12#[cfg(feature = "op_spammer")]
13use crate::operation_injector::start_operation_injector;
14use crate::settings::SETTINGS;
15use crate::survey::MassaSurvey;
16
17use cfg_if::cfg_if;
18use clap::{crate_version, Parser};
19use crossbeam_channel::TryRecvError;
20use dialoguer::Password;
21use massa_api::{ApiServer, ApiV2, Private, Public, RpcServer, StopHandle, API};
22use massa_api_exports::config::APIConfig;
23use massa_async_pool::AsyncPoolConfig;
24use massa_bootstrap::BootstrapError;
25use massa_bootstrap::{
26    get_state, start_bootstrap_server, BootstrapConfig, BootstrapManager, BootstrapTcpListener,
27    DefaultConnector,
28};
29use massa_channel::receiver::MassaReceiver;
30use massa_channel::MassaChannel;
31use massa_consensus_exports::events::ConsensusEvent;
32use massa_consensus_exports::{
33    ConsensusBroadcasts, ConsensusChannels, ConsensusConfig, ConsensusManager,
34};
35use massa_consensus_worker::start_consensus_worker;
36use massa_db_exports::{MassaDBConfig, MassaDBController};
37use massa_db_worker::MassaDB;
38use massa_deferred_calls::config::DeferredCallsConfig;
39use massa_executed_ops::{ExecutedDenunciationsConfig, ExecutedOpsConfig};
40use massa_execution_exports::{
41    CondomLimits, ExecutionChannels, ExecutionConfig, ExecutionManager, GasCosts,
42    StorageCostsConstants,
43};
44use massa_execution_worker::start_execution_worker;
45#[cfg(all(
46    feature = "dump-block",
47    feature = "file_storage_backend",
48    not(feature = "db_storage_backend")
49))]
50use massa_execution_worker::storage_backend::FileStorageBackend;
51#[cfg(all(feature = "dump-block", feature = "db_storage_backend"))]
52use massa_execution_worker::storage_backend::RocksDBStorageBackend;
53
54use massa_factory_exports::{FactoryChannels, FactoryConfig, FactoryManager};
55use massa_factory_worker::start_factory;
56use massa_final_state::{FinalState, FinalStateConfig, FinalStateController};
57use massa_grpc::config::{GrpcConfig, ServiceName};
58use massa_grpc::server::{MassaPrivateGrpc, MassaPublicGrpc};
59use massa_ledger_exports::LedgerConfig;
60use massa_ledger_worker::FinalLedger;
61use massa_logging::massa_trace;
62use massa_metrics::{MassaMetrics, MetricsStopper};
63use massa_models::address::Address;
64use massa_models::amount::Amount;
65use massa_models::config::constants::{
66    ASYNC_MSG_CST_GAS_COST, BLOCK_REWARD, BOOTSTRAP_RANDOMNESS_SIZE_BYTES, CHANNEL_SIZE,
67    CONSENSUS_BOOTSTRAP_PART_SIZE, DEFERRED_CALL_MAX_FUTURE_SLOTS, DELTA_F0,
68    DENUNCIATION_EXPIRE_PERIODS, ENDORSEMENT_COUNT, END_TIMESTAMP, GENESIS_KEY, GENESIS_TIMESTAMP,
69    INITIAL_DRAW_SEED, LEDGER_COST_PER_BYTE, LEDGER_ENTRY_BASE_COST,
70    LEDGER_ENTRY_DATASTORE_BASE_SIZE, MAX_ADVERTISE_LENGTH, MAX_ASYNC_GAS, MAX_ASYNC_POOL_LENGTH,
71    MAX_BLOCK_SIZE, MAX_BOOTSTRAP_BLOCKS, MAX_BOOTSTRAP_ERROR_LENGTH, MAX_BYTECODE_LENGTH,
72    MAX_CONSENSUS_BLOCKS_IDS, MAX_DATASTORE_ENTRY_COUNT, MAX_DATASTORE_KEY_LENGTH,
73    MAX_DATASTORE_VALUE_LENGTH, MAX_DEFERRED_CREDITS_LENGTH, MAX_DENUNCIATIONS_PER_BLOCK_HEADER,
74    MAX_DENUNCIATION_CHANGES_LENGTH, MAX_ENDORSEMENTS_PER_MESSAGE, MAX_EXECUTED_OPS_CHANGES_LENGTH,
75    MAX_EXECUTED_OPS_LENGTH, MAX_FUNCTION_NAME_LENGTH, MAX_GAS_PER_BLOCK, MAX_LEDGER_CHANGES_COUNT,
76    MAX_LISTENERS_PER_PEER, MAX_OPERATIONS_PER_BLOCK, MAX_OPERATIONS_PER_MESSAGE,
77    MAX_OPERATION_DATASTORE_ENTRY_COUNT, MAX_OPERATION_DATASTORE_KEY_LENGTH,
78    MAX_OPERATION_DATASTORE_VALUE_LENGTH, MAX_OPERATION_STORAGE_TIME, MAX_PARAMETERS_SIZE,
79    MAX_PEERS_IN_ANNOUNCEMENT_LIST, MAX_PRODUCTION_STATS_LENGTH, MAX_ROLLS_COUNT_LENGTH,
80    MAX_SIZE_CHANNEL_COMMANDS_CONNECTIVITY, MAX_SIZE_CHANNEL_COMMANDS_PEERS,
81    MAX_SIZE_CHANNEL_COMMANDS_PEER_TESTERS, MAX_SIZE_CHANNEL_COMMANDS_PROPAGATION_BLOCKS,
82    MAX_SIZE_CHANNEL_COMMANDS_PROPAGATION_ENDORSEMENTS,
83    MAX_SIZE_CHANNEL_COMMANDS_PROPAGATION_OPERATIONS, MAX_SIZE_CHANNEL_COMMANDS_RETRIEVAL_BLOCKS,
84    MAX_SIZE_CHANNEL_COMMANDS_RETRIEVAL_ENDORSEMENTS,
85    MAX_SIZE_CHANNEL_COMMANDS_RETRIEVAL_OPERATIONS, MAX_SIZE_CHANNEL_NETWORK_TO_BLOCK_HANDLER,
86    MAX_SIZE_CHANNEL_NETWORK_TO_ENDORSEMENT_HANDLER, MAX_SIZE_CHANNEL_NETWORK_TO_OPERATION_HANDLER,
87    MAX_SIZE_CHANNEL_NETWORK_TO_PEER_HANDLER, MIN_COMPATIBLE_VERSION,
88    MIP_STORE_STATS_BLOCK_CONSIDERED, OPERATION_VALIDITY_PERIODS, PERIODS_PER_CYCLE,
89    POS_MISS_RATE_DEACTIVATION_THRESHOLD, POS_SAVED_CYCLES, PROTOCOL_CONTROLLER_CHANNEL_SIZE,
90    PROTOCOL_EVENT_CHANNEL_SIZE, ROLL_COUNT_TO_SLASH_ON_DENUNCIATION, ROLL_PRICE,
91    SELECTOR_DRAW_CACHE_SIZE, T0, THREAD_COUNT, VERSION,
92};
93use massa_models::config::{
94    bootstrap_batch_allocation_budget, handle_disclaimer, BASE_OPERATION_GAS_COST, CHAINID,
95    DEFERRED_CALL_BASE_FEE_MAX_CHANGE_DENOMINATOR, DEFERRED_CALL_CST_GAS_COST,
96    DEFERRED_CALL_GLOBAL_OVERBOOKING_PENALTY, DEFERRED_CALL_MAX_ASYNC_GAS,
97    DEFERRED_CALL_MAX_POOL_CHANGES, DEFERRED_CALL_MIN_GAS_COST, DEFERRED_CALL_MIN_GAS_INCREMENT,
98    DEFERRED_CALL_SLOT_OVERBOOKING_PENALTY, KEEP_EXECUTED_HISTORY_EXTRA_PERIODS,
99    MAX_BOOTSTRAP_FINAL_STATE_PARTS_SIZE, MAX_BOOTSTRAP_VERSIONING_ELEMENTS_SIZE,
100    MAX_EVENT_DATA_SIZE, MAX_EVENT_DATA_SIZE_V0, MAX_EVENT_PER_OPERATION, MAX_MESSAGE_SIZE,
101    MAX_RECURSIVE_CALLS_DEPTH, MAX_RUNTIME_MODULE_CUSTOM_SECTION_DATA_LEN,
102    MAX_RUNTIME_MODULE_CUSTOM_SECTION_LEN, MAX_RUNTIME_MODULE_EXPORTS,
103    MAX_RUNTIME_MODULE_FUNCTIONS, MAX_RUNTIME_MODULE_FUNCTION_NAME_LEN,
104    MAX_RUNTIME_MODULE_GLOBAL_INITIALIZER, MAX_RUNTIME_MODULE_IMPORTS, MAX_RUNTIME_MODULE_MEMORIES,
105    MAX_RUNTIME_MODULE_NAME_LEN, MAX_RUNTIME_MODULE_PASSIVE_DATA,
106    MAX_RUNTIME_MODULE_PASSIVE_ELEMENT, MAX_RUNTIME_MODULE_SIGNATURE_LEN, MAX_RUNTIME_MODULE_TABLE,
107    MAX_RUNTIME_MODULE_TABLE_INITIALIZER,
108};
109use massa_models::slot::Slot;
110use massa_models::timeslots::get_block_slot_timestamp;
111use massa_pool_exports::{PoolBroadcasts, PoolChannels, PoolConfig, PoolManager};
112use massa_pool_worker::start_pool_controller;
113use massa_pos_exports::{PoSConfig, SelectorConfig, SelectorManager};
114use massa_pos_worker::start_selector_worker;
115use massa_protocol_exports::{ProtocolConfig, ProtocolManager, TransportType};
116use massa_protocol_worker::{create_protocol_controller, start_protocol_controller};
117use massa_signature::KeyPair;
118use massa_storage::Storage;
119use massa_time::MassaTime;
120use massa_versioning::keypair_factory::KeyPairFactory;
121use massa_versioning::mips::get_mip_list;
122use massa_versioning::versioning::{MipStatsConfig, MipStore};
123use massa_wallet::Wallet;
124use num::rational::Ratio;
125use parking_lot::RwLock;
126use settings::GrpcSettings;
127use std::collections::HashMap;
128use std::path::PathBuf;
129use std::sync::atomic::{AtomicUsize, Ordering};
130use std::sync::{Condvar, Mutex};
131use std::time::Duration;
132use std::{path::Path, process, sync::Arc};
133
134use massa_event_cache::config::EventCacheConfig;
135use massa_event_cache::worker::{start_event_cache_writer_worker, EventCacheManager};
136use survey::MassaSurveyStopper;
137use tokio::sync::broadcast;
138use tracing::{debug, error, info, warn};
139use tracing_subscriber::filter::{filter_fn, LevelFilter};
140
141#[cfg(feature = "op_spammer")]
142mod operation_injector;
143mod settings;
144mod survey;
145
146async fn launch(
147    args: &Args,
148    node_wallet: Arc<RwLock<Wallet>>,
149    sig_int_toggled: Arc<(Mutex<bool>, Condvar)>,
150) -> (
151    MassaReceiver<ConsensusEvent>,
152    Option<BootstrapManager>,
153    Box<dyn ConsensusManager>,
154    Box<dyn ExecutionManager>,
155    Box<dyn SelectorManager>,
156    Box<dyn PoolManager>,
157    Box<dyn ProtocolManager>,
158    Box<dyn FactoryManager>,
159    Box<dyn EventCacheManager>,
160    StopHandle,
161    StopHandle,
162    StopHandle,
163    Option<massa_grpc::server::StopHandle>,
164    Option<massa_grpc::server::StopHandle>,
165    MetricsStopper,
166    MassaSurveyStopper,
167) {
168    let now = MassaTime::now();
169
170    if let Some(end) = *END_TIMESTAMP {
171        if now > end {
172            panic!("This episode has come to an end, please get the latest testnet node version to continue");
173        }
174    }
175
176    // Storage shared by multiple components.
177    let shared_storage: Storage = Storage::create_root();
178
179    // init final state
180    let ledger_config = LedgerConfig {
181        thread_count: THREAD_COUNT,
182        initial_ledger_path: SETTINGS.ledger.initial_ledger_path.clone(),
183        max_key_length: MAX_DATASTORE_KEY_LENGTH,
184        max_datastore_value_length: MAX_DATASTORE_VALUE_LENGTH,
185        max_bytecode_size: MAX_BYTECODE_LENGTH,
186    };
187    let async_pool_config = AsyncPoolConfig {
188        max_length: MAX_ASYNC_POOL_LENGTH,
189        thread_count: THREAD_COUNT,
190        max_function_length: MAX_FUNCTION_NAME_LENGTH,
191        max_function_params_length: MAX_PARAMETERS_SIZE as u64,
192        max_key_length: MAX_DATASTORE_KEY_LENGTH as u32,
193    };
194    let pos_config = PoSConfig {
195        periods_per_cycle: PERIODS_PER_CYCLE,
196        thread_count: THREAD_COUNT,
197        cycle_history_length: POS_SAVED_CYCLES,
198        max_rolls_length: MAX_ROLLS_COUNT_LENGTH,
199        max_production_stats_length: MAX_PRODUCTION_STATS_LENGTH,
200        max_credit_length: MAX_DEFERRED_CREDITS_LENGTH,
201        initial_deferred_credits_path: SETTINGS.ledger.initial_deferred_credits_path.clone(),
202    };
203    let executed_ops_config = ExecutedOpsConfig {
204        thread_count: THREAD_COUNT,
205        keep_executed_history_extra_periods: KEEP_EXECUTED_HISTORY_EXTRA_PERIODS,
206    };
207    let executed_denunciations_config = ExecutedDenunciationsConfig {
208        denunciation_expire_periods: DENUNCIATION_EXPIRE_PERIODS,
209        thread_count: THREAD_COUNT,
210        endorsement_count: ENDORSEMENT_COUNT,
211        keep_executed_history_extra_periods: KEEP_EXECUTED_HISTORY_EXTRA_PERIODS,
212    };
213    let deferred_calls_config = DeferredCallsConfig {
214        thread_count: THREAD_COUNT,
215        max_function_name_length: MAX_FUNCTION_NAME_LENGTH,
216        max_parameter_size: MAX_PARAMETERS_SIZE,
217        max_pool_changes: DEFERRED_CALL_MAX_POOL_CHANGES,
218        max_gas: DEFERRED_CALL_MAX_ASYNC_GAS,
219        max_future_slots: DEFERRED_CALL_MAX_FUTURE_SLOTS,
220        base_fee_max_max_change_denominator: DEFERRED_CALL_BASE_FEE_MAX_CHANGE_DENOMINATOR,
221        min_gas_increment: DEFERRED_CALL_MIN_GAS_INCREMENT,
222        min_gas_cost: DEFERRED_CALL_MIN_GAS_COST,
223        global_overbooking_penalty: DEFERRED_CALL_GLOBAL_OVERBOOKING_PENALTY,
224        slot_overbooking_penalty: DEFERRED_CALL_SLOT_OVERBOOKING_PENALTY,
225        call_cst_gas_cost: DEFERRED_CALL_CST_GAS_COST,
226        ledger_cost_per_byte: LEDGER_COST_PER_BYTE,
227    };
228    let final_state_config = FinalStateConfig {
229        ledger_config: ledger_config.clone(),
230        async_pool_config,
231        deferred_calls_config,
232        pos_config,
233        executed_ops_config,
234        executed_denunciations_config,
235        final_history_length: SETTINGS.ledger.final_history_length,
236        thread_count: THREAD_COUNT,
237        periods_per_cycle: PERIODS_PER_CYCLE,
238        initial_seed_string: INITIAL_DRAW_SEED.into(),
239        initial_rolls_path: SETTINGS.selector.initial_rolls_path.clone(),
240        endorsement_count: ENDORSEMENT_COUNT,
241        max_executed_denunciations_length: MAX_DENUNCIATION_CHANGES_LENGTH,
242        max_denunciations_per_block_header: MAX_DENUNCIATIONS_PER_BLOCK_HEADER,
243        ledger_backup_periods_interval: SETTINGS.ledger.ledger_backup_periods_interval,
244        t0: T0,
245        genesis_timestamp: *GENESIS_TIMESTAMP,
246    };
247
248    // Start massa metrics
249    let (massa_metrics, metrics_stopper) = MassaMetrics::new(
250        SETTINGS.metrics.enabled,
251        SETTINGS.metrics.bind,
252        THREAD_COUNT,
253        SETTINGS.metrics.tick_delay.to_duration(),
254    );
255
256    // Remove current disk ledger if there is one and we don't want to restart from snapshot
257    // NOTE: this is temporary, since we cannot currently handle bootstrap from remaining ledger
258    if args.keep_ledger || args.restart_from_snapshot_at_period.is_some() {
259        info!("Loading old ledger for next episode");
260    } else {
261        if SETTINGS.ledger.disk_ledger_path.exists() {
262            std::fs::remove_dir_all(SETTINGS.ledger.disk_ledger_path.clone())
263                .expect("disk ledger delete failed");
264        }
265        if SETTINGS.execution.hd_cache_path.exists() {
266            std::fs::remove_dir_all(SETTINGS.execution.hd_cache_path.clone())
267                .expect("disk hd cache delete failed");
268        }
269    }
270
271    let db_config = MassaDBConfig {
272        path: SETTINGS.ledger.disk_ledger_path.clone(),
273        max_history_length: SETTINGS.ledger.final_history_length,
274        max_final_state_elements_size: MAX_BOOTSTRAP_FINAL_STATE_PARTS_SIZE.try_into().unwrap(),
275        max_versioning_elements_size: MAX_BOOTSTRAP_VERSIONING_ELEMENTS_SIZE.try_into().unwrap(),
276        thread_count: THREAD_COUNT,
277        max_ledger_backups: SETTINGS.ledger.max_ledger_backups,
278        enable_metrics: SETTINGS.metrics.enabled,
279    };
280    let db = Arc::new(RwLock::new(
281        Box::new(MassaDB::new(db_config)) as Box<dyn MassaDBController + 'static>
282    ));
283
284    // Create final ledger
285    let ledger = FinalLedger::new(ledger_config.clone(), db.clone());
286
287    // launch selector worker
288    let (selector_manager, selector_controller) = start_selector_worker(SelectorConfig {
289        max_draw_cache: SELECTOR_DRAW_CACHE_SIZE,
290        channel_size: CHANNEL_SIZE,
291        thread_count: THREAD_COUNT,
292        endorsement_count: ENDORSEMENT_COUNT,
293        periods_per_cycle: PERIODS_PER_CYCLE,
294        genesis_address: Address::from_public_key(&GENESIS_KEY.get_public_key()),
295    })
296    .expect("could not start selector worker");
297
298    // Creates an empty default store
299    let mip_stats_config = MipStatsConfig {
300        block_count_considered: MIP_STORE_STATS_BLOCK_CONSIDERED,
301        warn_announced_version_ratio: Ratio::new(
302            u64::from(SETTINGS.versioning.mip_stats_warn_announced_version),
303            100,
304        ),
305    };
306    // Ratio::new_raw(*SETTINGS.versioning.warn_announced_version_ratio, 100),
307
308    // Create final state, either from a snapshot, or from scratch
309    let final_state: Arc<RwLock<dyn FinalStateController>> = Arc::new(parking_lot::RwLock::new(
310        match args.restart_from_snapshot_at_period {
311            Some(last_start_period) => {
312                // The node is restarted from a snapshot:
313                // MIP store by reading from the db as it must have been updated by the massa ledger editor
314                // (to shift transitions that might have happened during the network shutdown)
315                // Note that FinalState::new_derived_from_snapshot will check if MIP store is consistent
316                // No Bootstrap are expected after this
317                let mip_store: MipStore = MipStore::try_from_db(db.clone(), mip_stats_config)
318                    .expect("MIP store creation failed");
319                debug!("After read from db, Mip store: {:?}", mip_store);
320
321                FinalState::new_derived_from_snapshot(
322                    db.clone(),
323                    final_state_config,
324                    Box::new(ledger),
325                    selector_controller.clone(),
326                    mip_store,
327                    last_start_period,
328                )
329                .expect("could not init final state")
330            }
331            None => {
332                // The node is started in a normal way
333                // Read the mip list supported by the current software
334                // The resulting MIP store will likely be updated by the bootstrap process in order
335                // to get the latest information for the MIP store (new states, votes...)
336
337                let mip_list = get_mip_list();
338                debug!("MIP list: {:?}", mip_list);
339                let mip_store = MipStore::try_from((mip_list, mip_stats_config))
340                    .expect("mip store creation failed");
341
342                FinalState::new(
343                    db.clone(),
344                    final_state_config,
345                    Box::new(ledger),
346                    selector_controller.clone(),
347                    mip_store,
348                    true,
349                )
350                .expect("could not init final state")
351            }
352        },
353    ));
354
355    let mip_store = final_state.read().get_mip_store().clone();
356
357    let bootstrap_config: BootstrapConfig = BootstrapConfig {
358        bootstrap_list: SETTINGS.bootstrap.bootstrap_list.clone(),
359        bootstrap_protocol: SETTINGS.bootstrap.bootstrap_protocol,
360        bootstrap_whitelist_path: SETTINGS.bootstrap.bootstrap_whitelist_path.clone(),
361        bootstrap_blacklist_path: SETTINGS.bootstrap.bootstrap_blacklist_path.clone(),
362        listen_addr: SETTINGS.bootstrap.bind,
363        min_compatible_version: *MIN_COMPATIBLE_VERSION,
364        connect_timeout: SETTINGS.bootstrap.connect_timeout,
365        bootstrap_timeout: SETTINGS.bootstrap.bootstrap_timeout,
366        read_timeout: SETTINGS.bootstrap.read_timeout,
367        write_timeout: SETTINGS.bootstrap.write_timeout,
368        read_error_timeout: SETTINGS.bootstrap.read_error_timeout,
369        write_error_timeout: SETTINGS.bootstrap.write_error_timeout,
370        retry_delay: SETTINGS.bootstrap.retry_delay,
371        max_ping: SETTINGS.bootstrap.max_ping,
372        max_clock_delta: SETTINGS.bootstrap.max_clock_delta,
373        cache_duration: SETTINGS.bootstrap.cache_duration,
374        keep_ledger: args.keep_ledger,
375        max_listeners_per_peer: MAX_LISTENERS_PER_PEER as u32,
376        max_simultaneous_bootstraps: SETTINGS.bootstrap.max_simultaneous_bootstraps,
377        per_ip_min_interval: SETTINGS.bootstrap.per_ip_min_interval,
378        ip_list_max_size: SETTINGS.bootstrap.ip_list_max_size,
379        rate_limit: SETTINGS.bootstrap.rate_limit,
380        max_datastore_key_length: MAX_DATASTORE_KEY_LENGTH,
381        randomness_size_bytes: BOOTSTRAP_RANDOMNESS_SIZE_BYTES,
382        thread_count: THREAD_COUNT,
383        t0: T0,
384        genesis_timestamp: *GENESIS_TIMESTAMP,
385        periods_per_cycle: PERIODS_PER_CYCLE,
386        endorsement_count: ENDORSEMENT_COUNT,
387        max_advertise_length: MAX_ADVERTISE_LENGTH,
388        max_bootstrap_blocks_length: MAX_BOOTSTRAP_BLOCKS,
389        max_bootstrap_error_length: MAX_BOOTSTRAP_ERROR_LENGTH,
390        max_final_state_elements_size: MAX_BOOTSTRAP_FINAL_STATE_PARTS_SIZE,
391        max_versioning_elements_size: MAX_BOOTSTRAP_VERSIONING_ELEMENTS_SIZE,
392        // Bound the in-memory footprint of a batch we *receive*. The sending side applies the
393        // same accounting to the batches it builds (see `get_batch_to_stream`), so a batch any
394        // peer of this version sends fits the budget a peer of this version parses with.
395        max_final_state_batch_allocation: bootstrap_batch_allocation_budget(
396            MAX_BOOTSTRAP_FINAL_STATE_PARTS_SIZE as usize,
397        ) as u64,
398        max_versioning_batch_allocation: bootstrap_batch_allocation_budget(
399            MAX_BOOTSTRAP_VERSIONING_ELEMENTS_SIZE as usize,
400        ) as u64,
401        max_operations_per_block: MAX_OPERATIONS_PER_BLOCK,
402        max_datastore_entry_count: MAX_DATASTORE_ENTRY_COUNT,
403        max_datastore_value_length: MAX_DATASTORE_VALUE_LENGTH,
404        max_function_name_length: MAX_FUNCTION_NAME_LENGTH,
405        max_ledger_changes_count: MAX_LEDGER_CHANGES_COUNT,
406        max_parameters_size: MAX_PARAMETERS_SIZE,
407        max_op_datastore_entry_count: MAX_OPERATION_DATASTORE_ENTRY_COUNT,
408        max_op_datastore_key_length: MAX_OPERATION_DATASTORE_KEY_LENGTH,
409        max_op_datastore_value_length: MAX_OPERATION_DATASTORE_VALUE_LENGTH,
410        max_changes_slot_count: SETTINGS.ledger.final_history_length as u64,
411        max_rolls_length: MAX_ROLLS_COUNT_LENGTH,
412        max_production_stats_length: MAX_PRODUCTION_STATS_LENGTH,
413        max_credits_length: MAX_DEFERRED_CREDITS_LENGTH,
414        max_executed_ops_length: MAX_EXECUTED_OPS_LENGTH,
415        max_ops_changes_length: MAX_EXECUTED_OPS_CHANGES_LENGTH,
416        consensus_bootstrap_part_size: CONSENSUS_BOOTSTRAP_PART_SIZE,
417        max_consensus_block_ids: MAX_CONSENSUS_BLOCKS_IDS,
418        mip_store_stats_block_considered: MIP_STORE_STATS_BLOCK_CONSIDERED,
419        max_denunciations_per_block_header: MAX_DENUNCIATIONS_PER_BLOCK_HEADER,
420        max_denunciation_changes_length: MAX_DENUNCIATION_CHANGES_LENGTH,
421        chain_id: *CHAINID,
422    };
423
424    let bootstrap_state = match get_state(
425        &bootstrap_config,
426        final_state.clone(),
427        DefaultConnector,
428        *VERSION,
429        *GENESIS_TIMESTAMP,
430        *END_TIMESTAMP,
431        args.restart_from_snapshot_at_period,
432        sig_int_toggled.clone(),
433        massa_metrics.clone(),
434    ) {
435        Ok(vals) => vals,
436        Err(BootstrapError::Interrupted(msg)) => {
437            info!("{}", msg);
438            process::exit(0);
439        }
440        Err(err) => panic!("critical error detected in the bootstrap process: {}", err),
441    };
442
443    if !final_state.read().is_db_valid() {
444        // TODO: Bootstrap again instead of panicking
445        panic!("critical: db is not valid after bootstrap");
446    }
447
448    if args.restart_from_snapshot_at_period.is_none() {
449        final_state.write().recompute_caches();
450
451        // give the controller to final state in order for it to feed the cycles
452        final_state
453            .write()
454            .compute_initial_draws()
455            .expect("could not compute initial draws"); // TODO: this might just mean a bad bootstrap, no need to panic, just reboot
456    }
457
458    // Note: the consistency of the network restart metadata with the MIP store has already been
459    // checked, either by the bootstrap client before storing what the server sent, or by
460    // `FinalState::new_derived_from_snapshot` when restarting from a local snapshot.
461    if final_state.read().get_last_slot_before_downtime().is_some() {
462        // If we are before a network restart, print the hash to make it easier to debug bootstrapping issues
463        let now = MassaTime::now();
464        let last_start_slot = Slot::new(
465            final_state.read().get_last_start_period(),
466            THREAD_COUNT.saturating_sub(1),
467        );
468        let last_start_slot_timestamp =
469            get_block_slot_timestamp(THREAD_COUNT, T0, *GENESIS_TIMESTAMP, last_start_slot)
470                .expect("Can't get timestamp for last_start_slot");
471        if now < last_start_slot_timestamp {
472            let final_state_hash = final_state.read().get_fingerprint();
473            info!(
474                "final_state hash before network restarts at slot {}: {}",
475                last_start_slot, final_state_hash
476            );
477        }
478    }
479
480    // Event cache thread
481    let event_cache_config = EventCacheConfig {
482        event_cache_path: SETTINGS.execution.event_cache_path.clone(),
483        max_event_cache_length: SETTINGS.execution.event_cache_size,
484        snip_amount: SETTINGS.execution.event_snip_amount,
485        // Note: we still use the v0 event data size for the event cache to be able to deserialize
486        // events that bypass the v1 event limitation
487        max_event_data_length: MAX_EVENT_DATA_SIZE_V0 as u64,
488        thread_count: THREAD_COUNT,
489        // Note: SCOutputEvent call stack comes from the execution module, and we assume
490        //       this should return a limited call stack length
491        //       The value remains for future use & limitations
492        max_call_stack_length: u16::MAX,
493
494        max_events_per_operation: MAX_EVENT_PER_OPERATION as u64,
495        max_operations_per_block: MAX_OPERATIONS_PER_BLOCK as u64,
496        max_events_per_query: SETTINGS.execution.max_event_per_query,
497    };
498    let (event_cache_manager, event_cache_controller) =
499        start_event_cache_writer_worker(event_cache_config);
500
501    // Storage costs constants
502    let storage_costs_constants = StorageCostsConstants {
503        ledger_cost_per_byte: LEDGER_COST_PER_BYTE,
504        ledger_entry_base_cost: LEDGER_ENTRY_BASE_COST,
505        ledger_entry_datastore_base_cost: LEDGER_COST_PER_BYTE
506            .checked_mul_u64(LEDGER_ENTRY_DATASTORE_BASE_SIZE as u64)
507            .expect("Overflow when creating constant ledger_entry_datastore_base_size"),
508    };
509
510    // gas costs
511    let gas_costs = GasCosts::new(SETTINGS.execution.abi_gas_costs_file.clone())
512        .expect("Failed to load gas costs");
513
514    // Limits imposed to wasm files so the compilation phase is smooth
515    let condom_limits = CondomLimits {
516        max_exports: Some(MAX_RUNTIME_MODULE_EXPORTS),
517        max_functions: Some(MAX_RUNTIME_MODULE_FUNCTIONS),
518        max_signature_len: Some(MAX_RUNTIME_MODULE_SIGNATURE_LEN),
519        max_name_len: Some(MAX_RUNTIME_MODULE_NAME_LEN),
520        max_imports_len: Some(MAX_RUNTIME_MODULE_IMPORTS),
521        max_table_initializers_len: Some(MAX_RUNTIME_MODULE_TABLE_INITIALIZER),
522        max_passive_elements_len: Some(MAX_RUNTIME_MODULE_PASSIVE_ELEMENT),
523        max_passive_data_len: Some(MAX_RUNTIME_MODULE_PASSIVE_DATA),
524        max_global_initializers_len: Some(MAX_RUNTIME_MODULE_GLOBAL_INITIALIZER),
525        max_function_names_len: Some(MAX_RUNTIME_MODULE_FUNCTION_NAME_LEN),
526        max_tables_count: Some(MAX_RUNTIME_MODULE_TABLE),
527        max_memories_len: Some(MAX_RUNTIME_MODULE_MEMORIES),
528        max_globals_len: Some(MAX_RUNTIME_MODULE_GLOBAL_INITIALIZER),
529        max_custom_sections_len: Some(MAX_RUNTIME_MODULE_CUSTOM_SECTION_LEN),
530        max_custom_sections_data_len: Some(MAX_RUNTIME_MODULE_CUSTOM_SECTION_DATA_LEN),
531    };
532
533    let block_dump_folder_path = SETTINGS.block_dump.block_dump_folder_path.clone();
534    if !block_dump_folder_path.exists() {
535        info!("Current folder: {:?}", std::env::current_dir().unwrap());
536        info!("Creating dump folder: {:?}", block_dump_folder_path);
537        std::fs::create_dir_all(block_dump_folder_path.clone())
538            .expect("Cannot create dump block folder");
539    }
540
541    // launch execution module
542    let execution_config = ExecutionConfig {
543        max_final_events: SETTINGS.execution.max_final_events,
544        readonly_queue_length: SETTINGS.execution.readonly_queue_length,
545        readonly_starvation_timeout: SETTINGS.execution.readonly_starvation_timeout,
546        cursor_delay: SETTINGS.execution.cursor_delay,
547        max_async_gas: MAX_ASYNC_GAS,
548        async_msg_cst_gas_cost: ASYNC_MSG_CST_GAS_COST,
549        max_gas_per_block: MAX_GAS_PER_BLOCK,
550        roll_price: ROLL_PRICE,
551        thread_count: THREAD_COUNT,
552        t0: T0,
553        genesis_timestamp: *GENESIS_TIMESTAMP,
554        block_reward: BLOCK_REWARD,
555        endorsement_count: ENDORSEMENT_COUNT as u64,
556        operation_validity_period: OPERATION_VALIDITY_PERIODS,
557        periods_per_cycle: PERIODS_PER_CYCLE,
558        stats_time_window_duration: SETTINGS.execution.stats_time_window_duration,
559        max_miss_ratio: *POS_MISS_RATE_DEACTIVATION_THRESHOLD,
560        max_datastore_key_length: MAX_DATASTORE_KEY_LENGTH,
561        max_bytecode_size: MAX_BYTECODE_LENGTH,
562        max_datastore_value_size: MAX_DATASTORE_VALUE_LENGTH,
563        storage_costs_constants,
564        max_read_only_gas: SETTINGS.execution.max_read_only_gas,
565        gas_costs: gas_costs.clone(),
566        base_operation_gas_cost: BASE_OPERATION_GAS_COST,
567        last_start_period: final_state.read().get_last_start_period(),
568        hd_cache_path: SETTINGS.execution.hd_cache_path.clone(),
569        lru_cache_size: SETTINGS.execution.lru_cache_size,
570        hd_cache_size: SETTINGS.execution.hd_cache_size,
571        snip_amount: SETTINGS.execution.snip_amount,
572        roll_count_to_slash_on_denunciation: ROLL_COUNT_TO_SLASH_ON_DENUNCIATION,
573        denunciation_expire_periods: DENUNCIATION_EXPIRE_PERIODS,
574        broadcast_enabled: SETTINGS.api.enable_broadcast,
575        broadcast_slot_execution_output_channel_capacity: SETTINGS
576            .execution
577            .broadcast_slot_execution_output_channel_capacity,
578        max_event_size_v0: MAX_EVENT_DATA_SIZE_V0,
579        max_event_size: MAX_EVENT_DATA_SIZE,
580        max_function_length: MAX_FUNCTION_NAME_LENGTH,
581        max_parameter_length: MAX_PARAMETERS_SIZE,
582        chain_id: *CHAINID,
583        #[cfg(feature = "execution-trace")]
584        broadcast_traces_enabled: true,
585        #[cfg(not(feature = "execution-trace"))]
586        broadcast_traces_enabled: false,
587        broadcast_slot_execution_traces_channel_capacity: SETTINGS
588            .execution
589            .broadcast_slot_execution_traces_channel_capacity,
590        max_execution_traces_slot_limit: SETTINGS.execution.execution_traces_limit,
591        block_dump_folder_path,
592        max_recursive_calls_depth: MAX_RECURSIVE_CALLS_DEPTH,
593        condom_limits,
594        deferred_calls_config,
595        max_event_per_operation: MAX_EVENT_PER_OPERATION,
596        event_cache_path: SETTINGS.execution.event_cache_path.clone(),
597        event_cache_size: SETTINGS.execution.event_cache_size,
598        event_snip_amount: SETTINGS.execution.event_snip_amount,
599        broadcast_slot_execution_info_channel_capacity: SETTINGS
600            .execution
601            .broadcast_slot_execution_info_channel_capacity,
602    };
603
604    let execution_channels = ExecutionChannels {
605        slot_execution_output_sender: broadcast::channel(
606            execution_config.broadcast_slot_execution_output_channel_capacity,
607        )
608        .0,
609        #[cfg(feature = "execution-trace")]
610        slot_execution_traces_sender: broadcast::channel(
611            execution_config.broadcast_slot_execution_traces_channel_capacity,
612        )
613        .0,
614        #[cfg(feature = "execution-info")]
615        slot_execution_info_sender: broadcast::channel(
616            execution_config.broadcast_slot_execution_info_channel_capacity,
617        )
618        .0,
619    };
620
621    cfg_if! {
622        if #[cfg(all(feature = "dump-block", feature = "db_storage_backend"))] {
623            let block_storage_backend = Arc::new(RwLock::new(
624                RocksDBStorageBackend::new(
625                execution_config.block_dump_folder_path.clone(), SETTINGS.block_dump.max_blocks),
626            ));
627        } else if #[cfg(all(feature = "dump-block", feature = "file_storage_backend"))] {
628            let block_storage_backend = Arc::new(RwLock::new(
629                FileStorageBackend::new(
630                execution_config.block_dump_folder_path.clone(), SETTINGS.block_dump.max_blocks),
631            ));
632        } else if #[cfg(feature = "dump-block")] {
633            compile_error!("feature dump-block requise either db_storage_backend or file_storage_backend");
634        }
635    }
636
637    let (execution_manager, execution_controller) = start_execution_worker(
638        execution_config,
639        final_state.clone(),
640        selector_controller.clone(),
641        mip_store.clone(),
642        execution_channels.clone(),
643        node_wallet.clone(),
644        massa_metrics.clone(),
645        event_cache_controller,
646        #[cfg(feature = "dump-block")]
647        block_storage_backend.clone(),
648    );
649
650    // launch pool controller
651    let pool_config = PoolConfig {
652        thread_count: THREAD_COUNT,
653        max_block_size: MAX_BLOCK_SIZE,
654        max_block_gas: MAX_GAS_PER_BLOCK,
655        base_operation_gas_cost: BASE_OPERATION_GAS_COST,
656        sp_compilation_cost: gas_costs.sp_compilation_cost,
657        roll_price: ROLL_PRICE,
658        max_block_endorsement_count: ENDORSEMENT_COUNT,
659        operation_validity_periods: OPERATION_VALIDITY_PERIODS,
660        max_operations_per_block: MAX_OPERATIONS_PER_BLOCK,
661        max_operation_pool_size: SETTINGS.pool.max_operation_pool_size,
662        max_operation_pool_excess_items: SETTINGS.pool.max_operation_pool_excess_items,
663        operation_pool_refresh_interval: SETTINGS.pool.operation_pool_refresh_interval,
664        operation_pool_swap_interval: SETTINGS.pool.operation_pool_swap_interval,
665        endorsement_pool_swap_interval: SETTINGS.pool.endorsement_pool_swap_interval,
666        denunciation_pool_refresh_interval: SETTINGS.pool.denunciation_pool_refresh_interval,
667        denunciation_pool_swap_interval: SETTINGS.pool.denunciation_pool_swap_interval,
668        operation_max_future_start_delay: SETTINGS.pool.operation_max_future_start_delay,
669        max_endorsements_pool_size_per_thread: SETTINGS.pool.max_endorsements_pool_size_per_thread,
670        operations_channel_size: SETTINGS.pool.operations_channel_capacity,
671        endorsements_channel_size: SETTINGS.pool.endorsements_channel_capacity,
672        denunciations_channel_size: SETTINGS.pool.denunciations_channel_capacity,
673        broadcast_enabled: SETTINGS.api.enable_broadcast,
674        broadcast_endorsements_channel_capacity: SETTINGS
675            .pool
676            .broadcast_endorsements_channel_capacity,
677        broadcast_operations_channel_capacity: SETTINGS.pool.broadcast_operations_channel_capacity,
678        genesis_timestamp: *GENESIS_TIMESTAMP,
679        t0: T0,
680        periods_per_cycle: PERIODS_PER_CYCLE,
681        denunciation_expire_periods: DENUNCIATION_EXPIRE_PERIODS,
682        max_denunciations_per_block_header: MAX_DENUNCIATIONS_PER_BLOCK_HEADER,
683        minimal_fees: SETTINGS.pool.minimal_fees,
684        last_start_period: final_state.read().get_last_start_period(),
685    };
686
687    let pool_channels = PoolChannels {
688        broadcasts: PoolBroadcasts {
689            endorsement_sender: broadcast::channel(
690                pool_config.broadcast_endorsements_channel_capacity,
691            )
692            .0,
693            operation_sender: broadcast::channel(pool_config.broadcast_operations_channel_capacity)
694                .0,
695        },
696        selector: selector_controller.clone(),
697        execution_controller: execution_controller.clone(),
698    };
699
700    let (pool_manager, pool_controller) = start_pool_controller(
701        pool_config,
702        &shared_storage,
703        pool_channels.clone(),
704        node_wallet.clone(),
705    );
706
707    // launch protocol controller
708    let mut listeners = HashMap::default();
709    listeners.insert(SETTINGS.protocol.bind, TransportType::Tcp);
710    let protocol_config = ProtocolConfig {
711        thread_count: THREAD_COUNT,
712        ask_block_timeout: SETTINGS.protocol.ask_block_timeout,
713        max_known_blocks_size: SETTINGS.protocol.max_known_blocks_size,
714        max_node_known_blocks_size: SETTINGS.protocol.max_node_known_blocks_size,
715        max_block_propagation_time: SETTINGS.protocol.max_block_propagation_time,
716        max_node_wanted_blocks_size: SETTINGS.protocol.max_node_wanted_blocks_size,
717        max_known_ops_size: SETTINGS.protocol.max_known_ops_size,
718        max_node_known_ops_size: SETTINGS.protocol.max_node_known_ops_size,
719        max_known_endorsements_size: SETTINGS.protocol.max_known_endorsements_size,
720        max_node_known_endorsements_size: SETTINGS.protocol.max_node_known_endorsements_size,
721        max_simultaneous_ask_blocks_per_node: SETTINGS
722            .protocol
723            .max_simultaneous_ask_blocks_per_node,
724        max_send_wait: SETTINGS.protocol.max_send_wait,
725        operation_batch_buffer_capacity: SETTINGS.protocol.operation_batch_buffer_capacity,
726        operation_announcement_buffer_capacity: SETTINGS
727            .protocol
728            .operation_announcement_buffer_capacity,
729        operation_batch_proc_period: SETTINGS.protocol.operation_batch_proc_period,
730        operation_announcement_interval: SETTINGS.protocol.operation_announcement_interval,
731        max_operations_per_message: SETTINGS.protocol.max_operations_per_message,
732        max_serialized_operations_size_per_block: MAX_BLOCK_SIZE as usize,
733        max_block_gas: MAX_GAS_PER_BLOCK,
734        base_operation_gas_cost: BASE_OPERATION_GAS_COST,
735        sp_compilation_cost: gas_costs.sp_compilation_cost,
736        max_operations_per_block: MAX_OPERATIONS_PER_BLOCK,
737        controller_channel_size: PROTOCOL_CONTROLLER_CHANNEL_SIZE,
738        event_channel_size: PROTOCOL_EVENT_CHANNEL_SIZE,
739        genesis_timestamp: *GENESIS_TIMESTAMP,
740        t0: T0,
741        endorsement_count: ENDORSEMENT_COUNT,
742        max_message_size: MAX_MESSAGE_SIZE as usize,
743        max_ops_kept_for_propagation: SETTINGS.protocol.max_ops_kept_for_propagation,
744        max_endorsements_per_propagation_round: SETTINGS
745            .protocol
746            .max_endorsements_per_propagation_round,
747        max_operations_propagation_time: SETTINGS.protocol.max_operations_propagation_time,
748        max_endorsements_propagation_time: SETTINGS.protocol.max_endorsements_propagation_time,
749        last_start_period: final_state.read().get_last_start_period(),
750        max_endorsements_per_message: MAX_ENDORSEMENTS_PER_MESSAGE as u64,
751        max_denunciations_in_block_header: MAX_DENUNCIATIONS_PER_BLOCK_HEADER,
752        initial_peers: SETTINGS.protocol.initial_peers_file.clone(),
753        listeners,
754        keypair_file: SETTINGS.protocol.keypair_file.clone(),
755        max_blocks_kept_for_propagation: SETTINGS.protocol.max_blocks_kept_for_propagation,
756        block_propagation_tick: SETTINGS.protocol.block_propagation_tick,
757        asked_operations_buffer_capacity: SETTINGS.protocol.asked_operations_buffer_capacity,
758        thread_tester_count: SETTINGS.protocol.thread_tester_count,
759        max_operation_storage_time: MAX_OPERATION_STORAGE_TIME,
760        max_size_channel_commands_propagation_blocks: MAX_SIZE_CHANNEL_COMMANDS_PROPAGATION_BLOCKS,
761        max_size_channel_commands_propagation_operations:
762            MAX_SIZE_CHANNEL_COMMANDS_PROPAGATION_OPERATIONS,
763        max_size_channel_commands_propagation_endorsements:
764            MAX_SIZE_CHANNEL_COMMANDS_PROPAGATION_ENDORSEMENTS,
765        max_size_channel_commands_retrieval_blocks: MAX_SIZE_CHANNEL_COMMANDS_RETRIEVAL_BLOCKS,
766        max_size_channel_commands_retrieval_operations:
767            MAX_SIZE_CHANNEL_COMMANDS_RETRIEVAL_OPERATIONS,
768        max_size_channel_commands_retrieval_endorsements:
769            MAX_SIZE_CHANNEL_COMMANDS_RETRIEVAL_ENDORSEMENTS,
770        max_size_channel_commands_connectivity: MAX_SIZE_CHANNEL_COMMANDS_CONNECTIVITY,
771        max_size_channel_commands_peers: MAX_SIZE_CHANNEL_COMMANDS_PEERS,
772        max_size_channel_commands_peer_testers: MAX_SIZE_CHANNEL_COMMANDS_PEER_TESTERS,
773        max_size_channel_network_to_block_handler: MAX_SIZE_CHANNEL_NETWORK_TO_BLOCK_HANDLER,
774        max_size_channel_network_to_operation_handler:
775            MAX_SIZE_CHANNEL_NETWORK_TO_OPERATION_HANDLER,
776        max_size_channel_network_to_endorsement_handler:
777            MAX_SIZE_CHANNEL_NETWORK_TO_ENDORSEMENT_HANDLER,
778        max_size_channel_network_to_peer_handler: MAX_SIZE_CHANNEL_NETWORK_TO_PEER_HANDLER,
779        max_bytecode_size: MAX_BYTECODE_LENGTH,
780        max_op_datastore_entry_count: MAX_OPERATION_DATASTORE_ENTRY_COUNT,
781        max_op_datastore_key_length: MAX_OPERATION_DATASTORE_KEY_LENGTH,
782        max_op_datastore_value_length: MAX_OPERATION_DATASTORE_VALUE_LENGTH,
783        max_size_function_name: MAX_FUNCTION_NAME_LENGTH,
784        max_size_call_sc_parameter: MAX_PARAMETERS_SIZE,
785        max_size_listeners_per_peer: MAX_LISTENERS_PER_PEER,
786        max_size_peers_announcement: MAX_PEERS_IN_ANNOUNCEMENT_LIST,
787        read_write_limit_bytes_per_second: SETTINGS.protocol.read_write_limit_bytes_per_second
788            as u128,
789        try_connection_timer: SETTINGS.protocol.try_connection_timer,
790        unban_everyone_timer: SETTINGS.protocol.unban_everyone_timer,
791        max_in_connections: SETTINGS.protocol.max_in_connections,
792        timeout_connection: SETTINGS.protocol.timeout_connection,
793        message_timeout: SETTINGS.protocol.message_timeout,
794        tester_timeout: SETTINGS.protocol.tester_timeout,
795        routable_ip: SETTINGS
796            .protocol
797            .routable_ip
798            .or(SETTINGS.network.routable_ip),
799        debug: false,
800        peers_categories: SETTINGS.protocol.peers_categories.clone(),
801        default_category_info: SETTINGS.protocol.default_category_info,
802        version: *VERSION,
803        min_compatible_version: *MIN_COMPATIBLE_VERSION,
804        try_connection_timer_same_peer: SETTINGS.protocol.try_connection_timer_same_peer,
805        test_oldest_peer_cooldown: SETTINGS.protocol.test_oldest_peer_cooldown,
806        rate_limit: SETTINGS.protocol.rate_limit,
807        chain_id: *CHAINID,
808    };
809
810    let (protocol_controller, protocol_channels) =
811        create_protocol_controller(protocol_config.clone());
812
813    let consensus_config = ConsensusConfig {
814        genesis_timestamp: *GENESIS_TIMESTAMP,
815        end_timestamp: *END_TIMESTAMP,
816        thread_count: THREAD_COUNT,
817        t0: T0,
818        genesis_key: GENESIS_KEY.clone(),
819        max_discarded_blocks: SETTINGS.consensus.max_discarded_blocks,
820        max_future_processing_blocks: SETTINGS.consensus.max_future_processing_blocks,
821        max_dependency_blocks: SETTINGS.consensus.max_dependency_blocks,
822        delta_f0: DELTA_F0,
823        operation_validity_periods: OPERATION_VALIDITY_PERIODS,
824        periods_per_cycle: PERIODS_PER_CYCLE,
825        stats_timespan: SETTINGS.consensus.stats_timespan,
826        force_keep_final_periods: SETTINGS.consensus.force_keep_final_periods,
827        endorsement_count: ENDORSEMENT_COUNT,
828        block_db_prune_interval: SETTINGS.consensus.block_db_prune_interval,
829        max_gas_per_block: MAX_GAS_PER_BLOCK,
830        channel_size: CHANNEL_SIZE,
831        bootstrap_part_size: CONSENSUS_BOOTSTRAP_PART_SIZE,
832        broadcast_enabled: SETTINGS.api.enable_broadcast,
833        broadcast_blocks_headers_channel_capacity: SETTINGS
834            .consensus
835            .broadcast_blocks_headers_channel_capacity,
836        broadcast_blocks_channel_capacity: SETTINGS.consensus.broadcast_blocks_channel_capacity,
837        broadcast_filled_blocks_channel_capacity: SETTINGS
838            .consensus
839            .broadcast_filled_blocks_channel_capacity,
840        last_start_period: final_state.read().get_last_start_period(),
841        force_keep_final_periods_without_ops: SETTINGS
842            .consensus
843            .force_keep_final_periods_without_ops,
844        chain_id: *CHAINID,
845    };
846
847    // Blocks are kept in the shared storage by consensus until they are older than
848    // `force_keep_final_periods` (past that point `strip_to_block` drops the storage
849    // reference, so the block is no longer servable to peers). That retention is what
850    // covers the tail of the propagation window when the propagation cache has to evict
851    // an entry early because of its `max_blocks_kept_for_propagation` size cap.
852    let consensus_block_retention = consensus_config
853        .t0
854        .saturating_mul(consensus_config.force_keep_final_periods);
855    if consensus_block_retention < protocol_config.max_block_propagation_time {
856        warn!(
857            "consensus only keeps blocks for {:?} (force_keep_final_periods={} * t0={:?}) which is shorter than max_block_propagation_time={:?}: a block evicted early by the max_blocks_kept_for_propagation cap will not be retrievable by peers anymore, consider raising force_keep_final_periods or lowering max_block_propagation_time",
858            consensus_block_retention.to_duration(),
859            consensus_config.force_keep_final_periods,
860            consensus_config.t0.to_duration(),
861            protocol_config.max_block_propagation_time.to_duration()
862        );
863    }
864
865    let (consensus_event_sender, consensus_event_receiver) =
866        MassaChannel::new("consensus_event".to_string(), Some(CHANNEL_SIZE));
867    let consensus_channels = ConsensusChannels {
868        execution_controller: execution_controller.clone(),
869        selector_controller: selector_controller.clone(),
870        pool_controller: pool_controller.clone(),
871        controller_event_tx: consensus_event_sender,
872        protocol_controller: protocol_controller.clone(),
873        broadcasts: ConsensusBroadcasts {
874            block_header_sender: broadcast::channel(
875                consensus_config.broadcast_blocks_headers_channel_capacity,
876            )
877            .0,
878            block_sender: broadcast::channel(consensus_config.broadcast_blocks_channel_capacity).0,
879            filled_block_sender: broadcast::channel(
880                consensus_config.broadcast_filled_blocks_channel_capacity,
881            )
882            .0,
883        },
884    };
885
886    let (consensus_controller, consensus_manager) = start_consensus_worker(
887        consensus_config,
888        consensus_channels.clone(),
889        bootstrap_state.graph,
890        shared_storage.clone(),
891        massa_metrics.clone(),
892    );
893
894    let (protocol_manager, keypair, node_id) = start_protocol_controller(
895        protocol_config.clone(),
896        selector_controller.clone(),
897        consensus_controller.clone(),
898        bootstrap_state.peers,
899        pool_controller.clone(),
900        shared_storage.clone(),
901        protocol_channels,
902        mip_store.clone(),
903        massa_metrics.clone(),
904    )
905    .expect("could not start protocol controller");
906
907    // launch factory
908    let factory_config = FactoryConfig {
909        thread_count: THREAD_COUNT,
910        genesis_timestamp: *GENESIS_TIMESTAMP,
911        t0: T0,
912        initial_delay: SETTINGS.factory.initial_delay,
913        max_block_size: MAX_BLOCK_SIZE as u64,
914        max_block_gas: MAX_GAS_PER_BLOCK,
915        max_operations_per_block: MAX_OPERATIONS_PER_BLOCK,
916        endorsement_count: ENDORSEMENT_COUNT,
917        last_start_period: final_state.read().get_last_start_period(),
918        periods_per_cycle: PERIODS_PER_CYCLE,
919        denunciation_expire_periods: DENUNCIATION_EXPIRE_PERIODS,
920        stop_production_when_zero_connections: SETTINGS
921            .factory
922            .stop_production_when_zero_connections,
923        chain_id: *CHAINID,
924        block_delay_warn: SETTINGS.factory.block_delay_warn_threshold,
925        block_opt_channel_timeout: SETTINGS.factory.block_opt_channel_timeout,
926    };
927    let factory_channels = FactoryChannels {
928        selector: selector_controller.clone(),
929        consensus: consensus_controller.clone(),
930        pool: pool_controller.clone(),
931        protocol: protocol_controller.clone(),
932        storage: shared_storage.clone(),
933    };
934    let factory_manager = start_factory(
935        factory_config,
936        node_wallet.clone(),
937        factory_channels,
938        mip_store.clone(),
939    );
940
941    let bootstrap_manager = bootstrap_config.listen_addr.map(|addr| {
942        let (listener_stopper, listener) =
943            BootstrapTcpListener::create(&addr).unwrap_or_else(|_| {
944                panic!(
945                    "{}",
946                    format!("Could not bind to address: {}", addr).as_str()
947                )
948            });
949
950        start_bootstrap_server(
951            listener,
952            listener_stopper,
953            consensus_controller.clone(),
954            protocol_controller.clone(),
955            final_state.clone(),
956            bootstrap_config,
957            keypair.clone(),
958            *VERSION,
959            massa_metrics.clone(),
960        )
961        .expect("Could not start bootstrap server")
962    });
963
964    let api_config: APIConfig = APIConfig {
965        bind_private: SETTINGS.api.bind_private,
966        bind_public: SETTINGS.api.bind_public,
967        bind_api: SETTINGS.api.bind_api,
968        draw_lookahead_period_count: SETTINGS.api.draw_lookahead_period_count,
969        max_arguments: SETTINGS.api.max_arguments,
970        openrpc_spec_path: SETTINGS.api.openrpc_spec_path.clone(),
971        bootstrap_whitelist_path: SETTINGS.bootstrap.bootstrap_whitelist_path.clone(),
972        bootstrap_blacklist_path: SETTINGS.bootstrap.bootstrap_blacklist_path.clone(),
973        max_request_body_size: SETTINGS.api.max_request_body_size,
974        max_response_body_size: SETTINGS.api.max_response_body_size,
975        max_connections: SETTINGS.api.max_connections,
976        max_subscriptions_per_connection: SETTINGS.api.max_subscriptions_per_connection,
977        max_log_length: SETTINGS.api.max_log_length,
978        allow_hosts: SETTINGS.api.allow_hosts.clone(),
979        batch_request_limit: SETTINGS.api.batch_request_limit,
980        ping_interval: SETTINGS.api.ping_interval,
981        enable_http: SETTINGS.api.enable_http,
982        enable_ws: SETTINGS.api.enable_ws,
983        max_bytecode_size: MAX_BYTECODE_LENGTH,
984        max_op_datastore_entry_count: MAX_OPERATION_DATASTORE_ENTRY_COUNT,
985        max_op_datastore_key_length: MAX_OPERATION_DATASTORE_KEY_LENGTH,
986        max_op_datastore_value_length: MAX_OPERATION_DATASTORE_VALUE_LENGTH,
987        max_gas_per_block: MAX_GAS_PER_BLOCK,
988        max_serialized_operation_size: MAX_BLOCK_SIZE as usize,
989        base_operation_gas_cost: BASE_OPERATION_GAS_COST,
990        sp_compilation_cost: gas_costs.sp_compilation_cost,
991        max_function_name_length: MAX_FUNCTION_NAME_LENGTH,
992        max_parameter_size: MAX_PARAMETERS_SIZE,
993        thread_count: THREAD_COUNT,
994        keypair: keypair.clone(),
995        genesis_timestamp: *GENESIS_TIMESTAMP,
996        t0: T0,
997        periods_per_cycle: PERIODS_PER_CYCLE,
998        last_start_period: final_state.read().get_last_start_period(),
999        chain_id: *CHAINID,
1000        deferred_credits_delta: SETTINGS.api.deferred_credits_delta,
1001        minimal_fees: SETTINGS.pool.minimal_fees,
1002        deferred_calls_config,
1003        max_datastore_keys_queries: SETTINGS.api.max_datastore_keys_query,
1004        max_event_per_query: SETTINGS.execution.max_event_per_query as u32,
1005        query_state_deadline_ms: SETTINGS.execution.query_state_deadline_ms,
1006        max_datastore_key_length: MAX_DATASTORE_KEY_LENGTH,
1007        max_addresses_datastore_keys_query: SETTINGS.api.max_addresses_datastore_keys_query,
1008        pool_api_timeout: SETTINGS.api.pool_api_timeout,
1009    };
1010
1011    // spawn Massa API
1012    let api = API::<ApiV2>::new(
1013        consensus_controller.clone(),
1014        consensus_channels.broadcasts.clone(),
1015        execution_controller.clone(),
1016        pool_channels.broadcasts.clone(),
1017        api_config.clone(),
1018        *VERSION,
1019    );
1020    let api_handle = api
1021        .serve(&SETTINGS.api.bind_api, &api_config)
1022        .await
1023        .expect("failed to start MASSA API");
1024
1025    info!(
1026        "API | EXPERIMENTAL JsonRPC | listening on: {}",
1027        &SETTINGS.api.bind_api
1028    );
1029
1030    // Disable WebSockets for Private and Public API's
1031    let mut api_config = api_config.clone();
1032    api_config.enable_ws = false;
1033
1034    // Whether to spawn gRPC PUBLIC API
1035    let grpc_public_handle = if SETTINGS.grpc.public.enabled {
1036        let grpc_public_config = configure_grpc(
1037            ServiceName::Public,
1038            &SETTINGS.grpc.public,
1039            keypair.clone(),
1040            &final_state,
1041            SETTINGS.pool.minimal_fees,
1042            gas_costs.sp_compilation_cost,
1043            SETTINGS.execution.max_event_per_query,
1044            SETTINGS.execution.query_state_deadline_ms,
1045        );
1046
1047        let grpc_public_api = MassaPublicGrpc {
1048            consensus_controller: consensus_controller.clone(),
1049            consensus_broadcasts: consensus_channels.broadcasts.clone(),
1050            execution_controller: execution_controller.clone(),
1051            execution_channels,
1052            pool_broadcasts: pool_channels.broadcasts.clone(),
1053            pool_controller: pool_controller.clone(),
1054            protocol_controller: protocol_controller.clone(),
1055            selector_controller: selector_controller.clone(),
1056            storage: shared_storage.clone(),
1057            grpc_config: grpc_public_config.clone(),
1058            protocol_config: protocol_config.clone(),
1059            node_id,
1060            version: *VERSION,
1061            keypair_factory: KeyPairFactory {
1062                mip_store: mip_store.clone(),
1063            },
1064        };
1065
1066        // Spawn gRPC PUBLIC API
1067        let grpc_public_stop_handle = grpc_public_api
1068            .serve(&grpc_public_config)
1069            .await
1070            .expect("failed to start gRPC PUBLIC API");
1071        info!("gRPC | PUBLIC | listening on: {}", grpc_public_config.bind);
1072
1073        Some(grpc_public_stop_handle)
1074    } else {
1075        None
1076    };
1077
1078    // Whether to spawn gRPC PRIVATE API
1079    let grpc_private_handle = if SETTINGS.grpc.private.enabled {
1080        let grpc_private_config = configure_grpc(
1081            ServiceName::Private,
1082            &SETTINGS.grpc.private,
1083            keypair.clone(),
1084            &final_state,
1085            SETTINGS.pool.minimal_fees,
1086            gas_costs.sp_compilation_cost,
1087            SETTINGS.execution.max_event_per_query,
1088            SETTINGS.execution.query_state_deadline_ms,
1089        );
1090
1091        let bs_white_black_list = bootstrap_manager
1092            .as_ref()
1093            .map(|manager| manager.white_black_list.clone());
1094
1095        let grpc_private_api = MassaPrivateGrpc {
1096            consensus_controller: consensus_controller.clone(),
1097            execution_controller: execution_controller.clone(),
1098            pool_controller: pool_controller.clone(),
1099            protocol_controller: protocol_controller.clone(),
1100            grpc_config: grpc_private_config.clone(),
1101            protocol_config: protocol_config.clone(),
1102            node_id,
1103            mip_store: mip_store.clone(),
1104            version: *VERSION,
1105            stop_cv: sig_int_toggled.clone(),
1106            node_wallet: node_wallet.clone(),
1107            bs_white_black_list,
1108        };
1109
1110        // Spawn gRPC PRIVATE API
1111        let grpc_private_stop_handle = grpc_private_api
1112            .serve(&grpc_private_config)
1113            .await
1114            .expect("failed to start gRPC PRIVATE API");
1115        info!(
1116            "gRPC | PRIVATE | listening on: {}",
1117            grpc_private_config.bind
1118        );
1119
1120        Some(grpc_private_stop_handle)
1121    } else {
1122        None
1123    };
1124
1125    #[cfg(feature = "op_spammer")]
1126    start_operation_injector(
1127        *GENESIS_TIMESTAMP,
1128        shared_storage.clone_without_refs(),
1129        node_wallet.read().clone(),
1130        pool_controller.clone(),
1131        protocol_controller.clone(),
1132        args.nb_op,
1133    );
1134
1135    // spawn private API
1136    let api_private = API::<Private>::new(
1137        protocol_controller.clone(),
1138        execution_controller.clone(),
1139        api_config.clone(),
1140        sig_int_toggled,
1141        node_wallet,
1142    );
1143    let api_private_handle = api_private
1144        .serve(&SETTINGS.api.bind_private, &api_config)
1145        .await
1146        .expect("failed to start PRIVATE API");
1147    info!(
1148        "API | PRIVATE JsonRPC | listening on: {}",
1149        api_config.bind_private
1150    );
1151
1152    // spawn public API
1153    let api_public = API::<Public>::new(
1154        consensus_controller.clone(),
1155        execution_controller.clone(),
1156        api_config.clone(),
1157        selector_controller.clone(),
1158        pool_controller.clone(),
1159        protocol_controller.clone(),
1160        protocol_config.clone(),
1161        *VERSION,
1162        node_id,
1163        shared_storage.clone(),
1164        mip_store.clone(),
1165    );
1166    let api_public_handle = api_public
1167        .serve(&SETTINGS.api.bind_public, &api_config)
1168        .await
1169        .expect("failed to start PUBLIC API");
1170    info!(
1171        "API | PUBLIC JsonRPC | listening on: {}",
1172        api_config.bind_public
1173    );
1174
1175    let massa_survey_stopper = MassaSurvey::run(
1176        SETTINGS.metrics.tick_delay.to_duration(),
1177        execution_controller,
1178        pool_controller,
1179        db.clone(),
1180        massa_metrics,
1181        (
1182            api_config.thread_count,
1183            api_config.t0,
1184            api_config.genesis_timestamp,
1185            api_config.periods_per_cycle,
1186            api_config.last_start_period,
1187        ),
1188        mip_store,
1189    );
1190
1191    #[cfg(feature = "deadlock_detection")]
1192    {
1193        // only for #[cfg]
1194        use parking_lot::deadlock;
1195        use std::thread;
1196
1197        let interval = Duration::from_secs(args.dl_interval);
1198        warn!("deadlocks detector will run every {:?}", interval);
1199
1200        // Create a background thread which checks for deadlocks at the defined interval
1201        let thread_builder = thread::Builder::new().name("deadlock-detection".into());
1202        thread_builder
1203            .spawn(move || loop {
1204                thread::sleep(interval);
1205                let deadlocks = deadlock::check_deadlock();
1206                if deadlocks.is_empty() {
1207                    continue;
1208                }
1209                warn!("{} deadlocks detected", deadlocks.len());
1210                for (i, threads) in deadlocks.iter().enumerate() {
1211                    warn!("Deadlock #{}", i);
1212                    for t in threads {
1213                        warn!("Thread Id {:#?}", t.thread_id());
1214                        warn!("{:#?}", t.backtrace());
1215                    }
1216                }
1217            })
1218            .expect("failed to spawn thread : deadlock-detection");
1219    }
1220    (
1221        consensus_event_receiver,
1222        bootstrap_manager,
1223        consensus_manager,
1224        execution_manager,
1225        selector_manager,
1226        pool_manager,
1227        protocol_manager,
1228        factory_manager,
1229        event_cache_manager,
1230        api_private_handle,
1231        api_public_handle,
1232        api_handle,
1233        grpc_private_handle,
1234        grpc_public_handle,
1235        metrics_stopper,
1236        massa_survey_stopper,
1237    )
1238}
1239
1240// Get the configuration of the gRPC server
1241#[allow(clippy::too_many_arguments)]
1242fn configure_grpc(
1243    name: ServiceName,
1244    settings: &GrpcSettings,
1245    keypair: KeyPair,
1246    final_state: &Arc<RwLock<dyn FinalStateController>>,
1247    minimal_fees: Amount,
1248    sp_compilation_cost: u64,
1249    max_event_per_query: usize,
1250    query_state_deadline_ms: Option<u64>,
1251) -> GrpcConfig {
1252    GrpcConfig {
1253        name,
1254        enabled: settings.enabled,
1255        accept_http1: settings.accept_http1,
1256        enable_cors: settings.enable_cors,
1257        enable_health: settings.enable_health,
1258        enable_reflection: settings.enable_reflection,
1259        enable_tls: settings.enable_tls,
1260        enable_mtls: settings.enable_mtls,
1261        generate_self_signed_certificates: settings.generate_self_signed_certificates,
1262        subject_alt_names: settings.subject_alt_names.clone(),
1263        bind: settings.bind,
1264        accept_compressed: settings.accept_compressed.clone(),
1265        send_compressed: settings.send_compressed.clone(),
1266        max_decoding_message_size: settings.max_decoding_message_size,
1267        max_encoding_message_size: settings.max_encoding_message_size,
1268        concurrency_limit_per_connection: settings.concurrency_limit_per_connection,
1269        timeout: settings.timeout.to_duration(),
1270        initial_stream_window_size: settings.initial_stream_window_size,
1271        initial_connection_window_size: settings.initial_connection_window_size,
1272        max_concurrent_streams: settings.max_concurrent_streams,
1273        max_arguments: settings.max_arguments,
1274        tcp_keepalive: settings.tcp_keepalive.map(|t| t.to_duration()),
1275        tcp_nodelay: settings.tcp_nodelay,
1276        http2_keepalive_interval: settings.http2_keepalive_interval.map(|t| t.to_duration()),
1277        http2_keepalive_timeout: settings.http2_keepalive_timeout.map(|t| t.to_duration()),
1278        http2_adaptive_window: settings.http2_adaptive_window,
1279        max_frame_size: settings.max_frame_size,
1280        thread_count: THREAD_COUNT,
1281        max_operations_per_block: MAX_OPERATIONS_PER_BLOCK,
1282        endorsement_count: ENDORSEMENT_COUNT,
1283        max_endorsements_per_message: MAX_ENDORSEMENTS_PER_MESSAGE,
1284        max_bytecode_size: MAX_BYTECODE_LENGTH,
1285        max_op_datastore_entry_count: MAX_OPERATION_DATASTORE_ENTRY_COUNT,
1286        max_datastore_entries_per_request: settings.max_datastore_entries_per_request,
1287        max_op_datastore_key_length: MAX_OPERATION_DATASTORE_KEY_LENGTH,
1288        max_op_datastore_value_length: MAX_OPERATION_DATASTORE_VALUE_LENGTH,
1289        max_function_name_length: MAX_FUNCTION_NAME_LENGTH,
1290        max_parameter_size: MAX_PARAMETERS_SIZE,
1291        max_operations_per_message: MAX_OPERATIONS_PER_MESSAGE,
1292        max_gas_per_block: MAX_GAS_PER_BLOCK,
1293        max_serialized_operation_size: MAX_BLOCK_SIZE as usize,
1294        base_operation_gas_cost: BASE_OPERATION_GAS_COST,
1295        sp_compilation_cost,
1296        genesis_timestamp: *GENESIS_TIMESTAMP,
1297        t0: T0,
1298        periods_per_cycle: PERIODS_PER_CYCLE,
1299        keypair,
1300        max_channel_size: settings.max_channel_size,
1301        draw_lookahead_period_count: settings.draw_lookahead_period_count,
1302        last_start_period: final_state.read().get_last_start_period(),
1303        max_denunciations_per_block_header: MAX_DENUNCIATIONS_PER_BLOCK_HEADER,
1304        max_addresses_per_request: settings.max_addresses_per_request,
1305        max_slot_ranges_per_request: settings.max_slot_ranges_per_request,
1306        max_block_ids_per_request: settings.max_block_ids_per_request,
1307        max_endorsement_ids_per_request: settings.max_endorsement_ids_per_request,
1308        max_operation_ids_per_request: settings.max_operation_ids_per_request,
1309        max_filters_per_request: settings.max_filters_per_request,
1310        max_query_items_per_request: settings.max_query_items_per_request,
1311        max_event_per_query: max_event_per_query as u32,
1312        query_state_deadline_ms,
1313        certificate_authority_root_path: settings.certificate_authority_root_path.clone(),
1314        server_certificate_path: settings.server_certificate_path.clone(),
1315        server_private_key_path: settings.server_private_key_path.clone(),
1316        client_certificate_authority_root_path: settings
1317            .client_certificate_authority_root_path
1318            .clone(),
1319        client_certificate_path: settings.client_certificate_path.clone(),
1320        client_private_key_path: settings.client_private_key_path.clone(),
1321        chain_id: *CHAINID,
1322        minimal_fees,
1323        max_datastore_keys_queries: settings.max_datastore_keys_query,
1324        max_datastore_key_length: MAX_DATASTORE_KEY_LENGTH,
1325        unidirectional_stream_interval_check: settings.unidirectional_stream_interval_check,
1326    }
1327}
1328
1329struct Managers {
1330    bootstrap_manager: Option<BootstrapManager>,
1331    consensus_manager: Box<dyn ConsensusManager>,
1332    execution_manager: Box<dyn ExecutionManager>,
1333    selector_manager: Box<dyn SelectorManager>,
1334    pool_manager: Box<dyn PoolManager>,
1335    protocol_manager: Box<dyn ProtocolManager>,
1336    factory_manager: Box<dyn FactoryManager>,
1337    event_cache_manager: Box<dyn EventCacheManager>,
1338}
1339
1340#[allow(clippy::too_many_arguments)]
1341async fn stop(
1342    _consensus_event_receiver: MassaReceiver<ConsensusEvent>,
1343    Managers {
1344        bootstrap_manager,
1345        mut execution_manager,
1346        mut consensus_manager,
1347        mut selector_manager,
1348        mut pool_manager,
1349        mut protocol_manager,
1350        mut factory_manager,
1351        mut event_cache_manager,
1352    }: Managers,
1353    api_private_handle: StopHandle,
1354    api_public_handle: StopHandle,
1355    api_handle: StopHandle,
1356    grpc_private_handle: Option<massa_grpc::server::StopHandle>,
1357    grpc_public_handle: Option<massa_grpc::server::StopHandle>,
1358    mut metrics_stopper: MetricsStopper,
1359    mut massa_survey_stopper: MassaSurveyStopper,
1360) {
1361    // stop bootstrap
1362    if let Some(bootstrap_manager) = bootstrap_manager {
1363        bootstrap_manager
1364            .stop()
1365            .expect("bootstrap server shutdown failed")
1366    }
1367
1368    info!("Start stopping API's: gRPC(PUBLIC, PRIVATE), EXPERIMENTAL, PUBLIC, PRIVATE");
1369
1370    // stop Massa gRPC PUBLIC API
1371    if let Some(handle) = grpc_public_handle {
1372        handle.stop();
1373    }
1374    info!("API | PUBLIC gRPC | stopped");
1375
1376    // stop Massa gRPC PRIVATE API
1377    if let Some(handle) = grpc_private_handle {
1378        handle.stop();
1379    }
1380    info!("API | PRIVATE gRPC | stopped");
1381
1382    // stop Massa API
1383    api_handle.stop().await;
1384    info!("API | EXPERIMENTAL JsonRPC | stopped");
1385
1386    // stop public API
1387    api_public_handle.stop().await;
1388    info!("API | PUBLIC JsonRPC | stopped");
1389
1390    // stop private API
1391    api_private_handle.stop().await;
1392    info!("API | PRIVATE JsonRPC | stopped");
1393
1394    // stop metrics
1395    metrics_stopper.stop();
1396
1397    // stop massa survey thread
1398    massa_survey_stopper.stop();
1399
1400    // stop factory
1401    factory_manager.stop();
1402
1403    // stop protocol controller
1404    protocol_manager.stop();
1405
1406    // stop consensus
1407    consensus_manager.stop();
1408
1409    // stop pool
1410    pool_manager.stop();
1411
1412    // stop execution controller
1413    execution_manager.stop();
1414
1415    // stop selector controller
1416    selector_manager.stop();
1417
1418    // stop pool controller
1419    // TODO
1420    //let protocol_pool_event_receiver = pool_manager.stop().await.expect("pool shutdown failed");
1421
1422    // note that FinalLedger gets destroyed as soon as its Arc count goes to zero
1423
1424    event_cache_manager.stop();
1425}
1426
1427#[derive(Parser)]
1428#[command(version = crate_version!())]
1429struct Args {
1430    #[arg(long = "keep-ledger")]
1431    keep_ledger: bool,
1432    /// Wallet password
1433    // #[arg(short = "p", long = "pwd")]
1434    #[arg(short = 'p', long = "pwd")]
1435    password: Option<String>,
1436
1437    /// restart_from_snapshot_at_period
1438    #[arg(long = "restart-from-snapshot-at-period")]
1439    restart_from_snapshot_at_period: Option<u64>,
1440
1441    #[arg(short = 'a', long = "accept-community-charter")]
1442    accept_community_charter: bool,
1443
1444    #[cfg(feature = "op_spammer")]
1445    /// number of operations
1446    #[arg(
1447        name = "number of operations",
1448        long_help = "Define the number of operations the node can spam.",
1449        short = 'n',
1450        long = "number-operations"
1451    )]
1452    nb_op: u64,
1453
1454    #[cfg(feature = "deadlock_detection")]
1455    /// Deadlocks detector
1456    #[structopt(
1457        name = "deadlocks interval",
1458        long_help = "Define the interval of launching a deadlocks checking.",
1459        short = 'i',
1460        long = "dli",
1461        default_value = "10"
1462    )]
1463    dl_interval: u64,
1464}
1465
1466/// Load wallet, asking for passwords if necessary
1467fn load_wallet(
1468    password: Option<String>,
1469    path: &Path,
1470    chain_id: u64,
1471) -> anyhow::Result<Arc<RwLock<Wallet>>> {
1472    let password = if path.is_dir() {
1473        password.unwrap_or_else(|| {
1474            Password::new()
1475                .with_prompt("Enter staking keys file password")
1476                .interact()
1477                .expect("IO error: Password reading failed, staking keys file couldn't be unlocked")
1478        })
1479    } else {
1480        password.unwrap_or_else(|| {
1481            Password::new()
1482                .with_prompt("Enter new password for staking keys file")
1483                .with_confirmation("Confirm password", "Passwords mismatching")
1484                .interact()
1485                .expect("IO error: Password reading failed, staking keys file couldn't be created")
1486        })
1487    };
1488    Ok(Arc::new(RwLock::new(Wallet::new(
1489        PathBuf::from(path),
1490        password,
1491        chain_id,
1492    )?)))
1493}
1494
1495fn main() -> anyhow::Result<()> {
1496    let args = Args::parse();
1497
1498    let tokio_rt = tokio::runtime::Builder::new_multi_thread()
1499        .thread_name_fn(|| {
1500            static ATOMIC_ID: AtomicUsize = AtomicUsize::new(0);
1501            let id = ATOMIC_ID.fetch_add(1, Ordering::SeqCst);
1502            format!("tokio-node-{}", id)
1503        })
1504        .enable_all()
1505        .build()
1506        .unwrap();
1507
1508    tokio_rt.block_on(run(args))
1509}
1510
1511async fn run(args: Args) -> anyhow::Result<()> {
1512    let mut cur_args = args;
1513    use tracing_subscriber::prelude::*;
1514    // spawn the console server in the background, returning a `Layer`:
1515    let tracing_layer = tracing_subscriber::fmt::layer()
1516        .with_filter(match SETTINGS.logging.level {
1517            4 => LevelFilter::TRACE,
1518            3 => LevelFilter::DEBUG,
1519            2 => LevelFilter::INFO,
1520            1 => LevelFilter::WARN,
1521            _ => LevelFilter::ERROR,
1522        })
1523        .with_filter(filter_fn(|metadata| {
1524            metadata.target().starts_with("massa") // ignore non-massa logs
1525        }));
1526    // build a `Subscriber` by combining layers with a `tracing_subscriber::Registry`:
1527    tracing_subscriber::registry()
1528        // add the console layer to the subscriber or default layers...
1529        .with(tracing_layer)
1530        .init();
1531
1532    // Setup panic handlers,
1533    // and when a panic occurs,
1534    // run default handler,
1535    // and then shutdown.
1536    let default_panic = std::panic::take_hook();
1537    std::panic::set_hook(Box::new(move |info| {
1538        default_panic(info);
1539        std::process::exit(1);
1540    }));
1541
1542    info!("Node version : {}", *VERSION);
1543
1544    handle_disclaimer(
1545        cur_args.accept_community_charter,
1546        &SETTINGS.cli.approved_community_charter_file_path,
1547    );
1548
1549    // load or create wallet, asking for password if necessary
1550    let node_wallet = load_wallet(
1551        cur_args.password.clone(),
1552        &SETTINGS.factory.staking_wallet_path,
1553        *CHAINID,
1554    )?;
1555
1556    // interrupt signal listener
1557    let sig_int_toggled = Arc::new((Mutex::new(false), Condvar::new()));
1558
1559    let sig_int_toggled_clone = Arc::clone(&sig_int_toggled);
1560    ctrlc::set_handler(move || {
1561        *sig_int_toggled_clone
1562            .0
1563            .lock()
1564            .expect("double-lock on interrupt bool in ctrl-c handler") = true;
1565        sig_int_toggled_clone.1.notify_all();
1566    })
1567    .expect("Error setting Ctrl-C handler");
1568
1569    #[cfg(feature = "resync_check")]
1570    let mut resync_check = Some(std::time::Instant::now() + std::time::Duration::from_secs(10));
1571
1572    loop {
1573        let (
1574            consensus_event_receiver,
1575            bootstrap_manager,
1576            consensus_manager,
1577            execution_manager,
1578            selector_manager,
1579            pool_manager,
1580            protocol_manager,
1581            factory_manager,
1582            event_cache_manager,
1583            api_private_handle,
1584            api_public_handle,
1585            api_handle,
1586            grpc_private_handle,
1587            grpc_public_handle,
1588            metrics_stopper,
1589            massa_survey_stopper,
1590        ) = launch(&cur_args, node_wallet.clone(), Arc::clone(&sig_int_toggled)).await;
1591
1592        // loop over messages
1593        let restart = loop {
1594            massa_trace!("massa-node.main.run.select", {});
1595            match consensus_event_receiver.try_recv() {
1596                Ok(evt) => match evt {
1597                    ConsensusEvent::NeedSync => {
1598                        warn!("in response to a desynchronization, the node is going to bootstrap again");
1599                        break true;
1600                    }
1601                    ConsensusEvent::Stop => {
1602                        break false;
1603                    }
1604                },
1605                Err(TryRecvError::Disconnected) => {
1606                    error!("consensus_event_receiver.wait_event disconnected");
1607                    break false;
1608                }
1609                _ => {}
1610            };
1611
1612            // every 100ms/or when alerted, check if sigint toggled
1613            // if toggled, break loop
1614            let int_sig = sig_int_toggled
1615                .0
1616                .lock()
1617                .expect("double-lock() on interrupted signal mutex");
1618            let wake = sig_int_toggled
1619                .1
1620                .wait_timeout(int_sig, Duration::from_millis(100))
1621                .expect("interrupt signal mutex poisoned");
1622            if *wake.0 {
1623                info!("interrupt signal received");
1624                break false;
1625            }
1626
1627            // Elements of the system that involve stopping and restarting should be checked by forcing a relaunch.
1628            // This check allows the system to start up as normal, wait 10s, then force a relaunch. If Things take too long
1629            // to shutdown, or does not allow for a clean relaunch, this feature flag can expose those issues.
1630            #[cfg(feature = "resync_check")]
1631            if let Some(resync_moment) = resync_check {
1632                if resync_moment < std::time::Instant::now() {
1633                    warn!("resync check triggered");
1634                    resync_check = None;
1635                    break true;
1636                }
1637            }
1638        };
1639        stop(
1640            consensus_event_receiver,
1641            Managers {
1642                bootstrap_manager,
1643                consensus_manager,
1644                execution_manager,
1645                selector_manager,
1646                pool_manager,
1647                protocol_manager,
1648                factory_manager,
1649                event_cache_manager,
1650            },
1651            api_private_handle,
1652            api_public_handle,
1653            api_handle,
1654            grpc_private_handle,
1655            grpc_public_handle,
1656            metrics_stopper,
1657            massa_survey_stopper,
1658        )
1659        .await;
1660
1661        if !restart {
1662            break;
1663        }
1664        // If we restart because of a desync, then we do not want to restart from a snapshot
1665        cur_args.restart_from_snapshot_at_period = None;
1666    }
1667    Ok(())
1668}