Skip to main content

sui_rpc_store/indexer/
event_bitmap.rs

1// Copyright (c) Mysten Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4//! Sequential pipeline that populates the
5//! [`schema::event_bitmap`](crate::schema::event_bitmap) CF.
6//!
7//! For every event in every transaction the pipeline:
8//!
9//! 1. Visits dimension candidates via
10//!    [`sui_inverted_index::for_each_event_dimension`].
11//! 2. Encodes the dimension key via
12//!    [`sui_inverted_index::encode_dimension_key`].
13//! 3. Packs `(tx_seq, event_idx)` into the event-seq space the
14//!    schema uses (`encode_event_seq(tx_seq, event_idx)`), checks the packing
15//!    doesn't overflow the per-tx event limit (`1 << EVENT_BITS`)
16//!    or the `tx_seq` ceiling (`u64::MAX >> EVENT_BITS`), and
17//!    groups the bit into `(dim_key, packed / EVENT_BUCKET_SIZE)`.
18
19use std::collections::HashMap;
20use std::sync::Arc;
21
22use async_trait::async_trait;
23use roaring::RoaringBitmap;
24use sui_indexer_alt_framework::pipeline::Processor;
25use sui_indexer_alt_framework::pipeline::sequential;
26use sui_inverted_index::encode_dimension_key;
27use sui_inverted_index::event_seq::{MAX_EVENTS_PER_TX, MAX_TX_SEQ};
28use sui_inverted_index::for_each_event_dimension;
29use sui_types::full_checkpoint_content::Checkpoint;
30use sui_types::transaction::TransactionDataAPI;
31
32use crate::indexer::Schema;
33use crate::indexer::Store;
34use crate::indexer::tx_seq_at;
35use crate::schema::event_bitmap;
36use crate::schema::event_bitmap::bit_of;
37use crate::schema::event_bitmap::bucket_of;
38use sui_inverted_index::event_seq::encode_event_seq;
39
40/// Pipeline marker for `event_bitmap`.
41pub struct EventBitmap;
42
43/// One pre-built bitmap for a single `(dimension_key, bucket)`
44/// pair, ready to be staged as a merge operand against the CF.
45pub struct Row {
46    pub dimension_key: Vec<u8>,
47    pub bucket: u64,
48    pub bitmap: RoaringBitmap,
49}
50
51#[async_trait]
52impl Processor for EventBitmap {
53    const NAME: &'static str = "event_bitmap";
54    type Value = Row;
55
56    async fn process(&self, checkpoint: &Arc<Checkpoint>) -> anyhow::Result<Vec<Row>> {
57        let mut groups: HashMap<(Vec<u8>, u64), RoaringBitmap> = HashMap::new();
58
59        for (i, tx) in checkpoint.transactions.iter().enumerate() {
60            let tx_seq = tx_seq_at(checkpoint, i);
61            if tx_seq > MAX_TX_SEQ {
62                anyhow::bail!("tx_seq {tx_seq} exceeds packed event-seq limit {MAX_TX_SEQ}",);
63            }
64            let sender = tx.transaction.sender();
65
66            // `for_each_event_dimension` can't propagate errors,
67            // so a packing failure has to be captured and
68            // surfaced afterwards.
69            let mut packing_error: Option<anyhow::Error> = None;
70            for_each_event_dimension(
71                sender,
72                &tx.effects,
73                tx.events.as_ref(),
74                |event_idx, dim, value| {
75                    if packing_error.is_some() {
76                        return;
77                    }
78                    if event_idx >= MAX_EVENTS_PER_TX {
79                        packing_error = Some(anyhow::anyhow!(
80                            "event_idx {event_idx} exceeds packed event-seq limit {}",
81                            MAX_EVENTS_PER_TX - 1,
82                        ));
83                        return;
84                    }
85                    let packed = encode_event_seq(tx_seq, event_idx);
86                    let bucket = bucket_of(packed);
87                    let bit = bit_of(packed);
88                    groups
89                        .entry((encode_dimension_key(dim, value), bucket))
90                        .or_default()
91                        .insert(bit);
92                },
93            );
94            if let Some(e) = packing_error {
95                return Err(e);
96            }
97        }
98
99        Ok(groups
100            .into_iter()
101            .map(|((dim_key, bucket), bitmap)| Row {
102                dimension_key: dim_key,
103                bucket,
104                bitmap,
105            })
106            .collect())
107    }
108}
109
110#[async_trait]
111impl sequential::Handler for EventBitmap {
112    type Store = Store;
113    /// Fold operands from multiple checkpoints together so the
114    /// commit path stages at most one merge operand per
115    /// `(dim_key, bucket)` per commit.
116    type Batch = HashMap<(Vec<u8>, u64), RoaringBitmap>;
117
118    fn batch(&self, batch: &mut Self::Batch, values: std::vec::IntoIter<Row>) {
119        for row in values {
120            let entry = batch.entry((row.dimension_key, row.bucket)).or_default();
121            *entry |= row.bitmap;
122        }
123    }
124
125    async fn commit<'a>(
126        &self,
127        batch: &Self::Batch,
128        conn: &mut sui_consistent_store::Connection<'a, Schema>,
129    ) -> anyhow::Result<usize> {
130        let cf = &conn.store.schema().event_bitmap;
131        for ((dim_key, bucket), bitmap) in batch {
132            let (k, v) = event_bitmap::store_bitmap(dim_key.clone(), *bucket, bitmap.clone());
133            conn.batch.merge(cf, &k, &v)?;
134        }
135        Ok(batch.len())
136    }
137}
138
139#[cfg(test)]
140mod tests {
141    use std::sync::Arc;
142
143    use sui_types::test_checkpoint_data_builder::TestCheckpointBuilder;
144
145    use super::*;
146
147    #[tokio::test]
148    async fn process_runs_against_synthetic_checkpoint() {
149        let checkpoint = Arc::new(TestCheckpointBuilder::new(1).build_checkpoint());
150        let _ = EventBitmap.process(&checkpoint).await.unwrap();
151    }
152}