1use 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
75const DEFAULT_MAX_DISPLAY_MOVE_VALUE_DEPTH: usize = 32;
77
78const DEFAULT_MAX_DISPLAY_OUTPUT_SIZE: usize = 1024 * 1024;
80
81static 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
95static 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
109static 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
123static 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
137static 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#[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#[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 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 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 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 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 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 .unique()
393 .collect::<Vec<CheckpointSequenceNumber>>();
394
395 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 let checkpoint_to_timestamp = unique_checkpoint_numbers
409 .into_iter()
410 .zip_debug_eq(timestamps)
411 .collect::<HashMap<_, _>>();
412
413 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 .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 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 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 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 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 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 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 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()) } 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()) })
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 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 if opts.require_input() {
887 temp_response.transaction = Some(transaction);
888 }
889
890 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 .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 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 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 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 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 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 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 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 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 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 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 #[tokio::test]
1786 async fn test_get_checkpoints_internal_preserves_alignment_across_pruned_contents() {
1787 let max_checkpoint: CheckpointSequenceNumber = 13;
1788 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 let summaries = summaries
1810 .iter()
1811 .map(|seq| Some(test_certified_summary(*seq, &store_contents[seq])))
1812 .collect();
1813 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 let checkpoints = ReadApi::get_checkpoints_internal(state, kv_store, Some(9), 4, false)
1831 .await
1832 .unwrap();
1833
1834 assert_eq!(
1836 checkpoints
1837 .iter()
1838 .map(|c| c.sequence_number)
1839 .collect::<Vec<_>>(),
1840 vec![10, 12, 13],
1841 );
1842
1843 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}