Skip to main content

iota_metrics/
metrics_network.rs

1// Copyright (c) Mysten Labs, Inc.
2// Modifications Copyright (c) 2024 IOTA Stiftung
3// SPDX-License-Identifier: Apache-2.0
4
5use 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    /// The connection status of known peers. 0 if not connected, 1 if
18    /// connected.
19    pub network_peer_connected: IntGaugeVec,
20    /// The number of connected peers
21    pub network_peers: IntGauge,
22    /// Number of disconnect events per peer.
23    pub network_peer_disconnects: IntCounterVec,
24    /// Receive buffer size of Anemo socket.
25    pub socket_receive_buffer_size: IntGauge,
26    /// Send buffer size of Anemo socket.
27    pub socket_send_buffer_size: IntGauge,
28
29    /// PathStats
30    /// The rtt for a peer connection in ms.
31    pub network_peer_rtt: IntGaugeVec,
32    /// The total number of lost packets for a peer connection.
33    pub network_peer_lost_packets: IntGaugeVec,
34    /// The total number of lost bytes for a peer connection.
35    pub network_peer_lost_bytes: IntGaugeVec,
36    /// The total number of packets sent for a peer connection.
37    pub network_peer_sent_packets: IntGaugeVec,
38    /// The total number of congestion events for a peer connection.
39    pub network_peer_congestion_events: IntGaugeVec,
40    /// The congestion window for a peer connection.
41    pub network_peer_congestion_window: IntGaugeVec,
42
43    /// FrameStats
44    /// The number of max data frames for a peer connection.
45    pub network_peer_max_data: IntGaugeVec,
46    /// The number of closed connections frames for a peer connection.
47    pub network_peer_closed_connections: IntGaugeVec,
48    /// The number of data blocked frames for a peer connection.
49    pub network_peer_data_blocked: IntGaugeVec,
50
51    /// UDPStats
52    /// The total number datagrams observed by the UDP peer connection.
53    pub network_peer_udp_datagrams: IntGaugeVec,
54    /// The total number bytes observed by the UDP peer connection.
55    pub network_peer_udp_bytes: IntGaugeVec,
56    /// The total number transmits observed by the UDP peer connection.
57    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            // PathStats
98            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            // FrameStats
142            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            // UDPStats
165            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    /// Counter of requests by route
193    requests: IntCounterVec,
194    /// Request latency by route
195    request_latency: HistogramVec,
196    /// Request size by route
197    request_size: HistogramVec,
198    /// Response size by route
199    response_size: HistogramVec,
200    /// Counter of requests exceeding the "excessive" size limit
201    excessive_size_requests: IntCounterVec,
202    /// Counter of responses exceeding the "excessive" size limit
203    excessive_size_responses: IntCounterVec,
204    /// Gauge of the number of inflight requests at any given time by route
205    inflight_requests: IntGaugeVec,
206    /// Failed requests by route
207    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
214// Arbitrarily chosen buckets for message size, with gradually-lowering exponent
215// to give us better resolution at high sizes.
216const SIZE_BYTE_BUCKETS: &[f64] = &[
217    2048., 8192., // *4
218    16384., 32768., 65536., 131072., 262144., 524288., 1048576., // *2
219    1572864., 2359256., 3538944., // *1.5
220    4600627., 5980815., 7775060., 10107578., 13139851., 17081807., 22206349., 28868253., 37528729.,
221    48787348., 63423553., // *1.3
222];
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    /// Size in bytes above which a request or response message is considered
312    /// excessively large
313    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    // The timer is held on to and "observed" once dropped
370    #[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}