Skip to main content

iota_indexer/processors/
network_metrics_processor.rs

1// Copyright (c) Mysten Labs, Inc.
2// Modifications Copyright (c) 2024 IOTA Stiftung
3// SPDX-License-Identifier: Apache-2.0
4
5use tap::tap::TapFallible;
6use tracing::{error, info};
7
8use crate::{
9    errors::IndexerError,
10    metrics::IndexerMetrics,
11    store::{IndexerAnalyticalStore, diesel_macro::spawn_blocking_task},
12    types::IndexerResult,
13};
14
15const MIN_NETWORK_METRICS_PROCESSOR_BATCH_SIZE: usize = 10;
16const MAX_NETWORK_METRICS_PROCESSOR_BATCH_SIZE: usize = 80000;
17const NETWORK_METRICS_PROCESSOR_PARALLELISM: usize = 1;
18
19pub struct NetworkMetricsProcessor<S> {
20    pub store: S,
21    metrics: IndexerMetrics,
22    pub min_network_metrics_processor_batch_size: usize,
23    pub max_network_metrics_processor_batch_size: usize,
24    pub network_metrics_processor_parallelism: usize,
25}
26
27impl<S> NetworkMetricsProcessor<S>
28where
29    S: IndexerAnalyticalStore + Clone + Sync + Send + 'static,
30{
31    pub fn new(store: S, metrics: IndexerMetrics) -> NetworkMetricsProcessor<S> {
32        let min_network_metrics_processor_batch_size =
33            std::env::var("MIN_NETWORK_METRICS_PROCESSOR_BATCH_SIZE")
34                .map(|s| {
35                    s.parse::<usize>()
36                        .unwrap_or(MIN_NETWORK_METRICS_PROCESSOR_BATCH_SIZE)
37                })
38                .unwrap_or(MIN_NETWORK_METRICS_PROCESSOR_BATCH_SIZE);
39        let max_network_metrics_processor_batch_size =
40            std::env::var("MAX_NETWORK_METRICS_PROCESSOR_BATCH_SIZE")
41                .map(|s| {
42                    s.parse::<usize>()
43                        .unwrap_or(MAX_NETWORK_METRICS_PROCESSOR_BATCH_SIZE)
44                })
45                .unwrap_or(MAX_NETWORK_METRICS_PROCESSOR_BATCH_SIZE);
46        let network_metrics_processor_parallelism =
47            std::env::var("NETWORK_METRICS_PROCESSOR_PARALLELISM")
48                .map(|s| {
49                    s.parse::<usize>()
50                        .unwrap_or(NETWORK_METRICS_PROCESSOR_PARALLELISM)
51                })
52                .unwrap_or(NETWORK_METRICS_PROCESSOR_PARALLELISM);
53        Self {
54            store,
55            metrics,
56            min_network_metrics_processor_batch_size,
57            max_network_metrics_processor_batch_size,
58            network_metrics_processor_parallelism,
59        }
60    }
61
62    pub async fn start(&self) -> IndexerResult<()> {
63        info!("Indexer network metrics async processor started...");
64        let latest_tx_count_metrics = self
65            .store
66            .get_latest_tx_count_metrics()
67            .await
68            .unwrap_or_default();
69        let latest_epoch_peak_tps = self
70            .store
71            .get_latest_epoch_peak_tps()
72            .await
73            .unwrap_or_default();
74        let mut last_processed_cp_seq = latest_tx_count_metrics
75            .unwrap_or_default()
76            .checkpoint_sequence_number;
77        let mut last_processed_peak_tps_epoch = latest_epoch_peak_tps.unwrap_or_default().epoch;
78
79        loop {
80            let latest_stored_checkpoint = loop {
81                if let Some(latest_stored_checkpoint) =
82                    self.store.get_latest_stored_checkpoint().await?
83                {
84                    if latest_stored_checkpoint.sequence_number
85                        >= last_processed_cp_seq
86                            + self.min_network_metrics_processor_batch_size as i64
87                    {
88                        break latest_stored_checkpoint;
89                    }
90                }
91                tokio::time::sleep(tokio::time::Duration::from_secs(1)).await;
92            };
93
94            let available_checkpoints =
95                latest_stored_checkpoint.sequence_number - last_processed_cp_seq;
96            let batch_size =
97                available_checkpoints.min(self.max_network_metrics_processor_batch_size as i64);
98
99            info!(
100                "Preparing tx count metrics for checkpoints [{}-{}]",
101                last_processed_cp_seq + 1,
102                last_processed_cp_seq + batch_size
103            );
104
105            let step_size =
106                (batch_size as usize / self.network_metrics_processor_parallelism).max(1);
107            let mut persist_tasks = vec![];
108
109            for chunk_start_cp in
110                (last_processed_cp_seq + 1..=last_processed_cp_seq + batch_size).step_by(step_size)
111            {
112                let chunk_end_cp =
113                    (chunk_start_cp + step_size as i64).min(last_processed_cp_seq + batch_size + 1);
114
115                let store = self.store.clone();
116                persist_tasks.push(spawn_blocking_task(move || {
117                    store.persist_tx_count_metrics(chunk_start_cp, chunk_end_cp)
118                }));
119            }
120
121            futures::future::join_all(persist_tasks)
122                .await
123                .into_iter()
124                .collect::<Result<Vec<_>, _>>()
125                .tap_err(|e| error!("error joining network persist tasks: {e:?}"))?
126                .into_iter()
127                .collect::<Result<Vec<_>, _>>()
128                .tap_err(|e| error!("error persisting tx count metrics: {e:?}"))?;
129
130            last_processed_cp_seq += batch_size;
131
132            self.metrics
133                .latest_network_metrics_cp_seq
134                .set(last_processed_cp_seq);
135
136            let end_cp = self
137                .store
138                .get_checkpoints_in_range(last_processed_cp_seq, last_processed_cp_seq + 1)
139                .await?
140                .first()
141                .ok_or(IndexerError::PostgresRead)
142                .inspect_err(|_| {
143                    tracing::error!("cannot read checkpoint from PG for epoch peak TPS")
144                })?
145                .clone();
146            for epoch in last_processed_peak_tps_epoch + 1..end_cp.epoch {
147                self.store.persist_epoch_peak_tps(epoch).await?;
148                last_processed_peak_tps_epoch = epoch;
149                info!("Persisted epoch peak TPS for epoch {}", epoch);
150            }
151        }
152    }
153}