1use ahash::AHashMap;
2use futures::future;
3use num_format::{Locale, ToFormattedString};
4use sled::Tree;
5use std::{collections::HashMap, future::Future, sync::Mutex};
6use subxt::{OnlineClient, backend::legacy::LegacyRpcMethods, blocks::Block, metadata::Metadata};
7use tokio::{
8 sync::{RwLock, mpsc, watch},
9 time::{self, Duration, Instant, MissedTickBehavior},
10};
11use tracing::{debug, error, info};
12use zerocopy::{AsBytes, BigEndian, FromBytes, byteorder::U32};
13
14use crate::{shared::*, websockets::process_msg_status};
15
16#[allow(clippy::type_complexity)]
17pub struct Indexer<R: RuntimeIndexer + ?Sized> {
18 trees: Trees<<R::ChainKey as IndexKey>::ChainTrees>,
19 api: Option<OnlineClient<R::RuntimeConfig>>,
20 rpc: Option<LegacyRpcMethods<R::RuntimeConfig>>,
21 index_variant: bool,
22 store_events: bool,
23 metadata_map_lock: RwLock<AHashMap<u32, Metadata>>,
24 status_sub: Mutex<Vec<mpsc::UnboundedSender<ResponseMessage<R::ChainKey>>>>,
25 events_sub_map:
26 Mutex<HashMap<Key<R::ChainKey>, Vec<mpsc::UnboundedSender<ResponseMessage<R::ChainKey>>>>>,
27}
28
29impl<R: RuntimeIndexer> Indexer<R> {
30 fn new(
31 trees: Trees<<R::ChainKey as IndexKey>::ChainTrees>,
32 api: OnlineClient<R::RuntimeConfig>,
33 rpc: LegacyRpcMethods<R::RuntimeConfig>,
34 index_variant: bool,
35 store_events: bool,
36 ) -> Self {
37 Indexer {
38 trees,
39 api: Some(api),
40 rpc: Some(rpc),
41 index_variant,
42 store_events,
43 metadata_map_lock: RwLock::new(AHashMap::new()),
44 status_sub: Vec::new().into(),
45 events_sub_map: HashMap::new().into(),
46 }
47 }
48
49 pub fn new_test(trees: Trees<<R::ChainKey as IndexKey>::ChainTrees>) -> Self {
50 Indexer {
51 trees,
52 api: None,
53 rpc: None,
54 index_variant: true,
55 store_events: true,
56 metadata_map_lock: RwLock::new(AHashMap::new()),
57 status_sub: Vec::new().into(),
58 events_sub_map: HashMap::new().into(),
59 }
60 }
61
62 async fn index_head(
63 &self,
64 next: impl Future<
65 Output = Option<
66 Result<Block<R::RuntimeConfig, OnlineClient<R::RuntimeConfig>>, subxt::Error>,
67 >,
68 >,
69 ) -> Result<(u32, u32, u32), IndexError> {
70 let block = next.await.unwrap()?;
71 self.index_block(block.number().into().try_into().unwrap())
72 .await
73 }
74
75 async fn index_block(&self, block_number: u32) -> Result<(u32, u32, u32), IndexError> {
76 let mut key_count = 0;
77 let api = self.api.as_ref().unwrap();
78 let rpc = self.rpc.as_ref().unwrap();
79
80 let block_hash = match rpc.chain_get_block_hash(Some(block_number.into())).await? {
81 Some(block_hash) => block_hash,
82 None => return Err(IndexError::BlockNotFound(block_number)),
83 };
84 let runtime_version = rpc.state_get_runtime_version(Some(block_hash)).await?;
86
87 let metadata_map = self.metadata_map_lock.read().await;
88 let metadata = match metadata_map.get(&runtime_version.spec_version) {
89 Some(metadata) => {
90 let metadata = metadata.clone();
91 drop(metadata_map);
92 metadata
93 }
94 None => {
95 drop(metadata_map);
96 let mut metadata_map = self.metadata_map_lock.write().await;
97
98 match metadata_map.get(&runtime_version.spec_version) {
99 Some(metadata) => metadata.clone(),
100 None => {
101 info!(
102 "Downloading metadata for spec version {}",
103 runtime_version.spec_version
104 );
105 let metadata: Metadata = rpc
106 .state_get_metadata(Some(block_hash))
107 .await?
108 .to_frame_metadata()?
109 .try_into()?;
110 info!(
111 "Finished downloading metadata for spec version {}",
112 runtime_version.spec_version
113 );
114 metadata_map.insert(runtime_version.spec_version, metadata.clone());
115 metadata
116 }
117 }
118 }
119 };
120
121 let events =
122 subxt::events::new_events_from_client(metadata, block_hash, api.clone()).await?;
123
124 for (i, event) in events.iter().enumerate() {
125 match event {
126 Ok(event) => {
127 let event_index = i.try_into().unwrap();
128 if self.index_variant {
129 self.index_event(
130 Key::Variant(event.pallet_index(), event.variant_index()),
131 block_number,
132 event_index,
133 )?;
134 key_count += 1;
135 }
136 if let Ok(event_key_count) =
137 R::process_event(self, block_number, event_index, event)
138 {
139 key_count += event_key_count;
140 }
141 }
142 Err(error) => error!("Block: {}, error: {}", block_number, error),
143 }
144 }
145
146 if self.store_events {
147 let key: U32<BigEndian> = block_number.into();
148 let spec_version: U32<BigEndian> = runtime_version.spec_version.into();
149
150 self.trees.block_events.insert(
151 key.as_bytes(),
152 [spec_version.as_bytes(), events.bytes()].concat(),
153 )?;
154 }
155
156 Ok((block_number, events.len(), key_count))
157 }
158
159 pub fn notify_status_subscribers(&self) {
160 let msg = process_msg_status::<R>(&self.trees.span);
161 let txs = self.status_sub.lock().unwrap();
162 for tx in txs.iter() {
163 if tx.send(msg.clone()).is_ok() {}
164 }
165 }
166
167 pub fn notify_subscribers(&self, search_key: Key<R::ChainKey>, event: Event) {
168 let events_sub_map = self.events_sub_map.lock().unwrap();
169 if let Some(txs) = events_sub_map.get(&search_key) {
170 let key: U32<BigEndian> = event.block_number.into();
171 let block_events =
172 if let Ok(Some(event_bytes)) = self.trees.block_events.get(key.as_bytes()) {
173 vec![crate::Block {
174 block_number: event.block_number,
175 bytes: event_bytes.to_vec(),
176 }]
177 } else {
178 vec![]
179 };
180
181 let msg = ResponseMessage::Events {
182 key: search_key,
183 events: vec![event],
184 block_events,
185 };
186 for tx in txs.iter() {
187 if tx.send(msg.clone()).is_ok() {}
188 }
189 }
190 }
191
192 pub fn index_event(
193 &self,
194 key: Key<R::ChainKey>,
195 block_number: u32,
196 event_index: u16,
197 ) -> Result<(), sled::Error> {
198 key.write_db_key(&self.trees, block_number, event_index)?;
199 self.notify_subscribers(
200 key,
201 Event {
202 block_number,
203 event_index,
204 },
205 );
206 Ok(())
207 }
208}
209
210pub fn load_spans<R: RuntimeIndexer>(
211 span_db: &Tree,
212 index_variant: bool,
213 store_events: bool,
214) -> Result<Vec<Span>, IndexError> {
215 let mut spans = vec![];
216 'span: for (key, value) in span_db.into_iter().flatten() {
217 let span_value = SpanDbValue::read_from(&value).unwrap();
218 let start: u32 = span_value.start.into();
219 let mut end: u32 = u32::from_be_bytes(key.as_ref().try_into().unwrap());
220 if index_variant && (span_value.index_variant != 1) {
222 span_db.remove(key)?;
224 info!(
225 "📚 Re-indexing span of blocks from #{} to #{}.",
226 start.to_formatted_string(&Locale::en),
227 end.to_formatted_string(&Locale::en)
228 );
229 info!("📚 Reason: event variants not indexed.");
230 continue;
231 }
232 if store_events && (span_value.store_events != 1) {
234 span_db.remove(key)?;
236 info!(
237 "📚 Re-indexing span of blocks from #{} to #{}.",
238 start.to_formatted_string(&Locale::en),
239 end.to_formatted_string(&Locale::en)
240 );
241 info!("📚 Reason: events not stored.");
242 continue;
243 }
244 let span_version: u16 = span_value.version.into();
245 for (version, block_number) in R::get_versions().iter().enumerate() {
247 if span_version < version.try_into().unwrap() && end >= *block_number {
248 span_db.remove(key)?;
249 if start >= *block_number {
250 info!(
251 "📚 Re-indexing span of blocks from #{} to #{}.",
252 start.to_formatted_string(&Locale::en),
253 end.to_formatted_string(&Locale::en)
254 );
255 continue 'span;
256 }
257 info!(
258 "📚 Re-indexing span of blocks from #{} to #{}.",
259 block_number.to_formatted_string(&Locale::en),
260 end.to_formatted_string(&Locale::en)
261 );
262 end = block_number - 1;
264 span_db.insert(end.to_be_bytes(), value)?;
265 break;
266 }
267 }
268 let span = Span { start, end };
269 info!(
270 "📚 Previous span of indexed blocks from #{} to #{}.",
271 start.to_formatted_string(&Locale::en),
272 end.to_formatted_string(&Locale::en)
273 );
274 spans.push(span);
275 }
276 Ok(spans)
277}
278
279pub fn check_span(
280 span_db: &Tree,
281 spans: &mut Vec<Span>,
282 current_span: &mut Span,
283) -> Result<(), IndexError> {
284 while let Some(span) = spans.last() {
285 if current_span.start > span.start && current_span.start - 1 <= span.end {
287 let skipped = span.end - span.start + 1;
288 info!(
289 "📚 Skipping {} blocks from #{} to #{}",
290 skipped.to_formatted_string(&Locale::en),
291 span.start.to_formatted_string(&Locale::en),
292 span.end.to_formatted_string(&Locale::en),
293 );
294 current_span.start = span.start;
295 span_db.remove(span.end.to_be_bytes())?;
297 spans.pop();
298 } else {
299 break;
300 }
301 }
302 Ok(())
303}
304
305pub fn check_next_batch_block(spans: &[Span], next_batch_block: &mut u32) {
306 let mut i = spans.len();
308 while i != 0 {
309 i -= 1;
310 if *next_batch_block >= spans[i].start && *next_batch_block <= spans[i].end {
311 *next_batch_block = spans[i].start - 1;
312 }
313 }
314}
315
316pub fn process_sub_msg<R: RuntimeIndexer>(
317 indexer: &Indexer<R>,
318 msg: SubscriptionMessage<R::ChainKey>,
319) {
320 match msg {
321 SubscriptionMessage::SubscribeStatus { sub_response_tx } => {
322 let mut txs = indexer.status_sub.lock().unwrap();
323 txs.push(sub_response_tx);
324 }
325 SubscriptionMessage::UnsubscribeStatus { sub_response_tx } => {
326 let mut txs = indexer.status_sub.lock().unwrap();
327 txs.retain(|value| !sub_response_tx.same_channel(value));
328 }
329 SubscriptionMessage::SubscribeEvents {
330 key,
331 sub_response_tx,
332 } => {
333 let mut events_sub_map = indexer.events_sub_map.lock().unwrap();
334 match events_sub_map.get_mut(&key) {
335 Some(txs) => {
336 txs.push(sub_response_tx);
337 }
338 None => {
339 let txs = vec![sub_response_tx];
340 events_sub_map.insert(key, txs);
341 }
342 };
343 }
344 SubscriptionMessage::UnsubscribeEvents {
345 key,
346 sub_response_tx,
347 } => {
348 let mut events_sub_map = indexer.events_sub_map.lock().unwrap();
349 if let Some(txs) = events_sub_map.get_mut(&key) {
350 txs.retain(|value| !sub_response_tx.same_channel(value));
351 };
352 }
353 };
354}
355
356pub async fn substrate_index<R: RuntimeIndexer>(
357 trees: Trees<<R::ChainKey as IndexKey>::ChainTrees>,
358 api: OnlineClient<R::RuntimeConfig>,
359 rpc: LegacyRpcMethods<R::RuntimeConfig>,
360 finalized: bool,
361 queue_depth: u32,
362 index_variant: bool,
363 store_events: bool,
364 mut exit_rx: watch::Receiver<bool>,
365 mut sub_rx: mpsc::UnboundedReceiver<SubscriptionMessage<R::ChainKey>>,
366) -> Result<(), IndexError> {
367 info!(
368 "📇 Only index finalized blocks: {}",
369 match finalized {
370 false => "disabled",
371 true => "enabled",
372 },
373 );
374
375 let mut blocks_sub = if finalized {
376 api.blocks().subscribe_finalized().await
377 } else {
378 api.blocks().subscribe_best().await
379 }?;
380
381 info!(
382 "📇 Event variant indexing: {}",
383 match index_variant {
384 false => "disabled",
385 true => "enabled",
386 },
387 );
388 info!(
389 "📇 Event storage: {}",
390 match store_events {
391 false => "disabled",
392 true => "enabled",
393 },
394 );
395
396 let mut next_batch_block: u32 = blocks_sub
398 .next()
399 .await
400 .ok_or(IndexError::BlockNotFound(0))??
401 .number()
402 .into()
403 .try_into()
404 .unwrap();
405 info!(
406 "📚 Indexing backwards from #{}",
407 next_batch_block.to_formatted_string(&Locale::en)
408 );
409 let mut spans = load_spans::<R>(&trees.span, index_variant, store_events)?;
411 let mut current_span = if let Some(span) = spans.last()
413 && span.end == next_batch_block
414 {
415 let span = span.clone();
416 let skipped = span.end - span.start + 1;
417 info!(
418 "📚 Skipping {} blocks from #{} to #{}",
419 skipped.to_formatted_string(&Locale::en),
420 span.start.to_formatted_string(&Locale::en),
421 span.end.to_formatted_string(&Locale::en),
422 );
423 trees.span.remove(span.end.to_be_bytes())?;
425 spans.pop();
426 next_batch_block = span.start - 1;
427 span
428 } else {
429 Span {
430 start: next_batch_block + 1,
431 end: next_batch_block + 1,
432 }
433 };
434
435 let indexer = Indexer::<R>::new(trees.clone(), api, rpc, index_variant, store_events);
436
437 let mut head_future = Box::pin(indexer.index_head(blocks_sub.next()));
438
439 info!("📚 Queue depth: {}", queue_depth);
440 let mut futures = Vec::with_capacity(queue_depth.try_into().unwrap());
441
442 for _ in 0..queue_depth {
443 check_next_batch_block(&spans, &mut next_batch_block);
444 futures.push(Box::pin(indexer.index_block(next_batch_block)));
445 debug!(
446 "⬆️ Block #{} queued.",
447 next_batch_block.to_formatted_string(&Locale::en)
448 );
449 next_batch_block -= 1;
450 }
451
452 let mut orphans: AHashMap<u32, ()> = AHashMap::new();
453
454 let mut stats_block_count = 0;
455 let mut stats_event_count = 0;
456 let mut stats_key_count = 0;
457 let mut stats_start_time = Instant::now();
458
459 let interval_duration = Duration::from_millis(2000);
460 let mut interval = time::interval_at(Instant::now() + interval_duration, interval_duration);
461 interval.set_missed_tick_behavior(MissedTickBehavior::Skip);
462
463 let mut is_batching = true;
464
465 loop {
466 tokio::select! {
467 biased;
468
469 _ = exit_rx.changed() => {
470 if current_span.start != current_span.end {
471 let value = SpanDbValue {
472 start: current_span.start.into(),
473 version: (R::get_versions().len() - 1).try_into().unwrap(),
474 index_variant: index_variant.into(),
475 store_events: store_events.into(),
476 };
477 trees.span.insert(current_span.end.to_be_bytes(), value.as_bytes())?;
478 info!(
479 "📚 Recording current indexed span from #{} to #{}",
480 current_span.start.to_formatted_string(&Locale::en),
481 current_span.end.to_formatted_string(&Locale::en)
482 );
483 }
484 return Ok(());
485 }
486 Some(msg) = sub_rx.recv() => process_sub_msg(&indexer, msg),
487 result = &mut head_future => {
488 match result {
489 Ok((block_number, event_count, key_count)) => {
490 trees.span.remove(current_span.end.to_be_bytes())?;
491 current_span.end = block_number;
492 let value = SpanDbValue {
493 start: current_span.start.into(),
494 version: (R::get_versions().len() - 1).try_into().unwrap(),
495 index_variant: index_variant.into(),
496 store_events: store_events.into(),
497 };
498 trees.span.insert(current_span.end.to_be_bytes(), value.as_bytes())?;
499 info!(
500 "✨ #{}: {} events, {} keys",
501 block_number.to_formatted_string(&Locale::en),
502 event_count.to_formatted_string(&Locale::en),
503 key_count.to_formatted_string(&Locale::en),
504 );
505 indexer.notify_status_subscribers();
506 drop(head_future);
507 head_future = Box::pin(indexer.index_head(blocks_sub.next()));
508 },
509 Err(error) => {
510 match error {
511 IndexError::BlockNotFound(block_number) => {
512 error!("✨ Block not found #{}", block_number.to_formatted_string(&Locale::en));
513 },
514 err => {
515 error!("✨ Indexing failed: {}", err);
516 },
517 }
518 },
519 };
520 }
521 _ = interval.tick(), if is_batching => {
522 let current_time = Instant::now();
523 let duration = (current_time.duration_since(stats_start_time)).as_micros();
524 if duration != 0 {
525 info!(
526 "📚 #{}: {} blocks/sec, {} events/sec, {} keys/sec",
527 current_span.start.to_formatted_string(&Locale::en),
528 (<u32 as Into<u128>>::into(stats_block_count) * 1_000_000 / duration).to_formatted_string(&Locale::en),
529 (<u32 as Into<u128>>::into(stats_event_count) * 1_000_000 / duration).to_formatted_string(&Locale::en),
530 (<u32 as Into<u128>>::into(stats_key_count) * 1_000_000 / duration).to_formatted_string(&Locale::en),
531 );
532 }
533 stats_block_count = 0;
534 stats_event_count = 0;
535 stats_key_count = 0;
536 stats_start_time = current_time;
537 }
538 (result, index, _) = future::select_all(&mut futures), if is_batching => {
539 match result {
540 Ok((block_number, event_count, key_count)) => {
541 if block_number == current_span.start - 1 {
543 current_span.start = block_number;
544 debug!("⬇️ Block #{} indexed.", block_number.to_formatted_string(&Locale::en));
545 check_span(&trees.span, &mut spans, &mut current_span)?;
546 while orphans.contains_key(&(current_span.start - 1)) {
548 current_span.start -= 1;
549 orphans.remove(¤t_span.start);
550 debug!("➡️ Block #{} unorphaned.", current_span.start.to_formatted_string(&Locale::en));
551 check_span(&trees.span, &mut spans, &mut current_span)?;
552 }
553 }
554 else {
555 orphans.insert(block_number, ());
556 debug!("⬇️ Block #{} indexed and orphaned.", block_number.to_formatted_string(&Locale::en));
557 }
558 stats_block_count += 1;
559 stats_event_count += event_count;
560 stats_key_count += key_count;
561 },
562 Err(error) => {
563 match error {
564 IndexError::BlockNotFound(block_number) => {
565 error!("📚 Block not found #{}", block_number.to_formatted_string(&Locale::en));
566 is_batching = false;
567 },
568 _ => {
569 error!("📚 Batch indexing failed: {:?}", error);
570 is_batching = false;
571 },
572 }
573 }
574 }
575 check_next_batch_block(&spans, &mut next_batch_block);
576 futures[index] = Box::pin(indexer.index_block(next_batch_block));
577 debug!("⬆️ Block #{} queued.", next_batch_block.to_formatted_string(&Locale::en));
578 next_batch_block -= 1;
579 }
580 }
581 }
582}