Skip to main content

sui_indexer_alt_reader/
events.rs

1// Copyright (c) Mysten Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4use 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/// Key for fetching transaction events contents (Events, TimestampMs) by digest.
17#[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}