1use 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 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#[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 enabled: bool,
191
192 process_available_processors: IntGauge,
194
195 consensus_vec: Vec<Gauge>,
198
199 stakers: IntGauge,
201 rolls: IntGauge,
203
204 current_time_thread: IntGauge,
206 current_time_period: IntGauge,
208
209 active_history: IntGauge,
211
212 operations_pool: IntGauge,
214 endorsements_pool: IntGauge,
216 denunciations_pool: IntGauge,
218
219 async_message_pool_size: IntGauge,
221
222 sc_messages_final: IntCounter,
224
225 bootstrap_counter: IntCounter,
227 bootstrap_peers_success: IntCounter,
229 bootstrap_peers_failed: IntCounter,
231
232 protocol_tester_success: IntCounter,
234 protocol_tester_failed: IntCounter,
236
237 protocol_known_peers: IntGauge,
239 protocol_banned_peers: IntGauge,
241
242 executed_final_slot: IntCounter,
244 executed_final_slot_with_block: IntCounter,
246
247 peernet_total_bytes_received: IntCounter,
249 peernet_total_bytes_sent: IntCounter,
251
252 block_slot_delay: Histogram,
254
255 active_in_connections: IntGauge,
257 active_out_connections: IntGauge,
259
260 operations_final_counter: IntCounter,
262
263 block_cache_checked_headers_size: IntGauge,
265 block_cache_blocks_known_by_peer: IntGauge,
266
267 operation_cache_checked_operations: IntGauge,
269 operation_cache_checked_operations_prefix: IntGauge,
270 operation_cache_ops_know_by_peer: IntGauge,
271
272 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_checked_endorsements: IntGauge,
281 endorsement_cache_known_by_peer: IntGauge,
282
283 active_cursor_thread: IntGauge,
285 active_cursor_period: IntGauge,
286
287 final_cursor_thread: IntGauge,
288 final_cursor_period: IntGauge,
289
290 peers_bandwidth: Arc<RwLock<HashMap<String, (IntCounter, IntCounter)>>>,
292
293 network_versions_votes: Arc<RwLock<HashMap<u32, IntGauge>>>,
295 network_current_version: IntGauge,
296
297 db_change_history_size: IntGauge,
299 db_change_versioning_history_size: IntGauge,
300
301 module_lru_cache_memory_usage: IntGauge,
303
304 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 let process_available_processors =
336 IntGauge::new("process_available_processors", "number of processors")
337 .expect("Failed to create available_processors counter");
338
339 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 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 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 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 let active_in_connections =
441 IntGauge::new("active_in_connections", "active connections IN len").unwrap();
442
443 let active_out_connections =
445 IntGauge::new("active_out_connections", "active connections OUT len").unwrap();
446
447 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 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 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 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 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 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 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 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 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 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 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 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}