Skip to main content

acuity_index_substrate/
substrate.rs

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        // Get the runtime version of the block.
85        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        // Check if variants are supposed to be indexed and they were not in this span.
221        if index_variant && (span_value.index_variant != 1) {
222            // Delete the span.
223            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        // Check if events are supposed to be stored and they were not in this span.
233        if store_events && (span_value.store_events != 1) {
234            // Delete the span.
235            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        // Loop through each indexer version.
246        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                // Truncate the span.
263                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        // Have we indexed all the blocks after the span?
286        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            // Remove the span.
296            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    // Figure out the next block to index, skipping the next span if we have reached it.
307    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    // Determine the correct block to start batch indexing.
397    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    // Load already indexed spans from the db.
410    let mut spans = load_spans::<R>(&trees.span, index_variant, store_events)?;
411    // If the first head block to be indexed will be touching the last span (the indexer was restarted), set the current span to the last span. Otherwise there will be no batch block indexed to connect the current span to the last span.
412    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        // Remove the span.
424        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                        // Is the new block contiguous to the current span or an orphan?
542                        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                            // Check if any orphans are now contiguous.
547                            while orphans.contains_key(&(current_span.start - 1)) {
548                                current_span.start -= 1;
549                                orphans.remove(&current_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}