sui_indexer_alt_reader/
events.rs1use std::collections::HashMap;
5
6use anyhow::Context;
7use prost_types::FieldMask;
8use sui_rpc::field::FieldMaskUtil;
9use sui_rpc::proto::sui::rpc::v2 as proto;
10use sui_types::digests::TransactionDigest;
11
12use crate::error::Error;
13use crate::ledger_grpc_reader::ChunkedLoader;
14use crate::ledger_grpc_reader::LedgerGrpcReader;
15
16#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
18pub struct TransactionEventsKey(pub TransactionDigest);
19
20#[async_trait::async_trait]
21impl ChunkedLoader<TransactionEventsKey> for LedgerGrpcReader {
22 type Value = proto::ExecutedTransaction;
23 type Error = Error;
24
25 fn chunk_size(&self) -> usize {
26 self.max_batch_get_transactions()
27 }
28
29 async fn load_chunk(
30 &self,
31 keys: &[TransactionEventsKey],
32 ) -> Result<HashMap<TransactionEventsKey, proto::ExecutedTransaction>, Error> {
33 let digests = keys.iter().map(|key| key.0.to_string()).collect();
34
35 let mut request = proto::BatchGetTransactionsRequest::default();
36 request.digests = digests;
37 request.read_mask = Some(FieldMask::from_paths(["digest", "events.bcs", "timestamp"]));
38
39 let batch_response = self.batch_get_transactions(request).await?;
40
41 batch_response
42 .transactions
43 .into_iter()
44 .filter_map(|tx_result| match tx_result.result {
45 Some(proto::get_transaction_result::Result::Transaction(executed)) => {
46 Some(executed)
47 }
48 _ => None,
49 })
50 .map(|executed| {
51 let digest: TransactionDigest = executed
52 .digest
53 .as_ref()
54 .context("Missing transaction digest")?
55 .parse()
56 .context("Failed to parse transaction digest")?;
57
58 Ok((TransactionEventsKey(digest), executed))
59 })
60 .collect::<anyhow::Result<_>>()
61 .map_err(Error::from)
62 }
63}
64
65#[cfg(test)]
66mod tests {
67 use async_graphql::dataloader::Loader;
68
69 use super::*;
70 use crate::ledger_grpc_reader::test_support::assert_chunked;
71 use crate::ledger_grpc_reader::test_support::mock_reader;
72
73 #[tokio::test]
74 async fn load_chunks_oversized_batches() {
75 let (reader, mock, server) = mock_reader().await;
76 let limit = reader.max_batch_get_transactions();
77
78 let keys: Vec<TransactionEventsKey> = (0..limit + 50)
79 .map(|_| TransactionEventsKey(TransactionDigest::random()))
80 .collect();
81
82 let result = reader.load(&keys).await.expect("load should succeed");
83 assert!(result.is_empty());
84
85 let expected: Vec<String> = keys.iter().map(|key| key.0.to_string()).collect();
86 assert_chunked(mock.transaction_batches(), limit, &expected);
87
88 server.abort();
89 }
90}