massa_node/
survey.rs

1#![allow(unused_imports)]
2use std::thread::JoinHandle;
3
4use crossbeam_channel::{select, tick};
5use massa_channel::{sender::MassaSender, MassaChannel};
6use massa_db_exports::ShareableMassaDBController;
7use massa_execution_exports::ExecutionController;
8use massa_metrics::MassaMetrics;
9use massa_models::{address::Address, slot::Slot, timeslots::get_latest_block_slot_at_timestamp};
10use massa_pool_exports::PoolController;
11use massa_time::MassaTime;
12use massa_versioning::versioning::MipStore;
13use tracing::info;
14use tracing::warn;
15
16pub struct MassaSurvey {}
17
18pub struct MassaSurveyStopper {
19    tx_stopper: Option<MassaSender<()>>,
20    handle: Option<JoinHandle<()>>,
21}
22
23impl MassaSurveyStopper {
24    pub fn stop(&mut self) {
25        if let Some(tx) = self.tx_stopper.take() {
26            info!("MassaSurvey | Stopping");
27            if let Err(e) = tx.send(()) {
28                warn!("failed to send stop signal to massa survey thread: {:?}", e);
29            }
30        }
31        if let Some(handle) = self.handle.take() {
32            match handle.join() {
33                Ok(_) => info!("MassaSurvey | Stopped"),
34                Err(_) => warn!("failed to join massa survey thread"),
35            }
36        }
37    }
38}
39
40impl MassaSurvey {
41    #[allow(unused_variables)]
42    // config : (thread_count, t0, genesis_timestamp, periods_per_cycle, last_start_period)
43    pub fn run(
44        tick_delay: std::time::Duration,
45        execution_controller: Box<dyn ExecutionController>,
46        pool_controller: Box<dyn PoolController>,
47        db: ShareableMassaDBController,
48        massa_metrics: MassaMetrics,
49        config: (u8, MassaTime, MassaTime, u64, u64),
50        mip_store: MipStore,
51    ) -> MassaSurveyStopper {
52        if massa_metrics.is_enabled() {
53            #[cfg(all(not(feature = "sandbox"), not(test)))]
54            {
55                // massa-survey
56                const THREAD_NAME: &str = "massa-survey";
57
58                let mut data_sent = 0;
59                let mut data_received = 0;
60                let (tx_stop, rx_stop) =
61                    MassaChannel::new("massa_survey_stop".to_string(), Some(1));
62                let update_tick = tick(tick_delay);
63                match std::thread::Builder::new()
64                    .name(THREAD_NAME.to_string())
65                    .spawn(move || loop {
66                        select! {
67                            recv(rx_stop) -> _ => {
68                                break;
69                            },
70                            recv(update_tick) -> _ => {
71                                let (
72                                    active_in_connections,
73                                    active_out_connections,
74                                    new_data_sent,
75                                    new_data_received,
76                                ) = massa_metrics.get_metrics_for_survey_thread();
77
78                                if active_in_connections + active_out_connections == 0 {
79                                    warn!("PEERNET | No active connections");
80                                }
81
82                                if new_data_sent == data_sent && new_data_received == data_received {
83                                    let now = MassaTime::now();
84                                    if now > config.2 {
85                                        warn!("PEERNET | No data sent or received since 5s");
86                                    }
87                                } else {
88                                    data_sent = new_data_sent;
89                                    data_received = new_data_received;
90                                }
91
92                                {
93                                           // update stakers / rolls
94                                    let now = MassaTime::now();
95
96                                    let curr_cycle =
97                                        match get_latest_block_slot_at_timestamp(config.0, config.1, config.2, now)
98                                        {
99                                            Ok(Some(cur_slot)) if cur_slot.period <= config.4 => {
100                                                Slot::new(config.4, 0).get_cycle(config.3)
101                                            }
102                                            Ok(Some(cur_slot)) => cur_slot.get_cycle(config.3),
103                                            Ok(None) => 0,
104                                            Err(e) => {
105                                                warn!(
106                                                "MassaSurvey | Failed to get latest block slot at timestamp: {:?}",
107                                                e
108                                            );
109                                                continue;
110                                            }
111                                        };
112
113                                    let staker_vec = execution_controller
114                                        .get_cycle_active_rolls(curr_cycle)
115                                        .into_iter()
116                                        .collect::<Vec<(Address, u64)>>();
117
118                                    massa_metrics.set_stakers(staker_vec.len());
119                                    let rolls_count = staker_vec.iter().map(|(_, r)| *r).sum::<u64>();
120                                    massa_metrics.set_rolls(rolls_count as usize);
121                                    let current_slot = get_latest_block_slot_at_timestamp(config.0, config.1, config.2, now).unwrap_or(None).unwrap_or(Slot::new(0, 0));
122                                    massa_metrics.set_current_time_thread(current_slot.thread);
123                                    massa_metrics.set_current_time_period(current_slot.period);
124                                }
125
126                                {
127                                    massa_metrics.set_operations_pool(pool_controller.get_operation_count(None).unwrap_or(0));
128                                    massa_metrics.set_endorsements_pool(pool_controller.get_endorsement_count(None).unwrap_or(0));
129                                    massa_metrics.set_denunciations_pool(pool_controller.get_denunciation_count(None).unwrap_or(0));
130
131                                    let count = std::thread::available_parallelism()
132                                    .unwrap_or(std::num::NonZeroUsize::MIN)
133                                    .get();
134                                    massa_metrics.set_available_processors(count);
135                                }
136
137                                {
138                                    massa_metrics.set_network_current_version(mip_store.get_network_version_current());
139                                    let network_stats= mip_store.0.read().get_network_versions_stats();
140                                    massa_metrics.update_network_version_vote(network_stats);
141                                }
142
143                                {
144                                    let change_history_sizes = db.read().get_change_history_sizes();
145                                    massa_metrics.set_db_change_history_size(change_history_sizes.0);
146                                    massa_metrics.set_db_change_versioning_history_size(change_history_sizes.1);
147                                }
148
149                                {
150                                    let module_lru_cache_memory_usage = execution_controller.get_module_lru_cache_memory_usage();
151                                    massa_metrics.set_module_lru_cache_memory_usage(module_lru_cache_memory_usage);
152                                }
153
154                                {
155                                    let active_history_total_event_len = execution_controller.get_active_history_total_event_len();
156                                    massa_metrics.set_active_history_total_event_len(active_history_total_event_len);
157                                }
158                            }
159                        }
160                    }) {
161                    Ok(handle) => MassaSurveyStopper { handle: Some(handle), tx_stopper: Some(tx_stop) },
162                    Err(e) => {
163                        warn!("MassaSurvey | Failed to spawn survey thread: {:?}", e);
164                        MassaSurveyStopper { handle: None, tx_stopper: None}
165                    }
166                }
167            }
168
169            #[cfg(any(feature = "sandbox", test))]
170            {
171                MassaSurveyStopper {
172                    handle: None,
173                    tx_stopper: None,
174                }
175            }
176        } else {
177            MassaSurveyStopper {
178                handle: None,
179                tx_stopper: None,
180            }
181        }
182    }
183}