1use std::sync::Arc;
6
7use anemo_tower::callback::{MakeCallbackHandler, ResponseHandler};
8use prometheus_filtered::{
9 HistogramTimer, HistogramVec, IntCounterVec, IntGauge, IntGaugeVec, MetricLevel, Registry,
10 register_histogram_vec_with_registry, register_int_counter_vec_with_registry,
11 register_int_gauge_vec_with_registry, register_int_gauge_with_registry,
12};
13use tracing::warn;
14
15#[derive(Clone)]
16pub struct NetworkConnectionMetrics {
17 pub network_peer_connected: IntGaugeVec,
20 pub network_peers: IntGauge,
22 pub network_peer_disconnects: IntCounterVec,
24 pub socket_receive_buffer_size: IntGauge,
26 pub socket_send_buffer_size: IntGauge,
28
29 pub network_peer_rtt: IntGaugeVec,
32 pub network_peer_lost_packets: IntGaugeVec,
34 pub network_peer_lost_bytes: IntGaugeVec,
36 pub network_peer_sent_packets: IntGaugeVec,
38 pub network_peer_congestion_events: IntGaugeVec,
40 pub network_peer_congestion_window: IntGaugeVec,
42
43 pub network_peer_max_data: IntGaugeVec,
46 pub network_peer_closed_connections: IntGaugeVec,
48 pub network_peer_data_blocked: IntGaugeVec,
50
51 pub network_peer_udp_datagrams: IntGaugeVec,
54 pub network_peer_udp_bytes: IntGaugeVec,
56 pub network_peer_udp_transmits: IntGaugeVec,
58}
59
60impl NetworkConnectionMetrics {
61 pub fn new(node: &'static str, registry: &Registry) -> Self {
62 Self {
63 network_peer_connected: register_int_gauge_vec_with_registry!(
64 format!("{node}_network_peer_connected"),
65 "The connection status of a peer. 0 if not connected, 1 if connected",
66 &["peer_id", "type"],
67 registry
68 )
69 .unwrap(),
70 network_peers: register_int_gauge_with_registry!(
71 format!("{node}_network_peers"),
72 "The number of connected peers.",
73 registry;
74 MetricLevel::Warn,
75 )
76 .unwrap(),
77 network_peer_disconnects: register_int_counter_vec_with_registry!(
78 format!("{node}_network_peer_disconnects"),
79 "Number of disconnect events per peer.",
80 &["peer_id", "reason"],
81 registry
82 )
83 .unwrap(),
84 socket_receive_buffer_size: register_int_gauge_with_registry!(
85 format!("{node}_socket_receive_buffer_size"),
86 "Receive buffer size of Anemo socket.",
87 registry
88 )
89 .unwrap(),
90 socket_send_buffer_size: register_int_gauge_with_registry!(
91 format!("{node}_socket_send_buffer_size"),
92 "Send buffer size of Anemo socket.",
93 registry
94 )
95 .unwrap(),
96
97 network_peer_rtt: register_int_gauge_vec_with_registry!(
99 format!("{node}_network_peer_rtt"),
100 "The rtt for a peer connection in ms.",
101 &["peer_id"],
102 registry
103 )
104 .unwrap(),
105 network_peer_lost_packets: register_int_gauge_vec_with_registry!(
106 format!("{node}_network_peer_lost_packets"),
107 "The total number of lost packets for a peer connection.",
108 &["peer_id"],
109 registry
110 )
111 .unwrap(),
112 network_peer_lost_bytes: register_int_gauge_vec_with_registry!(
113 format!("{node}_network_peer_lost_bytes"),
114 "The total number of lost bytes for a peer connection.",
115 &["peer_id"],
116 registry
117 )
118 .unwrap(),
119 network_peer_sent_packets: register_int_gauge_vec_with_registry!(
120 format!("{node}_network_peer_sent_packets"),
121 "The total number of sent packets for a peer connection.",
122 &["peer_id"],
123 registry
124 )
125 .unwrap(),
126 network_peer_congestion_events: register_int_gauge_vec_with_registry!(
127 format!("{node}_network_peer_congestion_events"),
128 "The total number of congestion events for a peer connection.",
129 &["peer_id"],
130 registry
131 )
132 .unwrap(),
133 network_peer_congestion_window: register_int_gauge_vec_with_registry!(
134 format!("{node}_network_peer_congestion_window"),
135 "The congestion window for a peer connection.",
136 &["peer_id"],
137 registry
138 )
139 .unwrap(),
140
141 network_peer_closed_connections: register_int_gauge_vec_with_registry!(
143 format!("{node}_network_peer_closed_connections"),
144 "The number of closed connections for a peer connection.",
145 &["peer_id", "direction"],
146 registry
147 )
148 .unwrap(),
149 network_peer_max_data: register_int_gauge_vec_with_registry!(
150 format!("{node}_network_peer_max_data"),
151 "The number of max data frames for a peer connection.",
152 &["peer_id", "direction"],
153 registry
154 )
155 .unwrap(),
156 network_peer_data_blocked: register_int_gauge_vec_with_registry!(
157 format!("{node}_network_peer_data_blocked"),
158 "The number of data blocked frames for a peer connection.",
159 &["peer_id", "direction"],
160 registry
161 )
162 .unwrap(),
163
164 network_peer_udp_datagrams: register_int_gauge_vec_with_registry!(
166 format!("{node}_network_peer_udp_datagrams"),
167 "The total number datagrams observed by the UDP peer connection.",
168 &["peer_id", "direction"],
169 registry
170 )
171 .unwrap(),
172 network_peer_udp_bytes: register_int_gauge_vec_with_registry!(
173 format!("{node}_network_peer_udp_bytes"),
174 "The total number bytes observed by the UDP peer connection.",
175 &["peer_id", "direction"],
176 registry
177 )
178 .unwrap(),
179 network_peer_udp_transmits: register_int_gauge_vec_with_registry!(
180 format!("{node}_network_peer_udp_transmits"),
181 "The total number transmits observed by the UDP peer connection.",
182 &["peer_id", "direction"],
183 registry
184 )
185 .unwrap(),
186 }
187 }
188}
189
190#[derive(Clone)]
191pub struct NetworkMetrics {
192 requests: IntCounterVec,
194 request_latency: HistogramVec,
196 request_size: HistogramVec,
198 response_size: HistogramVec,
200 excessive_size_requests: IntCounterVec,
202 excessive_size_responses: IntCounterVec,
204 inflight_requests: IntGaugeVec,
206 errors: IntCounterVec,
208}
209
210const LATENCY_SEC_BUCKETS: &[f64] = &[
211 0.001, 0.005, 0.01, 0.05, 0.1, 0.25, 0.5, 1., 2.5, 5., 10., 20., 30., 60., 90.,
212];
213
214const SIZE_BYTE_BUCKETS: &[f64] = &[
217 2048., 8192., 16384., 32768., 65536., 131072., 262144., 524288., 1048576., 1572864., 2359256., 3538944., 4600627., 5980815., 7775060., 10107578., 13139851., 17081807., 22206349., 28868253., 37528729.,
221 48787348., 63423553., ];
223
224impl NetworkMetrics {
225 pub fn new(node: &'static str, direction: &'static str, registry: &Registry) -> Self {
226 let requests = register_int_counter_vec_with_registry!(
227 format!("{node}_{direction}_requests"),
228 "The number of requests made on the network",
229 &["route"],
230 registry;
231 MetricLevel::Warn,
232 )
233 .unwrap();
234
235 let request_latency = register_histogram_vec_with_registry!(
236 format!("{node}_{direction}_request_latency"),
237 "Latency of a request by route",
238 &["route"],
239 LATENCY_SEC_BUCKETS.to_vec(),
240 registry;
241 MetricLevel::Warn,
242 )
243 .unwrap();
244
245 let request_size = register_histogram_vec_with_registry!(
246 format!("{node}_{direction}_request_size"),
247 "Size of a request by route",
248 &["route"],
249 SIZE_BYTE_BUCKETS.to_vec(),
250 registry,
251 )
252 .unwrap();
253
254 let response_size = register_histogram_vec_with_registry!(
255 format!("{node}_{direction}_response_size"),
256 "Size of a response by route",
257 &["route"],
258 SIZE_BYTE_BUCKETS.to_vec(),
259 registry,
260 )
261 .unwrap();
262
263 let excessive_size_requests = register_int_counter_vec_with_registry!(
264 format!("{node}_{direction}_excessive_size_requests"),
265 "The number of excessively large request messages sent",
266 &["route"],
267 registry
268 )
269 .unwrap();
270
271 let excessive_size_responses = register_int_counter_vec_with_registry!(
272 format!("{node}_{direction}_excessive_size_responses"),
273 "The number of excessively large response messages seen",
274 &["route"],
275 registry
276 )
277 .unwrap();
278
279 let inflight_requests = register_int_gauge_vec_with_registry!(
280 format!("{node}_{direction}_inflight_requests"),
281 "The number of inflight network requests",
282 &["route"],
283 registry
284 )
285 .unwrap();
286
287 let errors = register_int_counter_vec_with_registry!(
288 format!("{node}_{direction}_request_errors"),
289 "Number of errors by route",
290 &["route", "status"],
291 registry,
292 )
293 .unwrap();
294
295 Self {
296 requests,
297 request_latency,
298 request_size,
299 response_size,
300 excessive_size_requests,
301 excessive_size_responses,
302 inflight_requests,
303 errors,
304 }
305 }
306}
307
308#[derive(Clone)]
309pub struct MetricsMakeCallbackHandler {
310 metrics: Arc<NetworkMetrics>,
311 excessive_message_size: usize,
314}
315
316impl MetricsMakeCallbackHandler {
317 pub fn new(metrics: Arc<NetworkMetrics>, excessive_message_size: usize) -> Self {
318 Self {
319 metrics,
320 excessive_message_size,
321 }
322 }
323}
324
325impl MakeCallbackHandler for MetricsMakeCallbackHandler {
326 type Handler = MetricsResponseHandler;
327
328 fn make_handler(&self, request: &anemo::Request<bytes::Bytes>) -> Self::Handler {
329 let route = request.route().to_owned();
330
331 self.metrics.requests.with_label_values(&[&route]).inc();
332 self.metrics
333 .inflight_requests
334 .with_label_values(&[&route])
335 .inc();
336 let body_len = request.body().len();
337 self.metrics
338 .request_size
339 .with_label_values(&[&route])
340 .observe(body_len as f64);
341 if body_len > self.excessive_message_size {
342 warn!(
343 "Saw excessively large request with size {body_len} for {route} with peer {:?}",
344 request.peer_id()
345 );
346 self.metrics
347 .excessive_size_requests
348 .with_label_values(&[&route])
349 .inc();
350 }
351
352 let timer = self
353 .metrics
354 .request_latency
355 .with_label_values(&[&route])
356 .start_timer();
357
358 MetricsResponseHandler {
359 metrics: self.metrics.clone(),
360 timer,
361 route,
362 excessive_message_size: self.excessive_message_size,
363 }
364 }
365}
366
367pub struct MetricsResponseHandler {
368 metrics: Arc<NetworkMetrics>,
369 #[expect(unused)]
371 timer: HistogramTimer,
372 route: String,
373 excessive_message_size: usize,
374}
375
376impl ResponseHandler for MetricsResponseHandler {
377 fn on_response(self, response: &anemo::Response<bytes::Bytes>) {
378 let body_len = response.body().len();
379 self.metrics
380 .response_size
381 .with_label_values(&[&self.route])
382 .observe(body_len as f64);
383 if body_len > self.excessive_message_size {
384 warn!(
385 "Saw excessively large response with size {body_len} for {} with peer {:?}",
386 self.route,
387 response.peer_id()
388 );
389 self.metrics
390 .excessive_size_responses
391 .with_label_values(&[&self.route])
392 .inc();
393 }
394
395 if !response.status().is_success() {
396 let status = response.status().to_u16().to_string();
397 self.metrics
398 .errors
399 .with_label_values(&[&self.route, &status])
400 .inc();
401 }
402 }
403
404 fn on_error<E>(self, _error: &E) {
405 self.metrics
406 .errors
407 .with_label_values(&[self.route.as_str(), "unknown"])
408 .inc();
409 }
410}
411
412impl Drop for MetricsResponseHandler {
413 fn drop(&mut self) {
414 self.metrics
415 .inflight_requests
416 .with_label_values(&[&self.route])
417 .dec();
418 }
419}