sui_rpc_store/indexer/
event_bitmap.rs1use 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
40pub struct EventBitmap;
42
43pub 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 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 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}