massa_grpc/stream/
tx_throughput.rs

1// Copyright (c) 2023 MASSA LABS <info@massa.net>
2
3use crate::{error::GrpcError, server::MassaPublicGrpc};
4use futures_util::StreamExt;
5use massa_proto_rs::massa::api::v1 as grpc_api;
6use std::pin::Pin;
7use std::time::Duration;
8use tokio::{select, time};
9use tracing::error;
10
11/// default throughput interval in seconds
12///
13/// set 'high' value to avoid spamming the client with updates who doesn't need
14///
15/// end user can override this value by sending a request with a custom interval
16const DEFAULT_THROUGHPUT_INTERVAL: u64 = 10;
17
18/// Resolve the throughput interval (in seconds) from an optional client-supplied
19/// value, falling back to the default when it is absent or zero.
20///
21/// A zero interval must never reach `tokio::time::interval`, which panics when the
22/// period is zero and would take down the streaming task.
23fn resolve_interval_secs(requested: Option<u64>) -> u64 {
24    match requested {
25        Some(secs) if secs > 0 => secs,
26        _ => DEFAULT_THROUGHPUT_INTERVAL,
27    }
28}
29
30/// Type declaration for TransactionsThroughput
31pub type TransactionsThroughputStreamType = Pin<
32    Box<
33        dyn futures_util::Stream<
34                Item = Result<grpc_api::TransactionsThroughputResponse, tonic::Status>,
35            > + Send
36            + 'static,
37    >,
38>;
39
40/// Type declaration for TransactionsThroughput
41pub type TransactionsThroughputServerStreamType = Pin<
42    Box<
43        dyn futures_util::Stream<
44                Item = Result<grpc_api::TransactionsThroughputServerResponse, tonic::Status>,
45            > + Send
46            + 'static,
47    >,
48>;
49
50/// The function returns a stream of transaction throughput statistics
51pub(crate) async fn transactions_throughput(
52    grpc: &MassaPublicGrpc,
53    request: tonic::Request<tonic::Streaming<grpc_api::TransactionsThroughputRequest>>,
54) -> Result<TransactionsThroughputStreamType, GrpcError> {
55    let execution_controller = grpc.execution_controller.clone();
56
57    // Create a channel for sending responses to the client
58    let (tx, rx) = tokio::sync::mpsc::channel(grpc.grpc_config.max_channel_size);
59    // Extract the incoming stream of operations messages
60    let mut in_stream = request.into_inner();
61
62    // Spawn a new Tokio task to handle the stream processing
63    tokio::spawn(async move {
64        let mut interval = time::interval(Duration::from_secs(DEFAULT_THROUGHPUT_INTERVAL));
65
66        // Continuously loop until the stream ends or an error occurs
67        loop {
68            select! {
69                // Receive a new message from the in_stream
70                res = in_stream.next() => {
71                    match res {
72                        Some(Ok(req)) => {
73                            // Update the interval timer based on the request (or use the default).
74                            // A zero interval falls back to the default to avoid panicking `time::interval`.
75                            let new_timer = resolve_interval_secs(req.interval);
76                            interval = time::interval(Duration::from_secs(new_timer));
77                            interval.reset();
78                        },
79                        _ => {
80                            // Client disconnected
81                            break;
82                        }
83                    }
84                },
85                // Execute the code block whenever the timer ticks
86                _ = interval.tick() => {
87                    let stats = execution_controller.get_stats();
88                    // Calculate the throughput over the time window
89                    let nb_sec_range = stats
90                        .time_window_end
91                        .saturating_sub(stats.time_window_start)
92                        .to_duration()
93                        .as_secs();
94                    let throughput = stats
95                        .final_executed_operations_count
96                        .checked_div(nb_sec_range as usize)
97                        .unwrap_or_default() as u32;
98                    // Send the throughput response back to the client
99                    if let Err(e) = tx
100                        .send(Ok(grpc_api::TransactionsThroughputResponse {
101                            throughput,
102                        }))
103                        .await
104                    {
105                        // Log an error if sending the response fails
106                        error!("failed to send back transactions_throughput response: {}", e);
107                        break;
108                    }
109                }
110            }
111        }
112    });
113
114    let out_stream = tokio_stream::wrappers::ReceiverStream::new(rx);
115    Ok(Box::pin(out_stream) as TransactionsThroughputStreamType)
116}
117
118/// The function returns a stream unidirectional of transaction throughput statistics
119pub(crate) async fn transactions_throughput_server(
120    grpc: &MassaPublicGrpc,
121    request: tonic::Request<grpc_api::TransactionsThroughputServerRequest>,
122) -> Result<TransactionsThroughputServerStreamType, GrpcError> {
123    let execution_controller = grpc.execution_controller.clone();
124
125    // Create a channel for sending responses to the client
126    let (tx, rx) = tokio::sync::mpsc::channel(grpc.grpc_config.max_channel_size);
127    // Extract the incoming stream of operations messages
128    let request = request.into_inner();
129
130    // Spawn a new Tokio task to handle the stream processing
131    tokio::spawn(async move {
132        let mut interval =
133            time::interval(Duration::from_secs(resolve_interval_secs(request.interval)));
134
135        // Continuously loop until the stream ends or an error occurs
136        loop {
137            // Execute the code block whenever the timer ticks
138            interval.tick().await;
139
140            let stats = execution_controller.get_stats();
141            // Calculate the throughput over the time window
142            let nb_sec_range = stats
143                .time_window_end
144                .saturating_sub(stats.time_window_start)
145                .to_duration()
146                .as_secs();
147            let throughput = stats
148                .final_executed_operations_count
149                .checked_div(nb_sec_range as usize)
150                .unwrap_or_default() as u32;
151            // Send the throughput response back to the client
152            if let Err(e) = tx
153                .send(Ok(grpc_api::TransactionsThroughputServerResponse {
154                    throughput,
155                }))
156                .await
157            {
158                // Log an error if sending the response fails
159                error!(
160                    "failed to send back transactions_throughput response: {}",
161                    e
162                );
163                break;
164            }
165        }
166    });
167
168    let out_stream = tokio_stream::wrappers::ReceiverStream::new(rx);
169    Ok(Box::pin(out_stream) as TransactionsThroughputServerStreamType)
170}
171
172#[cfg(test)]
173mod tests {
174    use super::{resolve_interval_secs, DEFAULT_THROUGHPUT_INTERVAL};
175
176    #[test]
177    fn test_resolve_interval_secs() {
178        // Absent or zero -> default (zero would panic tokio's time::interval).
179        assert_eq!(resolve_interval_secs(None), DEFAULT_THROUGHPUT_INTERVAL);
180        assert_eq!(resolve_interval_secs(Some(0)), DEFAULT_THROUGHPUT_INTERVAL);
181        // Positive values are used as-is.
182        assert_eq!(resolve_interval_secs(Some(1)), 1);
183        assert_eq!(resolve_interval_secs(Some(42)), 42);
184    }
185}