massa_grpc/stream/
tx_throughput.rs1use 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
11const DEFAULT_THROUGHPUT_INTERVAL: u64 = 10;
17
18fn resolve_interval_secs(requested: Option<u64>) -> u64 {
24 match requested {
25 Some(secs) if secs > 0 => secs,
26 _ => DEFAULT_THROUGHPUT_INTERVAL,
27 }
28}
29
30pub 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
40pub 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
50pub(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 let (tx, rx) = tokio::sync::mpsc::channel(grpc.grpc_config.max_channel_size);
59 let mut in_stream = request.into_inner();
61
62 tokio::spawn(async move {
64 let mut interval = time::interval(Duration::from_secs(DEFAULT_THROUGHPUT_INTERVAL));
65
66 loop {
68 select! {
69 res = in_stream.next() => {
71 match res {
72 Some(Ok(req)) => {
73 let new_timer = resolve_interval_secs(req.interval);
76 interval = time::interval(Duration::from_secs(new_timer));
77 interval.reset();
78 },
79 _ => {
80 break;
82 }
83 }
84 },
85 _ = interval.tick() => {
87 let stats = execution_controller.get_stats();
88 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 if let Err(e) = tx
100 .send(Ok(grpc_api::TransactionsThroughputResponse {
101 throughput,
102 }))
103 .await
104 {
105 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
118pub(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 let (tx, rx) = tokio::sync::mpsc::channel(grpc.grpc_config.max_channel_size);
127 let request = request.into_inner();
129
130 tokio::spawn(async move {
132 let mut interval =
133 time::interval(Duration::from_secs(resolve_interval_secs(request.interval)));
134
135 loop {
137 interval.tick().await;
139
140 let stats = execution_controller.get_stats();
141 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 if let Err(e) = tx
153 .send(Ok(grpc_api::TransactionsThroughputServerResponse {
154 throughput,
155 }))
156 .await
157 {
158 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 assert_eq!(resolve_interval_secs(None), DEFAULT_THROUGHPUT_INTERVAL);
180 assert_eq!(resolve_interval_secs(Some(0)), DEFAULT_THROUGHPUT_INTERVAL);
181 assert_eq!(resolve_interval_secs(Some(1)), 1);
183 assert_eq!(resolve_interval_secs(Some(42)), 42);
184 }
185}