metrics.rs (7830B)
1 use std::sync::{ 2 Mutex as StdMutex, 3 atomic::{AtomicU64, Ordering}, 4 }; 5 6 use serde::Serialize; 7 8 #[derive(Default)] 9 pub(super) struct P2pMetricsCounters { 10 pub(super) inbound_sessions_started: AtomicU64, 11 pub(super) inbound_sessions_rejected: AtomicU64, 12 pub(super) outbound_connect_attempts: AtomicU64, 13 pub(super) outbound_connect_successes: AtomicU64, 14 pub(super) outbound_connect_failures: AtomicU64, 15 pub(super) outbound_sessions_started: AtomicU64, 16 pub(super) sessions_closed: AtomicU64, 17 pub(super) session_failures: AtomicU64, 18 pub(super) quiet_disconnects: AtomicU64, 19 pub(super) envelopes_received: AtomicU64, 20 pub(super) hello_envelopes_received: AtomicU64, 21 pub(super) peer_status_envelopes_received: AtomicU64, 22 pub(super) inventory_envelopes_received: AtomicU64, 23 pub(super) data_envelopes_received: AtomicU64, 24 pub(super) transaction_envelopes_received: AtomicU64, 25 pub(super) transactions_received: AtomicU64, 26 pub(super) burn_bundle_envelopes_received: AtomicU64, 27 pub(super) burn_bundles_received: AtomicU64, 28 pub(super) control_envelopes_received: AtomicU64, 29 pub(super) rejected_blocks: AtomicU64, 30 pub(super) rejected_block_batches: AtomicU64, 31 pub(super) rejected_snapshots: AtomicU64, 32 pub(super) bytes_received: AtomicU64, 33 pub(super) parse_errors: AtomicU64, 34 pub(super) empty_frames: AtomicU64, 35 pub(super) self_peer_rejections: AtomicU64, 36 pub(super) self_peer_skips: AtomicU64, 37 pub(super) outbound_queue_full: AtomicU64, 38 pub(super) outbound_queue_closed: AtomicU64, 39 pub(super) last_session_failure: StdMutex<Option<String>>, 40 pub(super) last_empty_frame_remote: StdMutex<Option<String>>, 41 pub(super) last_parse_error: StdMutex<Option<String>>, 42 pub(super) last_chain_payload_error: StdMutex<Option<String>>, 43 } 44 45 #[derive(Clone, Debug, Default, Eq, PartialEq, Serialize)] 46 pub struct P2pMetrics { 47 pub inbound_sessions_started: u64, 48 pub inbound_sessions_rejected: u64, 49 pub outbound_connect_attempts: u64, 50 pub outbound_connect_successes: u64, 51 pub outbound_connect_failures: u64, 52 pub outbound_sessions_started: u64, 53 pub sessions_closed: u64, 54 pub session_failures: u64, 55 pub quiet_disconnects: u64, 56 pub envelopes_received: u64, 57 pub hello_envelopes_received: u64, 58 pub peer_status_envelopes_received: u64, 59 pub inventory_envelopes_received: u64, 60 pub data_envelopes_received: u64, 61 pub transaction_envelopes_received: u64, 62 pub transactions_received: u64, 63 pub burn_bundle_envelopes_received: u64, 64 pub burn_bundles_received: u64, 65 pub control_envelopes_received: u64, 66 pub rejected_blocks: u64, 67 pub rejected_block_batches: u64, 68 pub rejected_snapshots: u64, 69 pub bytes_received: u64, 70 pub parse_errors: u64, 71 pub empty_frames: u64, 72 pub self_peer_rejections: u64, 73 pub self_peer_skips: u64, 74 pub outbound_queue_full: u64, 75 pub outbound_queue_closed: u64, 76 pub last_session_failure: Option<String>, 77 pub last_empty_frame_remote: Option<String>, 78 pub last_parse_error: Option<String>, 79 pub last_chain_payload_error: Option<String>, 80 } 81 82 impl P2pMetricsCounters { 83 pub(super) fn inc(counter: &AtomicU64) { 84 counter.fetch_add(1, Ordering::Relaxed); 85 } 86 87 pub(super) fn add(counter: &AtomicU64, amount: u64) { 88 counter.fetch_add(amount, Ordering::Relaxed); 89 } 90 91 pub(super) fn set_last(target: &StdMutex<Option<String>>, value: impl Into<String>) { 92 if let Ok(mut last) = target.lock() { 93 *last = Some(value.into()); 94 } 95 } 96 97 pub(super) fn snapshot(&self) -> P2pMetrics { 98 P2pMetrics { 99 inbound_sessions_started: self.inbound_sessions_started.load(Ordering::Relaxed), 100 inbound_sessions_rejected: self.inbound_sessions_rejected.load(Ordering::Relaxed), 101 outbound_connect_attempts: self.outbound_connect_attempts.load(Ordering::Relaxed), 102 outbound_connect_successes: self.outbound_connect_successes.load(Ordering::Relaxed), 103 outbound_connect_failures: self.outbound_connect_failures.load(Ordering::Relaxed), 104 outbound_sessions_started: self.outbound_sessions_started.load(Ordering::Relaxed), 105 sessions_closed: self.sessions_closed.load(Ordering::Relaxed), 106 session_failures: self.session_failures.load(Ordering::Relaxed), 107 quiet_disconnects: self.quiet_disconnects.load(Ordering::Relaxed), 108 envelopes_received: self.envelopes_received.load(Ordering::Relaxed), 109 hello_envelopes_received: self.hello_envelopes_received.load(Ordering::Relaxed), 110 peer_status_envelopes_received: self 111 .peer_status_envelopes_received 112 .load(Ordering::Relaxed), 113 inventory_envelopes_received: self.inventory_envelopes_received.load(Ordering::Relaxed), 114 data_envelopes_received: self.data_envelopes_received.load(Ordering::Relaxed), 115 transaction_envelopes_received: self 116 .transaction_envelopes_received 117 .load(Ordering::Relaxed), 118 transactions_received: self.transactions_received.load(Ordering::Relaxed), 119 burn_bundle_envelopes_received: self 120 .burn_bundle_envelopes_received 121 .load(Ordering::Relaxed), 122 burn_bundles_received: self.burn_bundles_received.load(Ordering::Relaxed), 123 control_envelopes_received: self.control_envelopes_received.load(Ordering::Relaxed), 124 rejected_blocks: self.rejected_blocks.load(Ordering::Relaxed), 125 rejected_block_batches: self.rejected_block_batches.load(Ordering::Relaxed), 126 rejected_snapshots: self.rejected_snapshots.load(Ordering::Relaxed), 127 bytes_received: self.bytes_received.load(Ordering::Relaxed), 128 parse_errors: self.parse_errors.load(Ordering::Relaxed), 129 empty_frames: self.empty_frames.load(Ordering::Relaxed), 130 self_peer_rejections: self.self_peer_rejections.load(Ordering::Relaxed), 131 self_peer_skips: self.self_peer_skips.load(Ordering::Relaxed), 132 outbound_queue_full: self.outbound_queue_full.load(Ordering::Relaxed), 133 outbound_queue_closed: self.outbound_queue_closed.load(Ordering::Relaxed), 134 last_session_failure: self 135 .last_session_failure 136 .lock() 137 .ok() 138 .and_then(|last| last.clone()), 139 last_empty_frame_remote: self 140 .last_empty_frame_remote 141 .lock() 142 .ok() 143 .and_then(|last| last.clone()), 144 last_parse_error: self 145 .last_parse_error 146 .lock() 147 .ok() 148 .and_then(|last| last.clone()), 149 last_chain_payload_error: self 150 .last_chain_payload_error 151 .lock() 152 .ok() 153 .and_then(|last| last.clone()), 154 } 155 } 156 } 157 158 #[cfg(test)] 159 mod tests { 160 use super::P2pMetricsCounters; 161 162 #[test] 163 fn snapshot_reports_rejected_chain_payloads() { 164 let metrics = P2pMetricsCounters::default(); 165 166 P2pMetricsCounters::inc(&metrics.rejected_blocks); 167 P2pMetricsCounters::inc(&metrics.rejected_block_batches); 168 P2pMetricsCounters::inc(&metrics.rejected_snapshots); 169 P2pMetricsCounters::set_last(&metrics.last_chain_payload_error, "block: invalid VDF"); 170 171 let snapshot = metrics.snapshot(); 172 assert_eq!(snapshot.rejected_blocks, 1); 173 assert_eq!(snapshot.rejected_block_batches, 1); 174 assert_eq!(snapshot.rejected_snapshots, 1); 175 assert_eq!( 176 snapshot.last_chain_payload_error.as_deref(), 177 Some("block: invalid VDF") 178 ); 179 } 180 }