Skip to main content

sui_json_rpc/
read_api.rs

1// Copyright (c) Mysten Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4use std::collections::{BTreeMap, HashMap};
5use std::sync::Arc;
6use std::time::Duration;
7
8use anyhow::anyhow;
9use async_trait::async_trait;
10use backoff::ExponentialBackoff;
11use backoff::future::retry;
12use fastcrypto::encoding::Base64;
13use fastcrypto_zkp::bn254::zk_login_api::ZkLoginEnv;
14use futures::future::join_all;
15use imbl::hashmap::HashMap as ImHashMap;
16use indexmap::map::IndexMap;
17use itertools::Itertools;
18use jsonrpsee::RpcModule;
19use jsonrpsee::core::RpcResult;
20use move_bytecode_utils::module_cache::GetModule;
21use move_core_types::account_address::AccountAddress;
22use move_core_types::annotated_value::{MoveStructLayout, MoveTypeLayout};
23use move_core_types::language_storage::StructTag;
24use mysten_common::ZipDebugEqIteratorExt;
25use once_cell::sync::Lazy;
26use serde_json::Value as Json;
27use shared_crypto::intent::{IntentMessage, PersonalMessage};
28use sui_display::v1::Format;
29use sui_json_rpc_types::ZkLoginIntentScope;
30use sui_types::base_types::SuiAddress;
31use sui_types::signature::{GenericSignature, VerifyParams};
32use sui_types::signature_verification::VerifiedDigestCache;
33use sui_types::storage::ObjectKey;
34use tap::TapFallible;
35use tracing::{debug, error, info, instrument, trace, warn};
36
37use mysten_metrics::add_server_timing;
38use sui_core::authority::AuthorityState;
39use sui_json_rpc_api::{
40    JsonRpcMetrics, QUERY_MAX_RESULT_LIMIT, QUERY_MAX_RESULT_LIMIT_CHECKPOINTS, ReadApiOpenRpc,
41    ReadApiServer, validate_limit,
42};
43use sui_json_rpc_types::{
44    BalanceChange, Checkpoint, CheckpointId, CheckpointPage, DisplayFieldsResponse, EventFilter,
45    ObjectChange, ProtocolConfigResponse, SuiEvent, SuiGetPastObjectRequest, SuiObjectDataOptions,
46    SuiObjectResponse, SuiPastObjectResponse, SuiTransactionBlock, SuiTransactionBlockEvents,
47    SuiTransactionBlockResponse, SuiTransactionBlockResponseOptions,
48};
49use sui_open_rpc::Module;
50use sui_protocol_config::{ProtocolConfig, ProtocolVersion};
51use sui_storage::key_value_store::TransactionKeyValueStore;
52use sui_types::base_types::{ObjectID, SequenceNumber, TransactionDigest};
53use sui_types::display::DisplayVersionUpdatedEvent;
54use sui_types::display_registry;
55use sui_types::effects::{TransactionEffects, TransactionEffectsAPI, TransactionEvents};
56use sui_types::error::{SuiError, SuiObjectResponseError};
57use sui_types::messages_checkpoint::{CheckpointSequenceNumber, CheckpointTimestamp};
58use sui_types::object::{Object, ObjectRead, PastObjectRead};
59use sui_types::sui_serde::BigInt;
60use sui_types::transaction::TransactionDataAPI;
61use sui_types::transaction::{Transaction, TransactionData};
62
63use crate::authority_state::{StateRead, StateReadError, StateReadResult};
64use crate::error::{Error, RpcInterimResult, SuiRpcInputError};
65use crate::{ObjectProvider, with_tracing};
66use crate::{
67    ObjectProviderCache, SuiRpcModule, get_balance_changes_from_effect, get_object_changes,
68};
69use fastcrypto::encoding::Encoding;
70use fastcrypto::traits::ToFromBytes;
71use shared_crypto::intent::Intent;
72use sui_json_rpc_types::ZkLoginVerifyResult;
73use sui_types::authenticator_state::{ActiveJwk, get_authenticator_state};
74
75/// Default max depth used while converting rendered Display values to JSON.
76const DEFAULT_MAX_DISPLAY_MOVE_VALUE_DEPTH: usize = 32;
77
78/// Default budget for Display output size.
79const DEFAULT_MAX_DISPLAY_OUTPUT_SIZE: usize = 1024 * 1024;
80
81/// A field access in a Display string cannot exceed this level of nesting.
82static MAX_DISPLAY_FIELD_DEPTH: Lazy<usize> = Lazy::new(|| {
83    let max_opt = std::env::var("MAX_DISPLAY_FIELD_DEPTH")
84        .ok()
85        .and_then(|s| s.parse().ok());
86
87    if let Some(max) = max_opt {
88        info!("Using custom value for 'MAX_DISPLAY_FIELD_DEPTH': {max}");
89        max
90    } else {
91        sui_display::v2::Limits::default().max_depth
92    }
93});
94
95/// Parser node budget for Display v2.
96static MAX_DISPLAY_FORMAT_NODES: Lazy<usize> = Lazy::new(|| {
97    let max_opt = std::env::var("MAX_DISPLAY_FORMAT_NODES")
98        .ok()
99        .and_then(|s| s.parse().ok());
100
101    if let Some(max) = max_opt {
102        info!("Using custom value for 'MAX_DISPLAY_FORMAT_NODES': {max}");
103        max
104    } else {
105        sui_display::v2::Limits::default().max_nodes
106    }
107});
108
109/// Max object loads budget for Display v2.
110static MAX_DISPLAY_OBJECT_LOADS: Lazy<usize> = Lazy::new(|| {
111    let max_opt = std::env::var("MAX_DISPLAY_OBJECT_LOADS")
112        .ok()
113        .and_then(|s| s.parse().ok());
114
115    if let Some(max) = max_opt {
116        info!("Using custom value for 'MAX_DISPLAY_OBJECT_LOADS': {max}");
117        max
118    } else {
119        sui_display::v2::Limits::default().max_loads
120    }
121});
122
123/// Maximum depth used while converting rendered Display values to JSON.
124static MAX_DISPLAY_MOVE_VALUE_DEPTH: Lazy<usize> = Lazy::new(|| {
125    let max_opt = std::env::var("MAX_MOVE_VALUE_DEPTH")
126        .ok()
127        .and_then(|s| s.parse().ok());
128
129    if let Some(max) = max_opt {
130        info!("Using custom value for 'MAX_MOVE_VALUE_DEPTH': {max}");
131        max
132    } else {
133        DEFAULT_MAX_DISPLAY_MOVE_VALUE_DEPTH
134    }
135});
136
137/// Overall display output cannot exceed this size.
138static MAX_DISPLAY_OUTPUT_SIZE: Lazy<usize> = Lazy::new(|| {
139    let max_opt = std::env::var("MAX_DISPLAY_OUTPUT_SIZE")
140        .ok()
141        .and_then(|s| s.parse().ok());
142
143    if let Some(max) = max_opt {
144        info!("Using custom value for 'MAX_DISPLAY_OUTPUT_SIZE': {max}");
145        max
146    } else {
147        DEFAULT_MAX_DISPLAY_OUTPUT_SIZE
148    }
149});
150
151struct DisplayStore<'s> {
152    state: &'s dyn StateRead,
153}
154
155impl<'s> DisplayStore<'s> {
156    fn new(state: &'s dyn StateRead) -> Self {
157        Self { state }
158    }
159}
160
161#[async_trait]
162impl sui_display::v2::Store for DisplayStore<'_> {
163    async fn latest(
164        &self,
165        id: AccountAddress,
166    ) -> anyhow::Result<Option<(MoveTypeLayout, Vec<u8>)>> {
167        let read = self.state.get_object_read(&id.into())?;
168        let ObjectRead::Exists(_, object, Some(layout)) = read else {
169            return Ok(None);
170        };
171
172        let Some(move_object) = object.data.try_as_move() else {
173            return Ok(None);
174        };
175
176        Ok(Some((
177            MoveTypeLayout::Struct(Box::new(layout)),
178            move_object.contents().to_vec(),
179        )))
180    }
181}
182
183// An implementation of the read portion of the JSON-RPC interface intended for use in
184// Fullnodes.
185#[derive(Clone)]
186pub struct ReadApi {
187    pub state: Arc<dyn StateRead>,
188    pub transaction_kv_store: Arc<TransactionKeyValueStore>,
189    pub metrics: Arc<JsonRpcMetrics>,
190}
191
192// Internal data structure to make it easy to work with data returned from
193// authority store and also enable code sharing between get_transaction_with_options,
194// multi_get_transaction_with_options, etc.
195#[derive(Default)]
196struct IntermediateTransactionResponse {
197    digest: TransactionDigest,
198    transaction: Option<Transaction>,
199    effects: Option<TransactionEffects>,
200    events: Option<SuiTransactionBlockEvents>,
201    checkpoint_seq: Option<CheckpointSequenceNumber>,
202    balance_changes: Option<Vec<BalanceChange>>,
203    object_changes: Option<Vec<ObjectChange>>,
204    timestamp: Option<CheckpointTimestamp>,
205    errors: Vec<String>,
206}
207
208impl IntermediateTransactionResponse {
209    pub fn new(digest: TransactionDigest) -> Self {
210        Self {
211            digest,
212            ..Default::default()
213        }
214    }
215
216    pub fn transaction(&self) -> &Option<Transaction> {
217        &self.transaction
218    }
219}
220
221impl ReadApi {
222    pub fn new(
223        state: Arc<AuthorityState>,
224        transaction_kv_store: Arc<TransactionKeyValueStore>,
225        metrics: Arc<JsonRpcMetrics>,
226    ) -> Self {
227        Self {
228            state,
229            transaction_kv_store,
230            metrics,
231        }
232    }
233
234    async fn get_checkpoint_internal(&self, id: CheckpointId) -> Result<Checkpoint, Error> {
235        Ok(match id {
236            CheckpointId::SequenceNumber(seq) => {
237                let verified_summary = self
238                    .transaction_kv_store
239                    .get_checkpoint_summary(seq)
240                    .await?;
241                let content = self
242                    .transaction_kv_store
243                    .get_checkpoint_contents(verified_summary.sequence_number)
244                    .await?;
245                let signature = verified_summary.auth_sig().signature.clone();
246                (verified_summary.into_data(), content, signature).into()
247            }
248            CheckpointId::Digest(digest) => {
249                let verified_summary = self
250                    .transaction_kv_store
251                    .get_checkpoint_summary_by_digest(digest)
252                    .await?;
253                let content = self
254                    .transaction_kv_store
255                    .get_checkpoint_contents(verified_summary.sequence_number)
256                    .await?;
257                let signature = verified_summary.auth_sig().signature.clone();
258                (verified_summary.into_data(), content, signature).into()
259            }
260        })
261    }
262
263    pub async fn get_checkpoints_internal(
264        state: Arc<dyn StateRead>,
265        transaction_kv_store: Arc<TransactionKeyValueStore>,
266        // If `Some`, the query will start from the next item after the specified cursor
267        cursor: Option<CheckpointSequenceNumber>,
268        limit: u64,
269        descending_order: bool,
270    ) -> StateReadResult<Vec<Checkpoint>> {
271        let max_checkpoint = state.get_latest_checkpoint_sequence_number()?;
272        let checkpoint_numbers =
273            calculate_checkpoint_numbers(cursor, limit, descending_order, max_checkpoint);
274
275        let verified_checkpoints = transaction_kv_store
276            .multi_get_checkpoints_summaries(&checkpoint_numbers)
277            .await?;
278        let checkpoint_contents = transaction_kv_store
279            .multi_get_checkpoints_contents(&checkpoint_numbers)
280            .await?;
281
282        // Summaries and contents are resolved from separate tables, and checkpoint
283        // pruning can delete a checkpoint's contents while leaving its
284        // sequence-addressable summary in place. Pair each summary with the
285        // contents for the *same* sequence number by zipping the two `Option`
286        // vectors index-by-index. Independently dropping the `None`s and zipping
287        // the dense remainders would shift later contents onto earlier summaries,
288        // yielding response rows whose summary and transaction list describe
289        // different checkpoints.
290        let mut checkpoints = Vec::with_capacity(checkpoint_numbers.len());
291        for (maybe_summary, maybe_contents) in verified_checkpoints
292            .into_iter()
293            .zip_debug_eq(checkpoint_contents)
294        {
295            // Skip any sequence number whose summary or contents are unavailable
296            // (e.g. pruned) rather than pairing it with another checkpoint's data.
297            let (Some(summary), Some(contents)) = (maybe_summary, maybe_contents) else {
298                continue;
299            };
300            let signature = summary.auth_sig().signature.clone();
301            let summary = summary.into_summary_and_sequence().1;
302            checkpoints.push(Checkpoint::from((summary, contents, signature)));
303        }
304
305        Ok(checkpoints)
306    }
307
308    #[instrument(skip_all)]
309    async fn multi_get_transaction_blocks_internal(
310        &self,
311        digests: Vec<TransactionDigest>,
312        opts: Option<SuiTransactionBlockResponseOptions>,
313    ) -> Result<Vec<SuiTransactionBlockResponse>, Error> {
314        trace!("start");
315
316        let num_digests = digests.len();
317        if num_digests > *QUERY_MAX_RESULT_LIMIT {
318            Err(SuiRpcInputError::SizeLimitExceeded(
319                QUERY_MAX_RESULT_LIMIT.to_string(),
320            ))?
321        }
322        self.metrics
323            .get_tx_blocks_limit
324            .observe(digests.len() as f64);
325
326        let opts = opts.unwrap_or_default();
327
328        // use LinkedHashMap to dedup and can iterate in insertion order.
329        let mut temp_response: IndexMap<&TransactionDigest, IntermediateTransactionResponse> =
330            IndexMap::from_iter(
331                digests
332                    .iter()
333                    .map(|k| (k, IntermediateTransactionResponse::new(*k))),
334            );
335        if temp_response.len() < num_digests {
336            Err(SuiRpcInputError::ContainsDuplicates)?
337        }
338
339        if opts.require_input() {
340            trace!("getting input");
341            let digests_clone = digests.clone();
342            let transactions =
343                self.transaction_kv_store.multi_get_tx(&digests_clone).await.tap_err(
344                    |err| debug!(digests=?digests_clone, "Failed to multi get transactions: {:?}", err),
345                )?;
346
347            for ((_digest, cache_entry), txn) in temp_response
348                .iter_mut()
349                .zip_debug_eq(transactions.into_iter())
350            {
351                cache_entry.transaction = txn;
352            }
353        }
354
355        // Fetch effects when `show_events` is true because events relies on effects
356        if opts.require_effects() {
357            trace!("getting effects");
358            let digests_clone = digests.clone();
359            let effects_list = self.transaction_kv_store
360                .multi_get_fx_by_tx_digest(&digests_clone)
361                .await
362                .tap_err(
363                    |err| debug!(digests=?digests_clone, "Failed to multi get effects for transactions: {:?}", err),
364                )?;
365            for ((_digest, cache_entry), e) in temp_response
366                .iter_mut()
367                .zip_debug_eq(effects_list.into_iter())
368            {
369                cache_entry.effects = e;
370            }
371        }
372
373        trace!("getting checkpoint sequence numbers");
374        let checkpoint_seq_list = self
375            .transaction_kv_store
376            .multi_get_transaction_checkpoint(&digests)
377            .await
378            .tap_err(
379                |err| debug!(digests=?digests, "Failed to multi get checkpoint sequence number: {:?}", err))?;
380        for ((_digest, cache_entry), seq) in temp_response
381            .iter_mut()
382            .zip_debug_eq(checkpoint_seq_list.into_iter())
383        {
384            cache_entry.checkpoint_seq = seq;
385        }
386
387        let unique_checkpoint_numbers = temp_response
388            .values()
389            .filter_map(|cache_entry| cache_entry.checkpoint_seq)
390            // It's likely that many transactions have the same checkpoint, so we don't
391            // need to over-fetch
392            .unique()
393            .collect::<Vec<CheckpointSequenceNumber>>();
394
395        // fetch timestamp from the DB
396        trace!("getting checkpoint summaries");
397        let timestamps = self
398            .transaction_kv_store
399            .multi_get_checkpoints_summaries(&unique_checkpoint_numbers)
400            .await
401            .map_err(|e| {
402                Error::UnexpectedError(format!("Failed to fetch checkpoint summaries by these checkpoint ids: {unique_checkpoint_numbers:?} with error: {e:?}"))
403            })?
404            .into_iter()
405            .map(|c| c.map(|checkpoint| checkpoint.timestamp_ms));
406
407        // construct a hashmap of checkpoint -> timestamp for fast lookup
408        let checkpoint_to_timestamp = unique_checkpoint_numbers
409            .into_iter()
410            .zip_debug_eq(timestamps)
411            .collect::<HashMap<_, _>>();
412
413        // fill cache with the timestamp
414        for (_, cache_entry) in temp_response.iter_mut() {
415            if let Some(checkpoint_seq) = cache_entry.checkpoint_seq.as_ref() {
416                cache_entry.timestamp = *checkpoint_to_timestamp
417                    .get(checkpoint_seq)
418                    // Safe to unwrap because checkpoint_seq is guaranteed to exist in checkpoint_to_timestamp
419                    .unwrap();
420            }
421        }
422
423        if opts.show_events {
424            trace!("getting events");
425            let mut non_empty_digests = vec![];
426            for cache_entry in temp_response.values() {
427                if let Some(effects) = &cache_entry.effects
428                    && effects.events_digest().is_some()
429                {
430                    non_empty_digests.push(cache_entry.digest);
431                }
432            }
433            // fetch events from the DB with retry, retry each 0.5s for 3s
434            let backoff = ExponentialBackoff {
435                max_elapsed_time: Some(Duration::from_secs(3)),
436                multiplier: 1.0,
437                ..ExponentialBackoff::default()
438            };
439            let mut events = retry(backoff, || async {
440                match self
441                    .transaction_kv_store
442                    .multi_get_events_by_tx_digests(&non_empty_digests)
443                    .await
444                {
445                    // Only return Ok when all the queried transaction events are found, otherwise retry
446                    // until timeout, then return Err.
447                    Ok(events) if !events.contains(&None) => Ok(events),
448                    Ok(_) => Err(backoff::Error::transient(Error::UnexpectedError(
449                        "Events not found, transaction execution may be incomplete.".into(),
450                    ))),
451                    Err(e) => Err(backoff::Error::permanent(Error::UnexpectedError(format!(
452                        "Failed to call multi_get_events: {e:?}"
453                    )))),
454                }
455            })
456            .await
457            .map_err(|e| {
458                Error::UnexpectedError(format!(
459                    "Retrieving events with retry failed for transaction digests {digests:?}: {e:?}"
460                ))
461            })?
462            .into_iter();
463
464            // fill cache with the events
465            for (_, cache_entry) in temp_response.iter_mut() {
466                let transaction_digest = cache_entry.digest;
467                if let Some(events_digest) =
468                    cache_entry.effects.as_ref().and_then(|e| e.events_digest())
469                {
470                    match events.next() {
471                        Some(Some(ev)) => {
472                            cache_entry.events =
473                                Some(to_sui_transaction_events(self, cache_entry.digest, ev)?)
474                        }
475                        None | Some(None) => {
476                            error!(
477                                "Failed to fetch events with event digest {events_digest:?} for txn {transaction_digest}"
478                            );
479                            cache_entry.errors.push(format!(
480                                "Failed to fetch events with event digest {events_digest:?}",
481                            ))
482                        }
483                    }
484                } else {
485                    // events field will be Some if and only if `show_events` is true and
486                    // there is no error in converting fetching events
487                    cache_entry.events = Some(SuiTransactionBlockEvents::default());
488                }
489            }
490        }
491
492        let mut object_cache =
493            ObjectProviderCache::new((self.state.clone(), self.transaction_kv_store.clone()));
494
495        // Prefetch the objects if we need to show balance or object changes
496        if opts.show_balance_changes || opts.show_object_changes {
497            let mut keys = vec![];
498            for resp in temp_response.values() {
499                let effects = resp.effects.as_ref().ok_or_else(|| {
500                    SuiRpcInputError::GenericNotFound(
501                        "unable to derive balance/object changes because effect is empty"
502                            .to_string(),
503                    )
504                })?;
505
506                for change in effects.object_changes() {
507                    if let Some(input_version) = change.input_version {
508                        keys.push(ObjectKey(change.id, input_version));
509                    }
510                    if let Some(output_version) = change.output_version {
511                        keys.push(ObjectKey(change.id, output_version));
512                    }
513                }
514            }
515
516            let objects = self
517                .transaction_kv_store
518                .multi_get_objects(&keys)
519                .await?
520                .into_iter()
521                .flatten()
522                .collect::<Vec<_>>();
523
524            object_cache.insert_objects_into_cache(objects);
525        }
526
527        if opts.show_balance_changes {
528            trace!("getting balance changes");
529
530            let mut results = vec![];
531            for resp in temp_response.values() {
532                let input_objects = if let Some(tx) = resp.transaction() {
533                    tx.data()
534                        .inner()
535                        .intent_message
536                        .value
537                        .input_objects()
538                        .unwrap_or_default()
539                } else {
540                    // don't have the input tx, so not much we can do. perhaps this is an Err?
541                    Vec::new()
542                };
543                results.push(get_balance_changes_from_effect(
544                    &object_cache,
545                    resp.effects.as_ref().ok_or_else(|| {
546                        SuiRpcInputError::GenericNotFound(
547                            "unable to derive balance changes because effect is empty".to_string(),
548                        )
549                    })?,
550                    input_objects,
551                    None,
552                ));
553            }
554            let results = join_all(results).await;
555            for (result, entry) in results.into_iter().zip_debug_eq(temp_response.iter_mut()) {
556                match result {
557                    Ok(balance_changes) => entry.1.balance_changes = Some(balance_changes),
558                    Err(e) => entry
559                        .1
560                        .errors
561                        .push(format!("Failed to fetch balance changes {e:?}")),
562                }
563            }
564        }
565
566        if opts.show_object_changes {
567            trace!("getting object changes");
568
569            let mut results = vec![];
570            for resp in temp_response.values() {
571                let effects = resp.effects.as_ref().ok_or_else(|| {
572                    SuiRpcInputError::GenericNotFound(
573                        "unable to derive object changes because effect is empty".to_string(),
574                    )
575                })?;
576
577                results.push(get_object_changes(
578                    &object_cache,
579                    effects,
580                    resp.transaction
581                        .as_ref()
582                        .ok_or_else(|| {
583                            SuiRpcInputError::GenericNotFound(
584                                "unable to derive object changes because transaction is empty"
585                                    .to_string(),
586                            )
587                        })?
588                        .data()
589                        .intent_message()
590                        .value
591                        .sender(),
592                    effects.modified_at_versions(),
593                    effects.all_changed_objects(),
594                    effects.all_removed_objects(),
595                ));
596            }
597            let results = join_all(results).await;
598            for (result, entry) in results.into_iter().zip_debug_eq(temp_response.iter_mut()) {
599                match result {
600                    Ok(object_changes) => entry.1.object_changes = Some(object_changes),
601                    Err(e) => entry
602                        .1
603                        .errors
604                        .push(format!("Failed to fetch object changes {e:?}")),
605                }
606            }
607        }
608
609        let epoch_store = self.state.load_epoch_store_one_call_per_task();
610        let converted_tx_block_resps = temp_response
611            .into_iter()
612            .map(|c| convert_to_response(c.1, &opts, epoch_store.module_cache()))
613            .collect::<Result<Vec<_>, _>>()?;
614
615        self.metrics
616            .get_tx_blocks_result_size
617            .observe(converted_tx_block_resps.len() as f64);
618        self.metrics
619            .get_tx_blocks_result_size_total
620            .inc_by(converted_tx_block_resps.len() as u64);
621
622        trace!("done");
623
624        Ok(converted_tx_block_resps)
625    }
626}
627
628#[async_trait]
629impl ReadApiServer for ReadApi {
630    #[instrument(skip(self))]
631    async fn get_object(
632        &self,
633        object_id: ObjectID,
634        options: Option<SuiObjectDataOptions>,
635    ) -> RpcResult<SuiObjectResponse> {
636        with_tracing!(async move {
637            let state = self.state.clone();
638            let object_read = state.get_object_read(&object_id).map_err(|e| {
639                warn!(?object_id, "Failed to get object: {:?}", e);
640                Error::from(e)
641            })?;
642            let options = options.unwrap_or_default();
643
644            match object_read {
645                ObjectRead::NotExists(id) => Ok(SuiObjectResponse::new_with_error(
646                    SuiObjectResponseError::NotExists { object_id: id },
647                )),
648                ObjectRead::Exists(object_ref, o, layout) => {
649                    let mut display_fields = None;
650                    if options.show_display {
651                        match get_display_fields(self, &self.transaction_kv_store, &o, &layout)
652                            .await
653                        {
654                            Ok(rendered_fields) => display_fields = Some(rendered_fields),
655                            Err(e) => {
656                                return Ok(SuiObjectResponse::new(
657                                    Some((object_ref, o, layout, options, None).try_into()?),
658                                    Some(SuiObjectResponseError::DisplayError {
659                                        error: e.to_string(),
660                                    }),
661                                ));
662                            }
663                        }
664                    }
665                    Ok(SuiObjectResponse::new_with_data(
666                        (object_ref, o, layout, options, display_fields).try_into()?,
667                    ))
668                }
669                ObjectRead::Deleted((object_id, version, digest)) => Ok(
670                    SuiObjectResponse::new_with_error(SuiObjectResponseError::Deleted {
671                        object_id,
672                        version,
673                        digest,
674                    }),
675                ),
676            }
677        })
678    }
679
680    #[instrument(skip(self))]
681    async fn multi_get_objects(
682        &self,
683        object_ids: Vec<ObjectID>,
684        options: Option<SuiObjectDataOptions>,
685    ) -> RpcResult<Vec<SuiObjectResponse>> {
686        with_tracing!(async move {
687            if object_ids.len() <= *QUERY_MAX_RESULT_LIMIT {
688                self.metrics
689                    .get_objects_limit
690                    .observe(object_ids.len() as f64);
691                let mut futures = vec![];
692                for object_id in object_ids {
693                    futures.push(self.get_object(object_id, options.clone()));
694                }
695                let results = join_all(futures).await;
696
697                let objects_result: Result<Vec<SuiObjectResponse>, String> = results
698                    .into_iter()
699                    .map(|result| match result {
700                        Ok(response) => Ok(response),
701                        Err(error) => {
702                            error!("Failed to fetch object with error: {error:?}");
703                            Err(format!("Error: {}", error))
704                        }
705                    })
706                    .collect();
707
708                let objects = objects_result.map_err(|err| {
709                    Error::UnexpectedError(format!("Failed to fetch objects with error: {}", err))
710                })?;
711
712                self.metrics
713                    .get_objects_result_size
714                    .observe(objects.len() as f64);
715                self.metrics
716                    .get_objects_result_size_total
717                    .inc_by(objects.len() as u64);
718                Ok(objects)
719            } else {
720                Err(SuiRpcInputError::SizeLimitExceeded(
721                    QUERY_MAX_RESULT_LIMIT.to_string(),
722                ))?
723            }
724        })
725    }
726
727    #[instrument(skip(self))]
728    async fn try_get_past_object(
729        &self,
730        object_id: ObjectID,
731        version: SequenceNumber,
732        options: Option<SuiObjectDataOptions>,
733    ) -> RpcResult<SuiPastObjectResponse> {
734        with_tracing!(async move {
735            let state = self.state.clone();
736            let past_read = state
737                .get_past_object_read(&object_id, version)
738                .map_err(|e| {
739                    error!("Failed to call try_get_past_object for object: {object_id:?} version: {version:?} with error: {e:?}");
740                    Error::from(e)
741                })?;
742            let options = options.unwrap_or_default();
743            match past_read {
744                PastObjectRead::ObjectNotExists(id) => {
745                    Ok(SuiPastObjectResponse::ObjectNotExists(id))
746                }
747                PastObjectRead::VersionFound(object_ref, o, layout) => {
748                    let display_fields = if options.show_display {
749                        // TODO (jian): api breaking change to also modify past objects.
750                        Some(
751                            get_display_fields(self, &self.transaction_kv_store, &o, &layout)
752                                .await
753                                .map_err(|e| {
754                                    Error::UnexpectedError(format!(
755                                        "Unable to render object at version {version}: {e}"
756                                    ))
757                                })?,
758                        )
759                    } else {
760                        None
761                    };
762                    Ok(SuiPastObjectResponse::VersionFound(
763                        (object_ref, o, layout, options, display_fields).try_into()?,
764                    ))
765                }
766                PastObjectRead::ObjectDeleted(oref) => {
767                    Ok(SuiPastObjectResponse::ObjectDeleted(oref.into()))
768                }
769                PastObjectRead::VersionNotFound(id, seq_num) => {
770                    Ok(SuiPastObjectResponse::VersionNotFound(id, seq_num))
771                }
772                PastObjectRead::VersionTooHigh {
773                    object_id,
774                    asked_version,
775                    latest_version,
776                } => Ok(SuiPastObjectResponse::VersionTooHigh {
777                    object_id,
778                    asked_version,
779                    latest_version,
780                }),
781            }
782        })
783    }
784
785    #[instrument(skip(self))]
786    async fn try_get_object_before_version(
787        &self,
788        object_id: ObjectID,
789        version: SequenceNumber,
790    ) -> RpcResult<SuiPastObjectResponse> {
791        let version = self
792            .state
793            .find_object_lt_or_eq_version(&object_id, &version)
794            .await
795            .map_err(Error::from)?
796            .map(|obj| obj.version())
797            .unwrap_or_default();
798        self.try_get_past_object(
799            object_id,
800            version,
801            Some(SuiObjectDataOptions::bcs_lossless()),
802        )
803        .await
804    }
805
806    #[instrument(skip(self))]
807    async fn try_multi_get_past_objects(
808        &self,
809        past_objects: Vec<SuiGetPastObjectRequest>,
810        options: Option<SuiObjectDataOptions>,
811    ) -> RpcResult<Vec<SuiPastObjectResponse>> {
812        with_tracing!(async move {
813            if past_objects.len() <= *QUERY_MAX_RESULT_LIMIT {
814                let mut futures = vec![];
815                for past_object in past_objects {
816                    futures.push(self.try_get_past_object(
817                        past_object.object_id,
818                        past_object.version,
819                        options.clone(),
820                    ));
821                }
822                let results = join_all(futures).await;
823
824                let (oks, errs): (Vec<_>, Vec<_>) = results.into_iter().partition(Result::is_ok);
825                let success = oks.into_iter().filter_map(Result::ok).collect();
826                let errors: Vec<_> = errs.into_iter().filter_map(Result::err).collect();
827                if !errors.is_empty() {
828                    let error_string = errors
829                        .iter()
830                        .map(|e| e.to_string())
831                        .collect::<Vec<String>>()
832                        .join("; ");
833                    Err(anyhow!("{error_string}").into()) // Collects errors not related to SuiPastObjectResponse variants
834                } else {
835                    Ok(success)
836                }
837            } else {
838                Err(SuiRpcInputError::SizeLimitExceeded(
839                    QUERY_MAX_RESULT_LIMIT.to_string(),
840                ))?
841            }
842        })
843    }
844
845    #[instrument(skip(self))]
846    async fn get_total_transaction_blocks(&self) -> RpcResult<BigInt<u64>> {
847        with_tracing!(async move {
848            Ok(self
849                .state
850                .get_total_transaction_blocks()
851                .map_err(Error::from)?
852                .into()) // converts into BigInt<u64>
853        })
854    }
855
856    #[instrument(skip(self))]
857    async fn get_transaction_block(
858        &self,
859        digest: TransactionDigest,
860        opts: Option<SuiTransactionBlockResponseOptions>,
861    ) -> RpcResult<SuiTransactionBlockResponse> {
862        with_tracing!(async move {
863            let opts = opts.unwrap_or_default();
864            let mut temp_response = IntermediateTransactionResponse::new(digest);
865
866            // Fetch transaction to determine existence
867            let transaction_kv_store = self.transaction_kv_store.clone();
868            let transaction = async move {
869                let ret = transaction_kv_store.get_tx(digest).await.map_err(|err| {
870                    debug!(tx_digest=?digest, "Failed to get transaction: {}", err);
871                    Error::from(err)
872                });
873                add_server_timing("tx_kv_lookup");
874                ret
875            }
876            .await?;
877            let input_objects = transaction
878                .data()
879                .inner()
880                .intent_message
881                .value
882                .input_objects()
883                .unwrap_or_default();
884
885            // the input is needed for object_changes to retrieve the sender address.
886            if opts.require_input() {
887                temp_response.transaction = Some(transaction);
888            }
889
890            // Fetch effects when `show_events` is true because events relies on effects
891            if opts.require_effects() {
892                let transaction_kv_store = self.transaction_kv_store.clone();
893                temp_response.effects = Some(
894                    transaction_kv_store
895                        .get_fx_by_tx_digest(digest)
896                        .await
897                        .map_err(|err| {
898                            debug!(tx_digest=?digest, "Failed to get effects: {:?}", err);
899                            Error::from(err)
900                        })?,
901                );
902            }
903
904            temp_response.checkpoint_seq = self
905                .transaction_kv_store
906                .deprecated_get_transaction_checkpoint(digest)
907                .await
908                .map_err(|e| {
909                    error!("Failed to retrieve checkpoint sequence for transaction {digest:?} with error: {e:?}");
910                    Error::from(e)
911                })?;
912
913            if let Some(checkpoint_seq) = &temp_response.checkpoint_seq {
914                let kv_store = self.transaction_kv_store.clone();
915                let checkpoint_seq = *checkpoint_seq;
916                let checkpoint = kv_store
917                    // safe to unwrap because we have checked `is_some` above
918                    .get_checkpoint_summary(checkpoint_seq)
919                    .await
920                    .map_err(|e| {
921                        error!("Failed to get checkpoint by sequence number: {checkpoint_seq:?} with error: {e:?}");
922                        Error::from(e)
923                    })?;
924                // TODO(chris): we don't need to fetch the whole checkpoint summary
925                temp_response.timestamp = Some(checkpoint.timestamp_ms);
926            }
927
928            if opts.show_events && temp_response.effects.is_some() {
929                let transaction_kv_store = self.transaction_kv_store.clone();
930                let events = transaction_kv_store
931                    .multi_get_events_by_tx_digests(&[digest])
932                    .await
933                    .map_err(|e| {
934                        error!("Failed to call get transaction events for transaction: {digest:?} with error {e:?}");
935                        Error::from(e)
936                    })?
937                    .pop()
938                    .flatten();
939                match events {
940                    None => temp_response.events = Some(SuiTransactionBlockEvents::default()),
941                    Some(events) => match to_sui_transaction_events(self, digest, events) {
942                        Ok(e) => temp_response.events = Some(e),
943                        Err(e) => temp_response.errors.push(e.to_string()),
944                    },
945                }
946            }
947
948            let object_cache =
949                ObjectProviderCache::new((self.state.clone(), self.transaction_kv_store.clone()));
950            if opts.show_balance_changes
951                && let Some(effects) = &temp_response.effects
952            {
953                let balance_changes =
954                    get_balance_changes_from_effect(&object_cache, effects, input_objects, None)
955                        .await;
956
957                if let Ok(balance_changes) = balance_changes {
958                    temp_response.balance_changes = Some(balance_changes);
959                } else {
960                    temp_response.errors.push(format!(
961                        "Cannot retrieve balance changes: {}",
962                        balance_changes.unwrap_err()
963                    ));
964                }
965            }
966
967            if opts.show_object_changes
968                && let (Some(effects), Some(input)) =
969                    (&temp_response.effects, &temp_response.transaction)
970            {
971                let sender = input.data().intent_message().value.sender();
972                let object_changes = get_object_changes(
973                    &object_cache,
974                    effects,
975                    sender,
976                    effects.modified_at_versions(),
977                    effects.all_changed_objects(),
978                    effects.all_removed_objects(),
979                )
980                .await;
981
982                if let Ok(object_changes) = object_changes {
983                    temp_response.object_changes = Some(object_changes);
984                } else {
985                    temp_response.errors.push(format!(
986                        "Cannot retrieve object changes: {}",
987                        object_changes.unwrap_err()
988                    ));
989                }
990            }
991            let epoch_store = self.state.load_epoch_store_one_call_per_task();
992            convert_to_response(temp_response, &opts, epoch_store.module_cache())
993        })
994    }
995
996    #[instrument(skip(self))]
997    async fn multi_get_transaction_blocks(
998        &self,
999        digests: Vec<TransactionDigest>,
1000        opts: Option<SuiTransactionBlockResponseOptions>,
1001    ) -> RpcResult<Vec<SuiTransactionBlockResponse>> {
1002        with_tracing!(async move {
1003            self.multi_get_transaction_blocks_internal(digests, opts)
1004                .await
1005        })
1006    }
1007
1008    #[instrument(skip(self))]
1009    async fn get_events(&self, transaction_digest: TransactionDigest) -> RpcResult<Vec<SuiEvent>> {
1010        with_tracing!(async move {
1011            let state = self.state.clone();
1012            let transaction_kv_store = self.transaction_kv_store.clone();
1013            async move {
1014                let store = state.load_epoch_store_one_call_per_task();
1015                let events = transaction_kv_store
1016                    .multi_get_events_by_tx_digests(&[transaction_digest])
1017                    .await
1018                    .map_err(|e| {
1019                        error!("Failed to get transaction events for transaction {transaction_digest:?} with error: {e:?}");
1020                        Error::StateReadError(e.into())
1021                    })?
1022                    .pop()
1023                    .flatten();
1024                Ok(match events {
1025                    Some(events) => events
1026                        .data
1027                        .into_iter()
1028                        .enumerate()
1029                        .map(|(seq, e)| {
1030                            let layout = store
1031                                .executor()
1032                                .type_layout_resolver(
1033                                    store.protocol_config(),
1034                                    Box::new(
1035                                        &state.get_backing_package_store().as_ref(),
1036                                    ))
1037                                .get_annotated_layout(&e.type_)?;
1038                            SuiEvent::try_from(e, transaction_digest, seq as u64, None, layout)
1039                        })
1040                        .collect::<Result<Vec<_>, _>>()
1041                        .map_err(Error::SuiError)?,
1042                    None => vec![],
1043                })
1044            }
1045            .await
1046        })
1047    }
1048
1049    #[instrument(skip(self))]
1050    async fn get_latest_checkpoint_sequence_number(&self) -> RpcResult<BigInt<u64>> {
1051        with_tracing!(async move {
1052            Ok(self
1053                .state
1054                .get_latest_checkpoint_sequence_number()
1055                .map_err(|e| {
1056                    SuiRpcInputError::GenericNotFound(format!(
1057                        "Latest checkpoint sequence number was not found with error :{e}"
1058                    ))
1059                })?
1060                .into())
1061        })
1062    }
1063
1064    #[instrument(skip(self))]
1065    async fn get_checkpoint(&self, id: CheckpointId) -> RpcResult<Checkpoint> {
1066        with_tracing!(self.get_checkpoint_internal(id))
1067    }
1068
1069    #[instrument(skip(self))]
1070    async fn get_checkpoints(
1071        &self,
1072        // If `Some`, the query will start from the next item after the specified cursor
1073        cursor: Option<BigInt<u64>>,
1074        limit: Option<usize>,
1075        descending_order: bool,
1076    ) -> RpcResult<CheckpointPage> {
1077        with_tracing!(async move {
1078            let limit = validate_limit(limit, QUERY_MAX_RESULT_LIMIT_CHECKPOINTS)
1079                .map_err(SuiRpcInputError::from)?;
1080
1081            let state = self.state.clone();
1082            let kv_store = self.transaction_kv_store.clone();
1083
1084            self.metrics.get_checkpoints_limit.observe(limit as f64);
1085
1086            let mut data = Self::get_checkpoints_internal(
1087                state,
1088                kv_store,
1089                cursor.map(|s| *s),
1090                limit as u64 + 1,
1091                descending_order,
1092            )
1093            .await
1094            .map_err(Error::from)?;
1095
1096            let has_next_page = data.len() > limit;
1097            data.truncate(limit);
1098
1099            let next_cursor = if has_next_page {
1100                data.last().cloned().map(|d| d.sequence_number.into())
1101            } else {
1102                None
1103            };
1104
1105            self.metrics
1106                .get_checkpoints_result_size
1107                .observe(data.len() as f64);
1108            self.metrics
1109                .get_checkpoints_result_size_total
1110                .inc_by(data.len() as u64);
1111
1112            Ok(CheckpointPage {
1113                data,
1114                next_cursor,
1115                has_next_page,
1116            })
1117        })
1118    }
1119
1120    #[instrument(skip(self))]
1121    async fn get_protocol_config(
1122        &self,
1123        version: Option<BigInt<u64>>,
1124    ) -> RpcResult<ProtocolConfigResponse> {
1125        with_tracing!(async move {
1126            version
1127                .map(|v| {
1128                    ProtocolConfig::get_for_version_if_supported(
1129                        (*v).into(),
1130                        self.state.get_chain_identifier()?.chain(),
1131                    )
1132                    .ok_or(SuiRpcInputError::ProtocolVersionUnsupported(
1133                        ProtocolVersion::MIN.as_u64(),
1134                        ProtocolVersion::MAX.as_u64(),
1135                    ))
1136                    .map_err(Error::from)
1137                })
1138                .unwrap_or(Ok(self
1139                    .state
1140                    .load_epoch_store_one_call_per_task()
1141                    .protocol_config()
1142                    .clone()))
1143                .map(ProtocolConfigResponse::from)
1144        })
1145    }
1146
1147    #[instrument(skip(self))]
1148    async fn get_chain_identifier(&self) -> RpcResult<String> {
1149        with_tracing!(async move {
1150            let ci = self.state.get_chain_identifier()?;
1151            Ok(ci.to_string())
1152        })
1153    }
1154    #[instrument(skip(self))]
1155    async fn verify_zklogin_signature(
1156        &self,
1157        bytes: String,
1158        signature: String,
1159        intent_scope: ZkLoginIntentScope,
1160        author: SuiAddress,
1161    ) -> RpcResult<ZkLoginVerifyResult> {
1162        let epoch_store = self.state.load_epoch_store_one_call_per_task();
1163        let curr_epoch = epoch_store.epoch();
1164        let zklogin_env_native = match self
1165            .state
1166            .get_chain_identifier()
1167            .expect("get chain identifier should not fail")
1168            .chain()
1169        {
1170            sui_protocol_config::Chain::Mainnet | sui_protocol_config::Chain::Testnet => {
1171                ZkLoginEnv::Prod
1172            }
1173            _ => ZkLoginEnv::Test,
1174        };
1175        let GenericSignature::ZkLoginAuthenticator(zklogin_sig) =
1176            GenericSignature::from_bytes(&Base64::decode(&signature).map_err(Error::from)?)
1177                .map_err(Error::from)?
1178        else {
1179            return Err(SuiRpcInputError::GenericNotFound(
1180                "Endpoint only supports zkLogin signature".to_string(),
1181            )
1182            .into());
1183        };
1184
1185        let new_jwks =
1186            match get_authenticator_state(self.state.get_object_store()).map_err(Error::from)? {
1187                Some(authenticator_state) => authenticator_state.active_jwks,
1188                None => {
1189                    return Err(SuiRpcInputError::GenericNotFound(
1190                        "Authenticator state not found".to_string(),
1191                    )
1192                    .into());
1193                }
1194            };
1195
1196        // construct verify params with active jwks and zklogin_env.
1197        let mut oidc_provider_jwks = ImHashMap::new();
1198        for active_jwk in new_jwks.iter() {
1199            let ActiveJwk { jwk_id, jwk, .. } = active_jwk;
1200            match oidc_provider_jwks.entry(jwk_id.clone()) {
1201                imbl::hashmap::Entry::Occupied(_) => {
1202                    warn!("JWK with kid {:?} already exists", jwk_id);
1203                }
1204                imbl::hashmap::Entry::Vacant(entry) => {
1205                    entry.insert(jwk.clone());
1206                }
1207            }
1208        }
1209        let verify_params = VerifyParams::new(
1210            oidc_provider_jwks,
1211            vec![],
1212            zklogin_env_native,
1213            epoch_store.protocol_config().zklogin_circuit_mode(),
1214            true,
1215            true,
1216            true,
1217            Some(30),
1218            true,
1219            true,
1220        );
1221        match intent_scope {
1222            ZkLoginIntentScope::TransactionData => {
1223                let tx_data: TransactionData =
1224                    bcs::from_bytes(&Base64::decode(&bytes).map_err(Error::from)?)
1225                        .map_err(Error::from)?;
1226                let intent_msg = IntentMessage::new(Intent::sui_transaction(), tx_data.clone());
1227                let sig = GenericSignature::ZkLoginAuthenticator(zklogin_sig);
1228                match sig.verify_authenticator(
1229                    &intent_msg,
1230                    author,
1231                    curr_epoch,
1232                    &verify_params,
1233                    Arc::new(VerifiedDigestCache::new_empty()),
1234                ) {
1235                    Ok(_) => Ok(ZkLoginVerifyResult {
1236                        success: true,
1237                        errors: vec![],
1238                    }),
1239                    Err(e) => Ok(ZkLoginVerifyResult {
1240                        success: false,
1241                        errors: vec![e.to_string()],
1242                    }),
1243                }
1244            }
1245            ZkLoginIntentScope::PersonalMessage => {
1246                let data = PersonalMessage {
1247                    message: Base64::decode(&bytes).map_err(Error::from)?,
1248                };
1249                let intent_msg = IntentMessage::new(Intent::personal_message(), data);
1250
1251                let sig = GenericSignature::ZkLoginAuthenticator(zklogin_sig);
1252                match sig.verify_authenticator(
1253                    &intent_msg,
1254                    author,
1255                    curr_epoch,
1256                    &verify_params,
1257                    Arc::new(VerifiedDigestCache::new_empty()),
1258                ) {
1259                    Ok(_) => Ok(ZkLoginVerifyResult {
1260                        success: true,
1261                        errors: vec![],
1262                    }),
1263                    Err(e) => Ok(ZkLoginVerifyResult {
1264                        success: false,
1265                        errors: vec![e.to_string()],
1266                    }),
1267                }
1268            }
1269        }
1270    }
1271}
1272
1273impl SuiRpcModule for ReadApi {
1274    fn rpc(self) -> RpcModule<Self> {
1275        self.into_rpc()
1276    }
1277
1278    fn rpc_doc_module() -> Module {
1279        ReadApiOpenRpc::module_doc()
1280    }
1281}
1282
1283#[instrument(skip_all)]
1284fn to_sui_transaction_events(
1285    fullnode_api: &ReadApi,
1286    tx_digest: TransactionDigest,
1287    events: TransactionEvents,
1288) -> Result<SuiTransactionBlockEvents, Error> {
1289    let epoch_store = fullnode_api.state.load_epoch_store_one_call_per_task();
1290    let backing_package_store = fullnode_api.state.get_backing_package_store();
1291    let mut layout_resolver = epoch_store.executor().type_layout_resolver(
1292        epoch_store.protocol_config(),
1293        Box::new(backing_package_store.as_ref()),
1294    );
1295    Ok(SuiTransactionBlockEvents::try_from(
1296        events,
1297        tx_digest,
1298        None,
1299        layout_resolver.as_mut(),
1300    )?)
1301}
1302
1303#[derive(Debug, thiserror::Error)]
1304pub enum ObjectDisplayError {
1305    #[error("Not a move struct")]
1306    NotMoveStruct,
1307
1308    #[error("Failed to extract layout")]
1309    Layout,
1310
1311    #[error("Failed to extract Move object")]
1312    MoveObject,
1313
1314    #[error(transparent)]
1315    Deserialization(#[from] SuiError),
1316
1317    #[error("Failed to deserialize 'VersionUpdatedEvent': {0}")]
1318    Bcs(#[from] bcs::Error),
1319
1320    #[error(transparent)]
1321    StateReadError(#[from] StateReadError),
1322
1323    #[error(transparent)]
1324    Internal(#[from] anyhow::Error),
1325}
1326
1327#[instrument(skip(fullnode_api, kv_store))]
1328async fn get_display_fields(
1329    fullnode_api: &ReadApi,
1330    kv_store: &Arc<TransactionKeyValueStore>,
1331    original_object: &Object,
1332    original_layout: &Option<MoveStructLayout>,
1333) -> Result<DisplayFieldsResponse, ObjectDisplayError> {
1334    let (layout, type_) = if let Some(layout) = original_layout {
1335        let type_ = &layout.type_;
1336        let layout = MoveTypeLayout::Struct(Box::new(layout.clone()));
1337        (layout, type_)
1338    } else {
1339        return Ok(DisplayFieldsResponse {
1340            data: None,
1341            error: None,
1342        });
1343    };
1344
1345    let Some(move_object) = original_object.data.try_as_move() else {
1346        return Err(ObjectDisplayError::MoveObject);
1347    };
1348
1349    let display: Vec<(String, Result<Json, anyhow::Error>)> =
1350        if let Some(display_object) = get_display_object_v2_by_type(fullnode_api, type_)? {
1351            let root = sui_display::v2::OwnedSlice::new(layout, move_object.contents().to_owned());
1352            let store = DisplayStore::new(fullnode_api.state.as_ref());
1353            let interpreter = sui_display::v2::Interpreter::new(root, store);
1354            let limits = sui_display::v2::Limits {
1355                max_depth: *MAX_DISPLAY_FIELD_DEPTH,
1356                max_nodes: *MAX_DISPLAY_FORMAT_NODES,
1357                max_loads: *MAX_DISPLAY_OBJECT_LOADS,
1358            };
1359
1360            match sui_display::v2::Display::parse(limits, display_object.fields()) {
1361                Ok(display) => match display
1362                    .display::<Json>(
1363                        *MAX_DISPLAY_MOVE_VALUE_DEPTH,
1364                        *MAX_DISPLAY_OUTPUT_SIZE,
1365                        &interpreter,
1366                    )
1367                    .await
1368                {
1369                    Ok(fields) => fields
1370                        .into_iter()
1371                        .map(|(field, value)| (field, value.map_err(anyhow::Error::from)))
1372                        .collect(),
1373                    Err(e) => {
1374                        return Ok(DisplayFieldsResponse {
1375                            data: None,
1376                            error: Some(SuiObjectResponseError::DisplayError {
1377                                error: e.to_string(),
1378                            }),
1379                        });
1380                    }
1381                },
1382
1383                Err(e) => {
1384                    return Ok(DisplayFieldsResponse {
1385                        data: None,
1386                        error: Some(SuiObjectResponseError::DisplayError {
1387                            error: e.to_string(),
1388                        }),
1389                    });
1390                }
1391            }
1392        } else if let Some(display_object) =
1393            get_display_object_v1_by_type(kv_store, fullnode_api, type_).await?
1394        {
1395            let format = match Format::parse(*MAX_DISPLAY_FIELD_DEPTH, &display_object.fields) {
1396                Ok(format) => format,
1397                Err(e) => {
1398                    return Ok(DisplayFieldsResponse {
1399                        data: None,
1400                        error: Some(SuiObjectResponseError::DisplayError {
1401                            error: e.to_string(),
1402                        }),
1403                    });
1404                }
1405            };
1406
1407            match format.display(*MAX_DISPLAY_OUTPUT_SIZE, move_object.contents(), &layout) {
1408                Ok(fields) => fields
1409                    .into_iter()
1410                    .map(|(field, value)| (field, value.map(Json::String)))
1411                    .collect(),
1412                Err(e) => {
1413                    return Ok(DisplayFieldsResponse {
1414                        data: None,
1415                        error: Some(SuiObjectResponseError::DisplayError {
1416                            error: e.to_string(),
1417                        }),
1418                    });
1419                }
1420            }
1421        } else {
1422            return Ok(DisplayFieldsResponse {
1423                data: None,
1424                error: None,
1425            });
1426        };
1427
1428    let mut fields = BTreeMap::new();
1429    let mut errors = vec![];
1430
1431    for (key, value) in display {
1432        match value {
1433            Ok(v) => {
1434                fields.insert(key, v);
1435            }
1436            Err(e) => {
1437                errors.push(e.to_string());
1438            }
1439        }
1440    }
1441
1442    Ok(DisplayFieldsResponse {
1443        data: (!fields.is_empty()).then_some(fields),
1444        error: (!errors.is_empty()).then(|| SuiObjectResponseError::DisplayError {
1445            error: errors.join("; "),
1446        }),
1447    })
1448}
1449
1450#[instrument(skip(kv_store, fullnode_api))]
1451async fn get_display_object_v1_by_type(
1452    kv_store: &Arc<TransactionKeyValueStore>,
1453    fullnode_api: &ReadApi,
1454    object_type: &StructTag,
1455) -> Result<Option<DisplayVersionUpdatedEvent>, ObjectDisplayError> {
1456    let mut events = fullnode_api
1457        .state
1458        .query_events(
1459            kv_store,
1460            EventFilter::MoveEventType(DisplayVersionUpdatedEvent::type_(object_type)),
1461            None,
1462            1,
1463            true,
1464        )
1465        .await?;
1466
1467    // If there's any recent version of Display, give it to the client.
1468    if let Some(event) = events.pop() {
1469        let display: DisplayVersionUpdatedEvent = bcs::from_bytes(&event.bcs.into_bytes())?;
1470        Ok(Some(display))
1471    } else {
1472        Ok(None)
1473    }
1474}
1475
1476#[instrument(skip(fullnode_api))]
1477fn get_display_object_v2_by_type(
1478    fullnode_api: &ReadApi,
1479    object_type: &StructTag,
1480) -> Result<Option<display_registry::Display>, ObjectDisplayError> {
1481    let object_id = display_registry::display_object_id(object_type.clone().into())?;
1482    let ObjectRead::Exists(_, object, _) = fullnode_api.state.get_object_read(&object_id)? else {
1483        return Ok(None);
1484    };
1485
1486    let Some(move_object) = object.data.try_as_move() else {
1487        return Ok(None);
1488    };
1489
1490    Ok(Some(bcs::from_bytes(move_object.contents())?))
1491}
1492
1493#[instrument(skip_all)]
1494fn convert_to_response(
1495    cache: IntermediateTransactionResponse,
1496    opts: &SuiTransactionBlockResponseOptions,
1497    module_cache: &impl GetModule,
1498) -> RpcInterimResult<SuiTransactionBlockResponse> {
1499    let mut response = SuiTransactionBlockResponse::new(cache.digest);
1500    response.errors = cache.errors;
1501
1502    if let Some(transaction) = cache.transaction {
1503        if opts.show_raw_input {
1504            response.raw_transaction = bcs::to_bytes(transaction.data()).map_err(|e| {
1505                // TODO: is this a client or server error?
1506                anyhow!("Failed to serialize raw transaction with error: {e}")
1507            })?;
1508        }
1509
1510        if opts.show_input {
1511            response.transaction = Some(SuiTransactionBlock::try_from(
1512                transaction.into_data(),
1513                module_cache,
1514            )?);
1515        }
1516    }
1517
1518    if let Some(effects) = cache.effects {
1519        if opts.show_raw_effects {
1520            response.raw_effects = bcs::to_bytes(&effects).map_err(|e| {
1521                // TODO: is this a client or server error?
1522                anyhow!("Failed to serialize transaction block effects with error: {e}")
1523            })?;
1524        }
1525
1526        if opts.show_effects {
1527            response.effects = Some(effects.try_into().map_err(|e| {
1528                // TODO: is this a client or server error?
1529                anyhow!("Failed to convert transaction block effects with error: {e}")
1530            })?);
1531        }
1532    }
1533
1534    response.checkpoint = cache.checkpoint_seq;
1535    response.timestamp_ms = cache.timestamp;
1536
1537    if opts.show_events {
1538        response.events = cache.events;
1539    }
1540
1541    if opts.show_balance_changes {
1542        response.balance_changes = cache.balance_changes;
1543    }
1544
1545    if opts.show_object_changes {
1546        response.object_changes = cache.object_changes;
1547    }
1548
1549    Ok(response)
1550}
1551
1552fn calculate_checkpoint_numbers(
1553    // If `Some`, the query will start from the next item after the specified cursor
1554    cursor: Option<CheckpointSequenceNumber>,
1555    limit: u64,
1556    descending_order: bool,
1557    max_checkpoint: CheckpointSequenceNumber,
1558) -> Vec<CheckpointSequenceNumber> {
1559    let (start_index, end_index) = match cursor {
1560        Some(t) => {
1561            if descending_order {
1562                let start = std::cmp::min(t.saturating_sub(1), max_checkpoint);
1563                let end = start.saturating_sub(limit - 1);
1564                (end, start)
1565            } else {
1566                let start =
1567                    std::cmp::min(t.checked_add(1).unwrap_or(max_checkpoint), max_checkpoint);
1568                let end = std::cmp::min(
1569                    start.checked_add(limit - 1).unwrap_or(max_checkpoint),
1570                    max_checkpoint,
1571                );
1572                (start, end)
1573            }
1574        }
1575        None => {
1576            if descending_order {
1577                (max_checkpoint.saturating_sub(limit - 1), max_checkpoint)
1578            } else {
1579                (0, std::cmp::min(limit - 1, max_checkpoint))
1580            }
1581        }
1582    };
1583
1584    if descending_order {
1585        (start_index..=end_index).rev().collect()
1586    } else {
1587        (start_index..=end_index).collect()
1588    }
1589}
1590
1591#[cfg(test)]
1592mod tests {
1593    use super::*;
1594    use crate::authority_state::MockStateRead;
1595    use mockall::mock;
1596    use roaring::RoaringBitmap;
1597    use std::collections::HashMap;
1598    use sui_storage::key_value_store::{
1599        KVStoreCheckpointData, KVStoreTransactionData, TransactionKeyValueStoreTrait,
1600    };
1601    use sui_storage::key_value_store_metrics::KeyValueStoreMetrics;
1602    use sui_types::base_types::ExecutionDigests;
1603    use sui_types::crypto::AuthorityStrongQuorumSignInfo;
1604    use sui_types::digests::TransactionEffectsDigest;
1605    use sui_types::effects::TransactionEvents;
1606    use sui_types::error::SuiResult;
1607    use sui_types::gas::GasCostSummary;
1608    use sui_types::message_envelope::Envelope;
1609    use sui_types::messages_checkpoint::{
1610        CertifiedCheckpointSummary, CheckpointContents, CheckpointDigest, CheckpointSummary,
1611    };
1612    use sui_types::object::Object;
1613    use sui_types::storage::ObjectKey;
1614
1615    #[test]
1616    fn test_calculate_checkpoint_numbers() {
1617        let cursor = Some(10);
1618        let limit = 5;
1619        let descending_order = true;
1620        let max_checkpoint = 15;
1621
1622        let checkpoint_numbers =
1623            calculate_checkpoint_numbers(cursor, limit, descending_order, max_checkpoint);
1624
1625        assert_eq!(checkpoint_numbers, vec![9, 8, 7, 6, 5]);
1626    }
1627
1628    #[test]
1629    fn test_calculate_checkpoint_numbers_descending_no_cursor() {
1630        let cursor = None;
1631        let limit = 5;
1632        let descending_order = true;
1633        let max_checkpoint = 15;
1634
1635        let checkpoint_numbers =
1636            calculate_checkpoint_numbers(cursor, limit, descending_order, max_checkpoint);
1637
1638        assert_eq!(checkpoint_numbers, vec![15, 14, 13, 12, 11]);
1639    }
1640
1641    #[test]
1642    fn test_calculate_checkpoint_numbers_ascending_no_cursor() {
1643        let cursor = None;
1644        let limit = 5;
1645        let descending_order = false;
1646        let max_checkpoint = 15;
1647
1648        let checkpoint_numbers =
1649            calculate_checkpoint_numbers(cursor, limit, descending_order, max_checkpoint);
1650
1651        assert_eq!(checkpoint_numbers, vec![0, 1, 2, 3, 4]);
1652    }
1653
1654    #[test]
1655    fn test_calculate_checkpoint_numbers_ascending_with_cursor() {
1656        let cursor = Some(10);
1657        let limit = 5;
1658        let descending_order = false;
1659        let max_checkpoint = 15;
1660
1661        let checkpoint_numbers =
1662            calculate_checkpoint_numbers(cursor, limit, descending_order, max_checkpoint);
1663
1664        assert_eq!(checkpoint_numbers, vec![11, 12, 13, 14, 15]);
1665    }
1666
1667    #[test]
1668    fn test_calculate_checkpoint_numbers_ascending_limit_exceeds_max() {
1669        let cursor = None;
1670        let limit = 20;
1671        let descending_order = false;
1672        let max_checkpoint = 15;
1673
1674        let checkpoint_numbers =
1675            calculate_checkpoint_numbers(cursor, limit, descending_order, max_checkpoint);
1676
1677        assert_eq!(checkpoint_numbers, (0..=15).collect::<Vec<_>>());
1678    }
1679
1680    #[test]
1681    fn test_calculate_checkpoint_numbers_descending_limit_exceeds_max() {
1682        let cursor = None;
1683        let limit = 20;
1684        let descending_order = true;
1685        let max_checkpoint = 15;
1686
1687        let checkpoint_numbers =
1688            calculate_checkpoint_numbers(cursor, limit, descending_order, max_checkpoint);
1689
1690        assert_eq!(checkpoint_numbers, (0..=15).rev().collect::<Vec<_>>());
1691    }
1692
1693    mock! {
1694        CheckpointKvStore {}
1695        #[async_trait]
1696        impl TransactionKeyValueStoreTrait for CheckpointKvStore {
1697            async fn multi_get(
1698                &self,
1699                transactions: &[TransactionDigest],
1700                effects: &[TransactionDigest],
1701            ) -> SuiResult<KVStoreTransactionData>;
1702
1703            async fn multi_get_checkpoints(
1704                &self,
1705                checkpoint_summaries: &[CheckpointSequenceNumber],
1706                checkpoint_contents: &[CheckpointSequenceNumber],
1707                checkpoint_summaries_by_digest: &[CheckpointDigest],
1708            ) -> SuiResult<KVStoreCheckpointData>;
1709
1710            async fn deprecated_get_transaction_checkpoint(
1711                &self,
1712                digest: TransactionDigest,
1713            ) -> SuiResult<Option<CheckpointSequenceNumber>>;
1714
1715            async fn get_object(
1716                &self,
1717                object_id: ObjectID,
1718                version: SequenceNumber,
1719            ) -> SuiResult<Option<Object>>;
1720
1721            async fn multi_get_objects(
1722                &self,
1723                object_keys: &[ObjectKey],
1724            ) -> SuiResult<Vec<Option<Object>>>;
1725
1726            async fn multi_get_transaction_checkpoint(
1727                &self,
1728                digests: &[TransactionDigest],
1729            ) -> SuiResult<Vec<Option<CheckpointSequenceNumber>>>;
1730
1731            async fn multi_get_events_by_tx_digests(
1732                &self,
1733                digests: &[TransactionDigest],
1734            ) -> SuiResult<Vec<Option<TransactionEvents>>>;
1735        }
1736    }
1737
1738    // Builds `CheckpointContents` whose single transaction digest uniquely
1739    // encodes `seq`, so a returned `Checkpoint` can be traced back to the
1740    // sequence number whose contents it actually carries.
1741    fn test_checkpoint_contents(seq: CheckpointSequenceNumber) -> CheckpointContents {
1742        let mut tx = [0u8; 32];
1743        tx[0] = 0xA;
1744        tx[1..9].copy_from_slice(&seq.to_le_bytes());
1745        let mut fx = [0u8; 32];
1746        fx[0] = 0xE;
1747        fx[1..9].copy_from_slice(&seq.to_le_bytes());
1748        CheckpointContents::new_with_digests_only_for_tests([ExecutionDigests::new(
1749            TransactionDigest::new(tx),
1750            TransactionEffectsDigest::new(fx),
1751        )])
1752    }
1753
1754    // Builds a certified summary for `seq`. The aggregate signature is a
1755    // placeholder; `get_checkpoints_internal` never verifies it.
1756    fn test_certified_summary(
1757        seq: CheckpointSequenceNumber,
1758        contents: &CheckpointContents,
1759    ) -> CertifiedCheckpointSummary {
1760        let summary = CheckpointSummary::new(
1761            &ProtocolConfig::get_for_max_version_UNSAFE(),
1762            0,
1763            seq,
1764            seq,
1765            contents,
1766            None,
1767            GasCostSummary::default(),
1768            None,
1769            0,
1770            Vec::new(),
1771            Vec::new(),
1772        );
1773        let auth_sig = AuthorityStrongQuorumSignInfo {
1774            epoch: 0,
1775            signature: Default::default(),
1776            signers_map: RoaringBitmap::new(),
1777        };
1778        Envelope::new_from_data_and_sig(summary, auth_sig)
1779    }
1780
1781    // Regression test for a pruning-induced misalignment: when a checkpoint's
1782    // contents are pruned but its sequence-addressable summary survives,
1783    // `get_checkpoints_internal` must not pair that summary (or any later one)
1784    // with a different checkpoint's contents.
1785    #[tokio::test]
1786    async fn test_get_checkpoints_internal_preserves_alignment_across_pruned_contents() {
1787        let max_checkpoint: CheckpointSequenceNumber = 13;
1788        // Contents for this sequence are pruned while its summary remains. It
1789        // sits in the interior of the requested range to show that alignment is
1790        // preserved regardless of where the hole falls.
1791        let pruned_seq: CheckpointSequenceNumber = 11;
1792
1793        let mut all_contents: HashMap<CheckpointSequenceNumber, CheckpointContents> =
1794            HashMap::new();
1795        for seq in 0..=max_checkpoint {
1796            all_contents.insert(seq, test_checkpoint_contents(seq));
1797        }
1798
1799        let mut mock_state = MockStateRead::new();
1800        mock_state
1801            .expect_get_latest_checkpoint_sequence_number()
1802            .returning(move || Ok(max_checkpoint));
1803
1804        let store_contents = all_contents.clone();
1805        let mut mock_kv = MockCheckpointKvStore::new();
1806        mock_kv.expect_multi_get_checkpoints().times(2).returning(
1807            move |summaries, contents, _by_digest| {
1808                // Summaries survive pruning for every requested sequence.
1809                let summaries = summaries
1810                    .iter()
1811                    .map(|seq| Some(test_certified_summary(*seq, &store_contents[seq])))
1812                    .collect();
1813                // Contents are missing for the pruned sequence only.
1814                let contents = contents
1815                    .iter()
1816                    .map(|seq| (*seq != pruned_seq).then(|| store_contents[seq].clone()))
1817                    .collect();
1818                Ok((summaries, contents, vec![]))
1819            },
1820        );
1821
1822        let state: Arc<dyn StateRead> = Arc::new(mock_state);
1823        let kv_store = Arc::new(TransactionKeyValueStore::new(
1824            "test",
1825            KeyValueStoreMetrics::new_for_tests(),
1826            Arc::new(mock_kv),
1827        ));
1828
1829        // cursor = 9, ascending, limit 4 => requested sequences [10, 11, 12, 13].
1830        let checkpoints = ReadApi::get_checkpoints_internal(state, kv_store, Some(9), 4, false)
1831            .await
1832            .unwrap();
1833
1834        // The pruned sequence is omitted; every other sequence is returned once.
1835        assert_eq!(
1836            checkpoints
1837                .iter()
1838                .map(|c| c.sequence_number)
1839                .collect::<Vec<_>>(),
1840            vec![10, 12, 13],
1841        );
1842
1843        // Crucially, each returned checkpoint still carries the transactions of
1844        // its own sequence number rather than a neighbor's.
1845        for checkpoint in &checkpoints {
1846            let expected: Vec<TransactionDigest> = all_contents[&checkpoint.sequence_number]
1847                .iter()
1848                .map(|digests| digests.transaction)
1849                .collect();
1850            assert_eq!(
1851                checkpoint.transactions, expected,
1852                "checkpoint {} was paired with another checkpoint's contents",
1853                checkpoint.sequence_number,
1854            );
1855        }
1856    }
1857}