massa_metrics/
lib.rs

1//! this library is used to collect metrics from the node and expose them to the prometheus server
2//!
3//! the metrics are collected from the node and from the survey
4//! the survey is a separate thread that is used to collect metrics from the network (active connections)
5//!
6
7use std::{
8    collections::HashMap,
9    net::SocketAddr,
10    sync::{Arc, RwLock},
11    thread::JoinHandle,
12    time::Duration,
13};
14
15use lazy_static::lazy_static;
16use prometheus::{
17    register_int_counter, register_int_gauge, Gauge, Histogram, IntCounter, IntGauge,
18};
19use tokio::sync::oneshot::Sender;
20use tracing::warn;
21
22mod server;
23
24lazy_static! {
25    // use lazy_static for these metrics because they are used in storage which implement default
26    static ref OPERATIONS_COUNTER: IntGauge = register_int_gauge!(
27        "operations_storage_counter",
28        "operations storage counter len"
29    )
30    .unwrap();
31    static ref BLOCKS_COUNTER: IntGauge =
32        register_int_gauge!("blocks_storage_counter", "blocks storage counter len").unwrap();
33    static ref ENDORSEMENTS_COUNTER: IntGauge =
34        register_int_gauge!("endorsements_storage_counter", "endorsements storage counter len").unwrap();
35
36        static ref DEFERRED_CALL_REGISTERED: IntGauge = register_int_gauge!(
37        "deferred_calls_registered", "number of deferred calls registered" ).unwrap();
38
39        static ref DEFERRED_CALL_CANCELLED: IntGauge = register_int_gauge!(
40            "deferred_calls_cancelled", "number of deferred calls cancelled" ).unwrap();
41
42        static ref DEFERRED_CALL_EXECUTED: IntGauge = register_int_gauge!(
43            "deferred_calls_executed", "number of deferred calls executed" ).unwrap();
44
45        static ref DEFERRED_CALL_FAILED: IntGauge = register_int_gauge!(
46            "deferred_calls_failed", "number of deferred calls failed" ).unwrap();
47
48        static ref DEFERRED_CALLS_TOTAL_GAS: IntGauge = register_int_gauge!(
49            "deferred_calls_total_gas", "total gas used by deferred calls" ).unwrap();
50
51        static ref EVENT_CACHE_VEC: IntGauge = register_int_gauge!(
52            "event_cache_vec_len", "vector len for events" ).unwrap();
53
54        static ref FACTORY_POOL_TIMEOUT_ENDORSEMENTS: IntCounter = register_int_counter!(
55            "factory_pool_timeout_endorsements",
56            "number of times block factory timed out while fetching endorsements from the pool"
57        ).unwrap();
58
59        static ref FACTORY_POOL_TIMEOUT_OPERATIONS: IntCounter = register_int_counter!(
60            "factory_pool_timeout_operations",
61            "number of times block factory timed out while fetching operations from the pool"
62        ).unwrap();
63
64        static ref FACTORY_POOL_TIMEOUT_DENUNCIATIONS: IntCounter = register_int_counter!(
65            "factory_pool_timeout_denunciations",
66            "number of times block factory timed out while fetching denunciations from the pool"
67        ).unwrap();
68
69        static ref FACTORY_BLOCKS_PRODUCED: IntCounter = register_int_counter!(
70            "factory_blocks_produced",
71            "number of blocks produced by the local block factory"
72        ).unwrap();
73
74        static ref FACTORY_BLOCKS_PRODUCED_WITH_POOL_TIMEOUT: IntCounter = register_int_counter!(
75            "factory_blocks_produced_with_pool_timeout",
76            "number of locally produced blocks where at least one pool read timed out during assembly"
77        ).unwrap();
78
79}
80
81pub fn inc_deferred_calls_registered() {
82    DEFERRED_CALL_REGISTERED.inc();
83}
84
85pub fn set_event_cache_vec_len(val: usize) {
86    EVENT_CACHE_VEC.set(val as i64);
87}
88
89pub fn set_deferred_calls_total_gas(val: u128) {
90    DEFERRED_CALLS_TOTAL_GAS.set(val as i64);
91}
92
93pub fn inc_deferred_calls_executed_by(val: u64) {
94    DEFERRED_CALL_EXECUTED.set(DEFERRED_CALL_EXECUTED.get().saturating_add(val as i64));
95}
96
97pub fn inc_deferred_calls_failed_by(val: u64) {
98    DEFERRED_CALL_FAILED.set(DEFERRED_CALL_FAILED.get().saturating_add(val as i64));
99}
100
101pub fn dec_deferred_calls_cancelled_by(val: u64) {
102    DEFERRED_CALL_CANCELLED.set(
103        DEFERRED_CALL_CANCELLED
104            .get()
105            .saturating_sub(val as i64)
106            .max(0),
107    );
108}
109
110pub fn dec_deferred_calls_registered_by(val: u64) {
111    DEFERRED_CALL_REGISTERED.set(
112        DEFERRED_CALL_REGISTERED
113            .get()
114            .saturating_sub(val as i64)
115            .max(0),
116    );
117}
118
119pub fn inc_deferred_calls_cancelled() {
120    DEFERRED_CALL_CANCELLED.inc();
121}
122
123pub fn set_blocks_counter(val: usize) {
124    BLOCKS_COUNTER.set(val as i64);
125}
126
127pub fn set_endorsements_counter(val: usize) {
128    ENDORSEMENTS_COUNTER.set(val as i64);
129}
130
131pub fn set_operations_counter(val: usize) {
132    OPERATIONS_COUNTER.set(val as i64);
133}
134
135pub fn inc_factory_pool_timeout_endorsements() {
136    FACTORY_POOL_TIMEOUT_ENDORSEMENTS.inc();
137}
138
139pub fn inc_factory_pool_timeout_operations() {
140    FACTORY_POOL_TIMEOUT_OPERATIONS.inc();
141}
142
143pub fn inc_factory_pool_timeout_denunciations() {
144    FACTORY_POOL_TIMEOUT_DENUNCIATIONS.inc();
145}
146
147pub fn inc_factory_blocks_produced(had_pool_timeout: bool) {
148    FACTORY_BLOCKS_PRODUCED.inc();
149    if had_pool_timeout {
150        FACTORY_BLOCKS_PRODUCED_WITH_POOL_TIMEOUT.inc();
151    }
152}
153
154/// Eagerly register lazy_static metrics so they appear on `/metrics` even before
155/// the first increment (e.g. factory pool timeouts that have not occurred yet).
156#[cfg(not(feature = "test-exports"))]
157fn register_lazy_static_metrics() {
158    lazy_static::initialize(&FACTORY_POOL_TIMEOUT_ENDORSEMENTS);
159    lazy_static::initialize(&FACTORY_POOL_TIMEOUT_OPERATIONS);
160    lazy_static::initialize(&FACTORY_POOL_TIMEOUT_DENUNCIATIONS);
161    lazy_static::initialize(&FACTORY_BLOCKS_PRODUCED);
162    lazy_static::initialize(&FACTORY_BLOCKS_PRODUCED_WITH_POOL_TIMEOUT);
163}
164
165#[derive(Default)]
166pub struct MetricsStopper {
167    pub(crate) stopper: Option<Sender<()>>,
168    pub(crate) stop_handle: Option<JoinHandle<()>>,
169}
170
171impl MetricsStopper {
172    pub fn stop(&mut self) {
173        if let Some(stopper) = self.stopper.take() {
174            if stopper.send(()).is_err() {
175                warn!("failed to send stop signal to metrics server");
176            }
177
178            if let Some(handle) = self.stop_handle.take() {
179                if let Err(_e) = handle.join() {
180                    warn!("failed to join metrics server thread");
181                }
182            }
183        }
184    }
185}
186
187#[derive(Clone)]
188pub struct MassaMetrics {
189    /// enable metrics
190    enabled: bool,
191
192    /// number of processors
193    process_available_processors: IntGauge,
194
195    /// consensus period for each thread
196    /// index 0 = thread 0 ...
197    consensus_vec: Vec<Gauge>,
198
199    /// number of stakers
200    stakers: IntGauge,
201    /// number of rolls
202    rolls: IntGauge,
203
204    // thread of actual slot
205    current_time_thread: IntGauge,
206    // period of actual slot
207    current_time_period: IntGauge,
208
209    /// number of elements in the active_history of execution
210    active_history: IntGauge,
211
212    /// number of operations in the operation pool
213    operations_pool: IntGauge,
214    /// number of endorsements in the endorsement pool
215    endorsements_pool: IntGauge,
216    /// number of elements in the denunciation pool
217    denunciations_pool: IntGauge,
218
219    // number of autonomous SCs messages in pool
220    async_message_pool_size: IntGauge,
221
222    // number of autonomous SC messages executed as final
223    sc_messages_final: IntCounter,
224
225    /// number of times our node (re-)bootstrapped
226    bootstrap_counter: IntCounter,
227    /// number of times we successfully bootstrapped someone
228    bootstrap_peers_success: IntCounter,
229    /// number of times we failed/refused to bootstrap someone
230    bootstrap_peers_failed: IntCounter,
231
232    /// number of times we successfully tested someone
233    protocol_tester_success: IntCounter,
234    /// number of times we failed to test someone
235    protocol_tester_failed: IntCounter,
236
237    /// know peers in protocol
238    protocol_known_peers: IntGauge,
239    /// banned peers in protocol
240    protocol_banned_peers: IntGauge,
241
242    /// executed final slot
243    executed_final_slot: IntCounter,
244    /// executed final slot with block (not miss)
245    executed_final_slot_with_block: IntCounter,
246
247    /// total bytes receive by peernet manager
248    peernet_total_bytes_received: IntCounter,
249    /// total bytes sent by peernet manager
250    peernet_total_bytes_sent: IntCounter,
251
252    /// block slot delay
253    block_slot_delay: Histogram,
254
255    /// active in connections peer
256    active_in_connections: IntGauge,
257    /// active out connections peer
258    active_out_connections: IntGauge,
259
260    /// counter of operations for final slot
261    operations_final_counter: IntCounter,
262
263    // block_cache
264    block_cache_checked_headers_size: IntGauge,
265    block_cache_blocks_known_by_peer: IntGauge,
266
267    // Operation cache
268    operation_cache_checked_operations: IntGauge,
269    operation_cache_checked_operations_prefix: IntGauge,
270    operation_cache_ops_know_by_peer: IntGauge,
271
272    // Consensus state
273    consensus_state_active_index: IntGauge,
274    consensus_state_active_index_without_ops: IntGauge,
275    consensus_state_incoming_index: IntGauge,
276    consensus_state_discarded_index: IntGauge,
277    consensus_state_block_statuses: IntGauge,
278
279    // endorsement cache
280    endorsement_cache_checked_endorsements: IntGauge,
281    endorsement_cache_known_by_peer: IntGauge,
282
283    // cursor
284    active_cursor_thread: IntGauge,
285    active_cursor_period: IntGauge,
286
287    final_cursor_thread: IntGauge,
288    final_cursor_period: IntGauge,
289
290    // peer bandwidth (bytes sent, bytes received)
291    peers_bandwidth: Arc<RwLock<HashMap<String, (IntCounter, IntCounter)>>>,
292
293    // network versions votes <version, votes>
294    network_versions_votes: Arc<RwLock<HashMap<u32, IntGauge>>>,
295    network_current_version: IntGauge,
296
297    // massa-db
298    db_change_history_size: IntGauge,
299    db_change_versioning_history_size: IntGauge,
300
301    // module cache
302    module_lru_cache_memory_usage: IntGauge,
303
304    // active history
305    active_history_total_event_len: IntGauge,
306
307    pub tick_delay: Duration,
308}
309
310impl MassaMetrics {
311    #[allow(unused_variables)]
312    #[allow(unused_mut)]
313    pub fn new(
314        enabled: bool,
315        addr: SocketAddr,
316        nb_thread: u8,
317        tick_delay: Duration,
318    ) -> (Self, MetricsStopper) {
319        let mut consensus_vec = vec![];
320        for i in 0..nb_thread {
321            let gauge = Gauge::new(
322                format!("consensus_thread_{}", i),
323                "consensus thread actual period",
324            )
325            .expect("Failed to create gauge");
326            #[cfg(not(feature = "test-exports"))]
327            {
328                let _ = prometheus::register(Box::new(gauge.clone()));
329            }
330
331            consensus_vec.push(gauge);
332        }
333
334        // set available processors
335        let process_available_processors =
336            IntGauge::new("process_available_processors", "number of processors")
337                .expect("Failed to create available_processors counter");
338
339        // stakers
340        let stakers = IntGauge::new("stakers", "number of stakers").unwrap();
341        let rolls = IntGauge::new("rolls", "number of rolls").unwrap();
342
343        let current_time_period =
344            IntGauge::new("current_time_period", "period of actual slot").unwrap();
345
346        let current_time_thread =
347            IntGauge::new("current_time_thread", "thread of actual slot").unwrap();
348
349        let executed_final_slot =
350            IntCounter::new("executed_final_slot", "number of executed final slot").unwrap();
351        let executed_final_slot_with_block = IntCounter::new(
352            "executed_final_slot_with_block",
353            "number of executed final slot with block (not miss)",
354        )
355        .unwrap();
356
357        let protocol_tester_success = IntCounter::new(
358            "protocol_tester_success",
359            "number of times we successfully tested someone",
360        )
361        .unwrap();
362        let protocol_tester_failed = IntCounter::new(
363            "protocol_tester_failed",
364            "number of times we failed to test someone",
365        )
366        .unwrap();
367
368        // pool
369        let operations_pool = IntGauge::new(
370            "operations_pool",
371            "number of operations in the operation pool",
372        )
373        .unwrap();
374        let endorsements_pool = IntGauge::new(
375            "endorsements_pool",
376            "number of endorsements in the endorsement pool",
377        )
378        .unwrap();
379        let denunciations_pool = IntGauge::new(
380            "denunciations_pool",
381            "number of elements in the denunciation pool",
382        )
383        .unwrap();
384
385        let async_message_pool_size = IntGauge::new(
386            "async_message_pool_size",
387            "number of autonomous SCs messages in pool",
388        )
389        .unwrap();
390
391        let sc_messages_final = IntCounter::new(
392            "sc_messages_final",
393            "number of autonomous SC messages executed as final",
394        )
395        .unwrap();
396
397        let bootstrap_counter = IntCounter::new(
398            "bootstrap_counter",
399            "number of times our node (re-)bootstrapped",
400        )
401        .unwrap();
402        let bootstrap_success = IntCounter::new(
403            "bootstrap_peers_success",
404            "number of times we successfully bootstrapped someone",
405        )
406        .unwrap();
407        let bootstrap_failed = IntCounter::new(
408            "bootstrap_peers_failed",
409            "number of times we failed/refused to bootstrap someone",
410        )
411        .unwrap();
412
413        let active_history = IntGauge::new(
414            "active_history",
415            "number of elements in the active_history of execution",
416        )
417        .unwrap();
418
419        let know_peers =
420            IntGauge::new("protocol_known_peers", "number of known peers in protocol").unwrap();
421        let banned_peers = IntGauge::new(
422            "protocol_banned_peers",
423            "number of banned peers in protocol",
424        )
425        .unwrap();
426
427        // active cursor
428        let active_cursor_thread =
429            IntGauge::new("active_cursor_thread", "execution active cursor thread").unwrap();
430        let active_cursor_period =
431            IntGauge::new("active_cursor_period", "execution active cursor period").unwrap();
432
433        // final cursor
434        let final_cursor_thread =
435            IntGauge::new("final_cursor_thread", "execution final cursor thread").unwrap();
436        let final_cursor_period =
437            IntGauge::new("final_cursor_period", "execution final cursor period").unwrap();
438
439        // active connections IN
440        let active_in_connections =
441            IntGauge::new("active_in_connections", "active connections IN len").unwrap();
442
443        // active connections OUT
444        let active_out_connections =
445            IntGauge::new("active_out_connections", "active connections OUT len").unwrap();
446
447        // block cache
448        let block_cache_checked_headers_size = IntGauge::new(
449            "block_cache_checked_headers_size",
450            "size of BlockCache checked_headers",
451        )
452        .unwrap();
453
454        let block_cache_blocks_known_by_peer = IntGauge::new(
455            "block_cache_blocks_known_by_peer_size",
456            "size of BlockCache blocks_known_by_peer",
457        )
458        .unwrap();
459
460        // operation cache
461        let operation_cache_checked_operations = IntGauge::new(
462            "operation_cache_checked_operations",
463            "size of OperationCache checked_operations",
464        )
465        .unwrap();
466
467        let operation_cache_checked_operations_prefix = IntGauge::new(
468            "operation_cache_checked_operations_prefix",
469            "size of OperationCache checked_operations_prefix",
470        )
471        .unwrap();
472
473        let operation_cache_ops_know_by_peer = IntGauge::new(
474            "operation_cache_ops_know_by_peer",
475            "size of OperationCache operation_cache_ops_know_by_peer",
476        )
477        .unwrap();
478
479        // consensus state from tick.rs
480        let consensus_state_active_index = IntGauge::new(
481            "consensus_state_active_index",
482            "consensus state active index size",
483        )
484        .unwrap();
485
486        let consensus_state_active_index_without_ops = IntGauge::new(
487            "consensus_state_active_index_without_ops",
488            "consensus state active index without ops size",
489        )
490        .unwrap();
491
492        let consensus_state_incoming_index = IntGauge::new(
493            "consensus_state_incoming_index",
494            "consensus state incoming index size",
495        )
496        .unwrap();
497
498        let consensus_state_discarded_index = IntGauge::new(
499            "consensus_state_discarded_index",
500            "consensus state discarded index size",
501        )
502        .unwrap();
503
504        let consensus_state_block_statuses = IntGauge::new(
505            "consensus_state_block_statuses",
506            "consensus state block statuses size",
507        )
508        .unwrap();
509
510        let endorsement_cache_checked_endorsements = IntGauge::new(
511            "endorsement_cache_checked_endorsements",
512            "endorsement cache checked endorsements size",
513        )
514        .unwrap();
515
516        let endorsement_cache_known_by_peer = IntGauge::new(
517            "endorsement_cache_known_by_peer",
518            "endorsement cache know by peer size",
519        )
520        .unwrap();
521
522        let peernet_total_bytes_received = IntCounter::new(
523            "peernet_total_bytes_received",
524            "total byte received by peernet",
525        )
526        .unwrap();
527
528        let peernet_total_bytes_sent =
529            IntCounter::new("peernet_total_bytes_sent", "total byte sent by peernet").unwrap();
530
531        let operations_final_counter =
532            IntCounter::new("operations_final_counter", "total final operations").unwrap();
533
534        let block_slot_delay = Histogram::with_opts(
535            prometheus::HistogramOpts::new("block_slot_delay", "block slot delay").buckets(vec![
536                0.100, 0.250, 0.500, 1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0,
537            ]),
538        )
539        .unwrap();
540
541        let network_current_version =
542            IntGauge::new("network_current_version", "current version of network").unwrap();
543
544        let db_change_history_size = IntGauge::new(
545            "db_change_history_size",
546            "total size of change history in DB",
547        )
548        .unwrap();
549
550        let db_change_versioning_history_size = IntGauge::new(
551            "db_change_versioning_history_size",
552            "total size of change versioning history in DB",
553        )
554        .unwrap();
555
556        let module_lru_cache_memory_usage = IntGauge::new(
557            "module_lru_cache_memory_usage",
558            "total size of the lru module cache in bytes",
559        )
560        .unwrap();
561
562        let active_history_total_event_len = IntGauge::new(
563            "active_history_total_event_len",
564            "total events len currently in the active history (1 event: max 1kB)",
565        )
566        .unwrap();
567
568        let mut stopper = MetricsStopper::default();
569
570        if enabled {
571            #[cfg(not(feature = "test-exports"))]
572            {
573                register_lazy_static_metrics();
574                let _ = prometheus::register(Box::new(final_cursor_thread.clone()));
575                let _ = prometheus::register(Box::new(final_cursor_period.clone()));
576                let _ = prometheus::register(Box::new(active_cursor_thread.clone()));
577                let _ = prometheus::register(Box::new(active_cursor_period.clone()));
578                let _ = prometheus::register(Box::new(active_out_connections.clone()));
579                let _ = prometheus::register(Box::new(block_cache_blocks_known_by_peer.clone()));
580                let _ = prometheus::register(Box::new(block_cache_checked_headers_size.clone()));
581                let _ = prometheus::register(Box::new(operation_cache_checked_operations.clone()));
582                let _ = prometheus::register(Box::new(active_in_connections.clone()));
583                let _ = prometheus::register(Box::new(operation_cache_ops_know_by_peer.clone()));
584                let _ = prometheus::register(Box::new(consensus_state_active_index.clone()));
585                let _ = prometheus::register(Box::new(
586                    consensus_state_active_index_without_ops.clone(),
587                ));
588                let _ = prometheus::register(Box::new(consensus_state_incoming_index.clone()));
589                let _ = prometheus::register(Box::new(consensus_state_discarded_index.clone()));
590                let _ = prometheus::register(Box::new(consensus_state_block_statuses.clone()));
591                let _ = prometheus::register(Box::new(
592                    operation_cache_checked_operations_prefix.clone(),
593                ));
594                let _ =
595                    prometheus::register(Box::new(endorsement_cache_checked_endorsements.clone()));
596                let _ = prometheus::register(Box::new(endorsement_cache_known_by_peer.clone()));
597                let _ = prometheus::register(Box::new(peernet_total_bytes_received.clone()));
598                let _ = prometheus::register(Box::new(peernet_total_bytes_sent.clone()));
599                let _ = prometheus::register(Box::new(operations_final_counter.clone()));
600                let _ = prometheus::register(Box::new(stakers.clone()));
601                let _ = prometheus::register(Box::new(rolls.clone()));
602                let _ = prometheus::register(Box::new(know_peers.clone()));
603                let _ = prometheus::register(Box::new(banned_peers.clone()));
604                let _ = prometheus::register(Box::new(executed_final_slot.clone()));
605                let _ = prometheus::register(Box::new(executed_final_slot_with_block.clone()));
606                let _ = prometheus::register(Box::new(active_history.clone()));
607                let _ = prometheus::register(Box::new(bootstrap_counter.clone()));
608                let _ = prometheus::register(Box::new(bootstrap_success.clone()));
609                let _ = prometheus::register(Box::new(bootstrap_failed.clone()));
610                let _ = prometheus::register(Box::new(process_available_processors.clone()));
611                let _ = prometheus::register(Box::new(operations_pool.clone()));
612                let _ = prometheus::register(Box::new(endorsements_pool.clone()));
613                let _ = prometheus::register(Box::new(denunciations_pool.clone()));
614                let _ = prometheus::register(Box::new(protocol_tester_success.clone()));
615                let _ = prometheus::register(Box::new(protocol_tester_failed.clone()));
616                let _ = prometheus::register(Box::new(sc_messages_final.clone()));
617                let _ = prometheus::register(Box::new(async_message_pool_size.clone()));
618                let _ = prometheus::register(Box::new(current_time_period.clone()));
619                let _ = prometheus::register(Box::new(current_time_thread.clone()));
620                let _ = prometheus::register(Box::new(block_slot_delay.clone()));
621                let _ = prometheus::register(Box::new(network_current_version.clone()));
622                let _ = prometheus::register(Box::new(db_change_history_size.clone()));
623                let _ = prometheus::register(Box::new(db_change_versioning_history_size.clone()));
624                let _ = prometheus::register(Box::new(module_lru_cache_memory_usage.clone()));
625                let _ = prometheus::register(Box::new(active_history_total_event_len.clone()));
626
627                stopper = server::bind_metrics(addr);
628            }
629        }
630
631        (
632            MassaMetrics {
633                enabled,
634                process_available_processors,
635                consensus_vec,
636                stakers,
637                rolls,
638                current_time_thread,
639                current_time_period,
640                active_history,
641                operations_pool,
642                endorsements_pool,
643                denunciations_pool,
644                async_message_pool_size,
645                sc_messages_final,
646                bootstrap_counter,
647                bootstrap_peers_success: bootstrap_success,
648                bootstrap_peers_failed: bootstrap_failed,
649                protocol_tester_success,
650                protocol_tester_failed,
651                protocol_known_peers: know_peers,
652                protocol_banned_peers: banned_peers,
653                executed_final_slot,
654                executed_final_slot_with_block,
655                peernet_total_bytes_received,
656                peernet_total_bytes_sent,
657                block_slot_delay,
658                active_in_connections,
659                active_out_connections,
660                operations_final_counter,
661                block_cache_checked_headers_size,
662                block_cache_blocks_known_by_peer,
663                operation_cache_checked_operations,
664                operation_cache_checked_operations_prefix,
665                operation_cache_ops_know_by_peer,
666                consensus_state_active_index,
667                consensus_state_active_index_without_ops,
668                consensus_state_incoming_index,
669                consensus_state_discarded_index,
670                consensus_state_block_statuses,
671                endorsement_cache_checked_endorsements,
672                endorsement_cache_known_by_peer,
673                // blocks_counter,
674                // endorsements_counter,
675                // operations_counter,
676                active_cursor_thread,
677                active_cursor_period,
678                final_cursor_thread,
679                final_cursor_period,
680                peers_bandwidth: Arc::new(RwLock::new(HashMap::new())),
681                network_versions_votes: Arc::new(RwLock::new(HashMap::new())),
682                network_current_version,
683                db_change_history_size,
684                db_change_versioning_history_size,
685                module_lru_cache_memory_usage,
686                active_history_total_event_len,
687                tick_delay,
688            },
689            stopper,
690        )
691    }
692
693    pub fn is_enabled(&self) -> bool {
694        self.enabled
695    }
696
697    pub fn get_metrics_for_survey_thread(&self) -> (i64, i64, u64, u64) {
698        (
699            self.active_in_connections.clone().get(),
700            self.active_out_connections.clone().get(),
701            self.peernet_total_bytes_sent.clone().get(),
702            self.peernet_total_bytes_received.clone().get(),
703        )
704    }
705
706    pub fn set_active_connections(&self, in_connections: usize, out_connections: usize) {
707        self.active_in_connections.set(in_connections as i64);
708        self.active_out_connections.set(out_connections as i64);
709    }
710
711    pub fn set_active_cursor(&self, period: u64, thread: u8) {
712        self.active_cursor_thread.set(thread as i64);
713        self.active_cursor_period.set(period as i64);
714    }
715
716    pub fn set_final_cursor(&self, period: u64, thread: u8) {
717        self.final_cursor_thread.set(thread as i64);
718        self.final_cursor_period.set(period as i64);
719    }
720
721    pub fn set_consensus_period(&self, thread: usize, period: u64) {
722        if let Some(g) = self.consensus_vec.get(thread) {
723            g.set(period as f64);
724        }
725    }
726
727    pub fn set_consensus_state(
728        &self,
729        active_index: usize,
730        incoming_index: usize,
731        discarded_index: usize,
732        block_statuses: usize,
733        active_index_without_ops: usize,
734    ) {
735        self.consensus_state_active_index.set(active_index as i64);
736        self.consensus_state_incoming_index
737            .set(incoming_index as i64);
738        self.consensus_state_discarded_index
739            .set(discarded_index as i64);
740        self.consensus_state_block_statuses
741            .set(block_statuses as i64);
742        self.consensus_state_active_index_without_ops
743            .set(active_index_without_ops as i64);
744    }
745
746    pub fn set_block_cache_metrics(&self, checked_header_size: usize, blocks_known_by_peer: usize) {
747        self.block_cache_checked_headers_size
748            .set(checked_header_size as i64);
749        self.block_cache_blocks_known_by_peer
750            .set(blocks_known_by_peer as i64);
751    }
752
753    pub fn set_operations_cache_metrics(
754        &self,
755        checked_operations: usize,
756        checked_operations_prefix: usize,
757        ops_know_by_peer: usize,
758    ) {
759        self.operation_cache_checked_operations
760            .set(checked_operations as i64);
761        self.operation_cache_checked_operations_prefix
762            .set(checked_operations_prefix as i64);
763        self.operation_cache_ops_know_by_peer
764            .set(ops_know_by_peer as i64);
765    }
766
767    pub fn set_endorsements_cache_metrics(
768        &self,
769        checked_endorsements: usize,
770        known_by_peer: usize,
771    ) {
772        self.endorsement_cache_checked_endorsements
773            .set(checked_endorsements as i64);
774        self.endorsement_cache_known_by_peer
775            .set(known_by_peer as i64);
776    }
777
778    pub fn set_peernet_total_bytes_received(&self, new_value: u64) {
779        let diff = new_value.saturating_sub(self.peernet_total_bytes_received.get());
780        self.peernet_total_bytes_received.inc_by(diff);
781    }
782
783    pub fn set_peernet_total_bytes_sent(&self, new_value: u64) {
784        let diff = new_value.saturating_sub(self.peernet_total_bytes_sent.get());
785        self.peernet_total_bytes_sent.inc_by(diff);
786    }
787
788    pub fn inc_operations_final_counter(&self, diff: u64) {
789        self.operations_final_counter.inc_by(diff);
790    }
791
792    pub fn set_known_peers(&self, nb: usize) {
793        self.protocol_known_peers.set(nb as i64);
794    }
795
796    pub fn set_banned_peers(&self, nb: usize) {
797        self.protocol_banned_peers.set(nb as i64);
798    }
799
800    pub fn inc_executed_final_slot(&self) {
801        self.executed_final_slot.inc();
802    }
803
804    pub fn inc_executed_final_slot_with_block(&self) {
805        self.executed_final_slot_with_block.inc();
806    }
807
808    pub fn set_active_history(&self, nb: usize) {
809        self.active_history.set(nb as i64);
810    }
811
812    pub fn inc_bootstrap_counter(&self) {
813        self.bootstrap_counter.inc();
814    }
815
816    pub fn inc_bootstrap_peers_success(&self) {
817        self.bootstrap_peers_success.inc();
818    }
819
820    pub fn inc_bootstrap_peers_failed(&self) {
821        self.bootstrap_peers_failed.inc();
822    }
823
824    pub fn set_operations_pool(&self, nb: usize) {
825        self.operations_pool.set(nb as i64);
826    }
827
828    pub fn set_endorsements_pool(&self, nb: usize) {
829        self.endorsements_pool.set(nb as i64);
830    }
831
832    pub fn set_denunciations_pool(&self, nb: usize) {
833        self.denunciations_pool.set(nb as i64);
834    }
835
836    pub fn inc_protocol_tester_success(&self) {
837        self.protocol_tester_success.inc();
838    }
839
840    pub fn inc_protocol_tester_failed(&self) {
841        self.protocol_tester_failed.inc();
842    }
843
844    pub fn set_stakers(&self, nb: usize) {
845        self.stakers.set(nb as i64);
846    }
847
848    pub fn set_rolls(&self, nb: usize) {
849        self.rolls.set(nb as i64);
850    }
851
852    pub fn inc_sc_messages_final_by(&self, diff: usize) {
853        self.sc_messages_final.inc_by(diff as u64);
854    }
855
856    pub fn set_async_message_pool_size(&self, nb: usize) {
857        self.async_message_pool_size.set(nb as i64);
858    }
859
860    pub fn set_available_processors(&self, nb: usize) {
861        self.process_available_processors.set(nb as i64);
862    }
863
864    pub fn set_current_time_period(&self, period: u64) {
865        self.current_time_period.set(period as i64);
866    }
867
868    pub fn set_current_time_thread(&self, thread: u8) {
869        self.current_time_thread.set(thread as i64);
870    }
871
872    pub fn set_block_slot_delay(&self, delay: f64) {
873        self.block_slot_delay.observe(delay);
874    }
875
876    pub fn set_network_current_version(&self, version: u32) {
877        self.network_current_version.set(version as i64);
878    }
879
880    pub fn set_db_change_history_size(&self, size: usize) {
881        self.db_change_history_size.set(size as i64);
882    }
883
884    pub fn set_db_change_versioning_history_size(&self, size: usize) {
885        self.db_change_versioning_history_size.set(size as i64);
886    }
887
888    pub fn set_module_lru_cache_memory_usage(&self, size: usize) {
889        self.module_lru_cache_memory_usage.set(size as i64);
890    }
891
892    pub fn set_active_history_total_event_len(&self, size: usize) {
893        self.active_history_total_event_len.set(size as i64);
894    }
895
896    // Update the network version vote metrics
897    pub fn update_network_version_vote(&self, data: HashMap<u32, u64>) {
898        if self.enabled {
899            let mut write = self.network_versions_votes.write().unwrap();
900
901            let current_version: u32 = self.network_current_version.get() as u32;
902
903            {
904                let missing_version = write
905                    .keys()
906                    .filter(|key| !data.contains_key(key))
907                    .cloned()
908                    .collect::<Vec<u32>>();
909
910                for key in missing_version {
911                    if let Some(counter) = write.remove(&key) {
912                        if let Err(e) = prometheus::unregister(Box::new(counter)) {
913                            warn!("Failed to unregister network_version_vote_{} : {}", key, e);
914                        }
915                    }
916                }
917
918                if current_version > 0 {
919                    // remove metrics for version 0 if we have a current version > 0
920                    // in this case 0 means no vote
921                    if let Some(counter) = write.remove(&0) {
922                        if let Err(e) = prometheus::unregister(Box::new(counter)) {
923                            warn!("Failed to unregister network_version_vote_0 : {}", e);
924                        }
925                    }
926                }
927            }
928
929            for (version, count) in data.into_iter() {
930                if version.eq(&0) && current_version > 0 {
931                    // skip version 0 if we have a current version
932                    continue;
933                }
934                if let Some(actual_counter) = write.get_mut(&version) {
935                    actual_counter.set(count as i64);
936                } else {
937                    let label = format!("network_version_votes_{}", version);
938                    let counter = IntGauge::new(label, "vote counter for network version").unwrap();
939                    counter.set(count as i64);
940                    let _ = prometheus::register(Box::new(counter.clone()));
941                    write.insert(version, counter);
942                }
943            }
944        }
945    }
946
947    /// Update the bandwidth metrics for all peers
948    /// HashMap<peer_id, (tx, rx)>
949    pub fn update_peers_tx_rx(&self, data: HashMap<String, (u64, u64)>) {
950        if self.enabled {
951            let mut write = self.peers_bandwidth.write().unwrap();
952
953            // metrics of peers that are not in the data HashMap are removed
954            let missing_peer: Vec<String> = write
955                .keys()
956                .filter(|key| !data.contains_key(key.as_str()))
957                .cloned()
958                .collect();
959
960            for key in missing_peer {
961                // remove peer and unregister metrics
962                if let Some((tx, rx)) = write.remove(&key) {
963                    if let Err(e) = prometheus::unregister(Box::new(tx)) {
964                        warn!("Failed to unregister tx metricfor peer {} : {}", key, e);
965                    }
966
967                    if let Err(e) = prometheus::unregister(Box::new(rx)) {
968                        warn!("Failed to unregister rx metric for peer {} : {}", key, e);
969                    }
970                }
971            }
972
973            for (k, (tx_peernet, rx_peernet)) in data {
974                if let Some((tx_metric, rx_metric)) = write.get_mut(&k) {
975                    // peer metrics exist
976                    // update tx and rx
977
978                    let to_add = tx_peernet.saturating_sub(tx_metric.get());
979                    tx_metric.inc_by(to_add);
980
981                    let to_add = rx_peernet.saturating_sub(rx_metric.get());
982                    rx_metric.inc_by(to_add);
983                } else {
984                    // peer metrics does not exist
985                    let label_rx = format!("peer_total_bytes_receive_{}", k);
986                    let label_tx = format!("peer_total_bytes_sent_{}", k);
987
988                    let peer_total_bytes_receive =
989                        IntCounter::new(label_rx, "total byte received by the peer").unwrap();
990
991                    let peer_total_bytes_sent =
992                        IntCounter::new(label_tx, "total byte sent by the peer").unwrap();
993
994                    peer_total_bytes_sent.inc_by(tx_peernet);
995                    peer_total_bytes_receive.inc_by(rx_peernet);
996
997                    let _ = prometheus::register(Box::new(peer_total_bytes_receive.clone()));
998                    let _ = prometheus::register(Box::new(peer_total_bytes_sent.clone()));
999
1000                    write.insert(k, (peer_total_bytes_sent, peer_total_bytes_receive));
1001                }
1002            }
1003        }
1004    }
1005}