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 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 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 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}