iota_indexer/processors/
network_metrics_processor.rs1use 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}