Skip to main content

sui_core/
storage.rs

1// Copyright (c) Mysten Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4use crate::authority::AuthorityState;
5use crate::checkpoints::CheckpointStore;
6use crate::epoch::committee_store::CommitteeStore;
7use crate::execution_cache::ExecutionCacheTraitPointers;
8use move_core_types::language_storage::StructTag;
9use parking_lot::Mutex;
10use std::sync::Arc;
11use sui_rpc_store::RpcStoreReader;
12use sui_types::base_types::ObjectID;
13use sui_types::base_types::SequenceNumber;
14use sui_types::base_types::SuiAddress;
15use sui_types::base_types::TransactionDigest;
16use sui_types::committee::Committee;
17use sui_types::committee::EpochId;
18use sui_types::effects::{TransactionEffects, TransactionEvents};
19use sui_types::error::{SuiErrorKind, SuiResult};
20use sui_types::full_checkpoint_content::ObjectSet;
21use sui_types::messages_checkpoint::CheckpointContentsDigest;
22use sui_types::messages_checkpoint::CheckpointDigest;
23use sui_types::messages_checkpoint::CheckpointSequenceNumber;
24use sui_types::messages_checkpoint::EndOfEpochData;
25use sui_types::messages_checkpoint::VerifiedCheckpoint;
26use sui_types::messages_checkpoint::VerifiedCheckpointContents;
27use sui_types::messages_checkpoint::VersionedFullCheckpointContents;
28use sui_types::object::Object;
29use sui_types::object::Owner;
30use sui_types::storage::BalanceInfo;
31use sui_types::storage::BalanceIterator;
32use sui_types::storage::CoinInfo;
33use sui_types::storage::DynamicFieldKey;
34use sui_types::storage::LedgerBitmapBucketIterator;
35use sui_types::storage::LedgerTxSeqDigest;
36use sui_types::storage::LedgerTxSeqDigestIterator;
37use sui_types::storage::ObjectStore;
38use sui_types::storage::OwnedObjectInfo;
39use sui_types::storage::RpcIndexes;
40use sui_types::storage::RpcStateReader;
41use sui_types::storage::RuntimeObjectResolver;
42use sui_types::storage::WriteStore;
43use sui_types::storage::error::Error as StorageError;
44use sui_types::storage::error::Result;
45use sui_types::storage::{BackingPackageStore, PackageObject};
46use sui_types::storage::{ObjectKey, OverlayBackingPackageStore, ReadStore};
47use sui_types::transaction::VerifiedTransaction;
48use tap::TapFallible;
49use tracing::error;
50use typed_store::TypedStoreError;
51
52#[derive(Clone)]
53pub struct RocksDbStore {
54    cache_traits: ExecutionCacheTraitPointers,
55
56    committee_store: Arc<CommitteeStore>,
57    checkpoint_store: Arc<CheckpointStore>,
58    // in memory checkpoint watermark sequence numbers
59    highest_verified_checkpoint: Arc<Mutex<Option<u64>>>,
60    highest_synced_checkpoint: Arc<Mutex<Option<u64>>>,
61}
62
63impl RocksDbStore {
64    pub fn new(
65        cache_traits: ExecutionCacheTraitPointers,
66        committee_store: Arc<CommitteeStore>,
67        checkpoint_store: Arc<CheckpointStore>,
68    ) -> Self {
69        Self {
70            cache_traits,
71            committee_store,
72            checkpoint_store,
73            highest_verified_checkpoint: Arc::new(Mutex::new(None)),
74            highest_synced_checkpoint: Arc::new(Mutex::new(None)),
75        }
76    }
77
78    pub fn get_objects(&self, object_keys: &[ObjectKey]) -> Vec<Option<Object>> {
79        self.cache_traits
80            .object_cache_reader
81            .multi_get_objects_by_key(object_keys)
82    }
83
84    pub fn get_last_executed_checkpoint(&self) -> Option<VerifiedCheckpoint> {
85        self.checkpoint_store
86            .get_highest_executed_checkpoint()
87            .expect("db error")
88    }
89}
90
91impl ReadStore for RocksDbStore {
92    fn get_checkpoint_by_digest(&self, digest: &CheckpointDigest) -> Option<VerifiedCheckpoint> {
93        self.checkpoint_store
94            .get_checkpoint_by_digest(digest)
95            .expect("db error")
96    }
97
98    fn get_checkpoint_by_sequence_number(
99        &self,
100        sequence_number: CheckpointSequenceNumber,
101    ) -> Option<VerifiedCheckpoint> {
102        self.checkpoint_store
103            .get_checkpoint_by_sequence_number(sequence_number)
104            .expect("db error")
105    }
106
107    fn multi_get_checkpoint_by_sequence_number(
108        &self,
109        sequence_numbers: &[CheckpointSequenceNumber],
110    ) -> Vec<Option<VerifiedCheckpoint>> {
111        self.checkpoint_store
112            .multi_get_checkpoint_by_sequence_number(sequence_numbers)
113            .expect("db error")
114    }
115
116    fn get_highest_verified_checkpoint(&self) -> Result<VerifiedCheckpoint, StorageError> {
117        self.checkpoint_store
118            .get_highest_verified_checkpoint()
119            .map(|maybe_checkpoint| {
120                maybe_checkpoint
121                    .expect("storage should have been initialized with genesis checkpoint")
122            })
123            .map_err(Into::into)
124    }
125
126    fn get_highest_synced_checkpoint(&self) -> Result<VerifiedCheckpoint, StorageError> {
127        self.checkpoint_store
128            .get_highest_synced_checkpoint()
129            .map(|maybe_checkpoint| {
130                maybe_checkpoint
131                    .expect("storage should have been initialized with genesis checkpoint")
132            })
133            .map_err(Into::into)
134    }
135
136    fn get_lowest_available_checkpoint(&self) -> Result<CheckpointSequenceNumber, StorageError> {
137        if let Some(highest_pruned_cp) = self
138            .checkpoint_store
139            .get_highest_pruned_checkpoint_seq_number()
140            .map_err(Into::<StorageError>::into)?
141        {
142            Ok(highest_pruned_cp + 1)
143        } else {
144            Ok(0)
145        }
146    }
147
148    fn get_full_checkpoint_contents(
149        &self,
150        sequence_number: Option<CheckpointSequenceNumber>,
151        digest: &CheckpointContentsDigest,
152    ) -> Option<VersionedFullCheckpointContents> {
153        #[cfg(debug_assertions)]
154        if let Some(sequence_number) = sequence_number {
155            // When sequence_number is provided as an optimization, we want to ensure that
156            // the sequence number we get from the db matches the one we provided.
157            // Only check this in debug mode though.
158            if let Some(loaded_sequence_number) = self
159                .checkpoint_store
160                .get_sequence_number_by_contents_digest(digest)
161                .expect("db error")
162            {
163                assert_eq!(loaded_sequence_number, sequence_number);
164            }
165        }
166
167        let sequence_number = sequence_number.or_else(|| {
168            self.checkpoint_store
169                .get_sequence_number_by_contents_digest(digest)
170                .expect("db error")
171        });
172        if let Some(sequence_number) = sequence_number {
173            // Note: We don't use `?` here because we want to tolerate
174            // potential db errors due to data corruption.
175            // In that case, we will fallback and construct the contents
176            // from the individual components as if we could not find the
177            // cached full contents.
178            if let Ok(Some(contents)) = self
179                .checkpoint_store
180                .get_full_checkpoint_contents_by_sequence_number(sequence_number)
181                .tap_err(|e| {
182                    error!(
183                        "error getting full checkpoint contents for checkpoint {:?}: {:?}",
184                        sequence_number, e
185                    )
186                })
187            {
188                return Some(contents);
189            }
190        }
191
192        // Otherwise gather it from the individual components.
193        // Note we can't insert the constructed contents into `full_checkpoint_content`,
194        // because it needs to be inserted along with `checkpoint_sequence_by_contents_digest`
195        // and `checkpoint_content`. However at this point it's likely we don't know the
196        // corresponding sequence number yet.
197        self.checkpoint_store
198            .get_checkpoint_contents(digest)
199            .expect("db error")
200            .and_then(|contents| {
201                let mut transactions = Vec::with_capacity(contents.size());
202                for tx in contents.iter() {
203                    if let (Some(t), Some(e)) = (
204                        self.get_transaction(&tx.transaction),
205                        self.cache_traits
206                            .transaction_cache_reader
207                            .get_effects(&tx.effects),
208                    ) {
209                        transactions.push(sui_types::base_types::ExecutionData::new(
210                            (*t).clone().into_inner(),
211                            e,
212                        ))
213                    } else {
214                        return None;
215                    }
216                }
217                Some(
218                    VersionedFullCheckpointContents::from_contents_and_execution_data(
219                        contents,
220                        transactions.into_iter(),
221                    ),
222                )
223            })
224    }
225
226    fn get_committee(&self, epoch: EpochId) -> Option<Arc<Committee>> {
227        self.committee_store.get_committee(&epoch).unwrap()
228    }
229
230    fn get_transaction(&self, digest: &TransactionDigest) -> Option<Arc<VerifiedTransaction>> {
231        self.cache_traits
232            .transaction_cache_reader
233            .get_transaction_block(digest)
234    }
235
236    fn multi_get_transactions(
237        &self,
238        digests: &[TransactionDigest],
239    ) -> Vec<Option<Arc<VerifiedTransaction>>> {
240        self.cache_traits
241            .transaction_cache_reader
242            .multi_get_transaction_blocks(digests)
243    }
244
245    fn get_transaction_effects(&self, digest: &TransactionDigest) -> Option<TransactionEffects> {
246        self.cache_traits
247            .transaction_cache_reader
248            .get_executed_effects(digest)
249    }
250
251    fn multi_get_transaction_effects(
252        &self,
253        digests: &[TransactionDigest],
254    ) -> Vec<Option<TransactionEffects>> {
255        self.cache_traits
256            .transaction_cache_reader
257            .multi_get_executed_effects(digests)
258    }
259
260    fn get_events(&self, digest: &TransactionDigest) -> Option<TransactionEvents> {
261        self.cache_traits
262            .transaction_cache_reader
263            .get_events(digest)
264    }
265
266    fn multi_get_events(&self, digests: &[TransactionDigest]) -> Vec<Option<TransactionEvents>> {
267        self.cache_traits
268            .transaction_cache_reader
269            .multi_get_events(digests)
270    }
271
272    fn get_unchanged_loaded_runtime_objects(
273        &self,
274        digest: &TransactionDigest,
275    ) -> Option<Vec<ObjectKey>> {
276        self.cache_traits
277            .transaction_cache_reader
278            .get_unchanged_loaded_runtime_objects(digest)
279    }
280
281    fn multi_get_unchanged_loaded_runtime_objects(
282        &self,
283        digests: &[TransactionDigest],
284    ) -> Vec<Option<Vec<ObjectKey>>> {
285        self.cache_traits
286            .transaction_cache_reader
287            .multi_get_unchanged_loaded_runtime_objects(digests)
288    }
289
290    fn get_transaction_checkpoint(
291        &self,
292        digest: &TransactionDigest,
293    ) -> Option<CheckpointSequenceNumber> {
294        self.cache_traits
295            .checkpoint_cache
296            .deprecated_get_transaction_checkpoint(digest)
297            .map(|(_epoch, checkpoint)| checkpoint)
298    }
299
300    fn get_latest_checkpoint(&self) -> sui_types::storage::error::Result<VerifiedCheckpoint> {
301        self.checkpoint_store
302            .get_highest_executed_checkpoint()
303            .expect("db error")
304            .ok_or_else(|| {
305                sui_types::storage::error::Error::missing("unable to get latest checkpoint")
306            })
307    }
308
309    fn get_checkpoint_contents_by_digest(
310        &self,
311        digest: &CheckpointContentsDigest,
312    ) -> Option<sui_types::messages_checkpoint::CheckpointContents> {
313        self.checkpoint_store
314            .get_checkpoint_contents(digest)
315            .expect("db error")
316    }
317
318    fn get_checkpoint_contents_by_sequence_number(
319        &self,
320        sequence_number: CheckpointSequenceNumber,
321    ) -> Option<sui_types::messages_checkpoint::CheckpointContents> {
322        match self.get_checkpoint_by_sequence_number(sequence_number) {
323            Some(checkpoint) => self.get_checkpoint_contents_by_digest(&checkpoint.content_digest),
324            None => None,
325        }
326    }
327}
328
329impl ObjectStore for RocksDbStore {
330    fn get_object(&self, object_id: &sui_types::base_types::ObjectID) -> Option<Object> {
331        self.cache_traits.object_store.get_object(object_id)
332    }
333
334    fn get_object_by_key(
335        &self,
336        object_id: &sui_types::base_types::ObjectID,
337        version: sui_types::base_types::VersionNumber,
338    ) -> Option<Object> {
339        self.cache_traits
340            .object_store
341            .get_object_by_key(object_id, version)
342    }
343
344    fn multi_get_objects_by_key(&self, object_keys: &[ObjectKey]) -> Vec<Option<Object>> {
345        self.cache_traits
346            .object_cache_reader
347            .multi_get_objects_by_key(object_keys)
348    }
349}
350
351impl WriteStore for RocksDbStore {
352    fn insert_checkpoint(
353        &self,
354        checkpoint: &VerifiedCheckpoint,
355    ) -> Result<(), sui_types::storage::error::Error> {
356        if let Some(EndOfEpochData {
357            next_epoch_committee,
358            ..
359        }) = checkpoint.end_of_epoch_data.as_ref()
360        {
361            let next_committee = next_epoch_committee.iter().cloned().collect();
362            let committee =
363                Committee::new(checkpoint.epoch().checked_add(1).unwrap(), next_committee);
364            self.insert_committee(committee)?;
365        }
366
367        self.checkpoint_store
368            .insert_verified_checkpoint(checkpoint)
369            .map_err(Into::into)
370    }
371
372    fn update_highest_synced_checkpoint(
373        &self,
374        checkpoint: &VerifiedCheckpoint,
375    ) -> Result<(), sui_types::storage::error::Error> {
376        let mut locked = self.highest_synced_checkpoint.lock();
377        if locked.is_some() && locked.unwrap() >= checkpoint.sequence_number {
378            return Ok(());
379        }
380        self.checkpoint_store
381            .update_highest_synced_checkpoint(checkpoint)
382            .map_err(sui_types::storage::error::Error::custom)?;
383        *locked = Some(checkpoint.sequence_number);
384        Ok(())
385    }
386
387    fn update_highest_verified_checkpoint(
388        &self,
389        checkpoint: &VerifiedCheckpoint,
390    ) -> Result<(), sui_types::storage::error::Error> {
391        let mut locked = self.highest_verified_checkpoint.lock();
392        if locked.is_some() && locked.unwrap() >= checkpoint.sequence_number {
393            return Ok(());
394        }
395        self.checkpoint_store
396            .update_highest_verified_checkpoint(checkpoint)
397            .map_err(sui_types::storage::error::Error::custom)?;
398        *locked = Some(checkpoint.sequence_number);
399        Ok(())
400    }
401
402    fn insert_checkpoint_contents(
403        &self,
404        checkpoint: &VerifiedCheckpoint,
405        contents: VerifiedCheckpointContents,
406    ) -> Result<(), sui_types::storage::error::Error> {
407        self.cache_traits
408            .state_sync_store
409            .multi_insert_transaction_and_effects(contents.transactions());
410        self.checkpoint_store
411            .insert_verified_checkpoint_contents(checkpoint, contents)
412            .map_err(Into::into)
413    }
414
415    fn insert_committee(
416        &self,
417        new_committee: Committee,
418    ) -> Result<(), sui_types::storage::error::Error> {
419        self.committee_store
420            .insert_new_committee(&new_committee)
421            .unwrap();
422        Ok(())
423    }
424}
425
426pub struct RestReadStore {
427    state: Arc<AuthorityState>,
428    rocks: RocksDbStore,
429}
430
431impl RestReadStore {
432    pub fn new(state: Arc<AuthorityState>, rocks: RocksDbStore) -> Self {
433        Self { state, rocks }
434    }
435}
436
437impl ObjectStore for RestReadStore {
438    fn get_object(&self, object_id: &sui_types::base_types::ObjectID) -> Option<Object> {
439        self.rocks.get_object(object_id)
440    }
441
442    fn get_object_by_key(
443        &self,
444        object_id: &sui_types::base_types::ObjectID,
445        version: sui_types::base_types::VersionNumber,
446    ) -> Option<Object> {
447        self.rocks.get_object_by_key(object_id, version)
448    }
449
450    fn multi_get_objects_by_key(&self, object_keys: &[ObjectKey]) -> Vec<Option<Object>> {
451        self.rocks.multi_get_objects_by_key(object_keys)
452    }
453}
454
455impl ReadStore for RestReadStore {
456    fn get_committee(&self, epoch: EpochId) -> Option<Arc<Committee>> {
457        self.rocks.get_committee(epoch)
458    }
459
460    fn get_latest_checkpoint(&self) -> sui_types::storage::error::Result<VerifiedCheckpoint> {
461        self.rocks.get_latest_checkpoint()
462    }
463
464    fn get_highest_verified_checkpoint(
465        &self,
466    ) -> sui_types::storage::error::Result<VerifiedCheckpoint> {
467        self.rocks.get_highest_verified_checkpoint()
468    }
469
470    fn get_highest_synced_checkpoint(
471        &self,
472    ) -> sui_types::storage::error::Result<VerifiedCheckpoint> {
473        self.rocks.get_highest_synced_checkpoint()
474    }
475
476    fn get_lowest_available_checkpoint(
477        &self,
478    ) -> sui_types::storage::error::Result<CheckpointSequenceNumber> {
479        self.rocks.get_lowest_available_checkpoint()
480    }
481
482    fn get_checkpoint_by_digest(&self, digest: &CheckpointDigest) -> Option<VerifiedCheckpoint> {
483        self.rocks.get_checkpoint_by_digest(digest)
484    }
485
486    fn get_checkpoint_by_sequence_number(
487        &self,
488        sequence_number: CheckpointSequenceNumber,
489    ) -> Option<VerifiedCheckpoint> {
490        self.rocks
491            .get_checkpoint_by_sequence_number(sequence_number)
492    }
493
494    fn multi_get_checkpoint_by_sequence_number(
495        &self,
496        sequence_numbers: &[CheckpointSequenceNumber],
497    ) -> Vec<Option<VerifiedCheckpoint>> {
498        self.rocks
499            .multi_get_checkpoint_by_sequence_number(sequence_numbers)
500    }
501
502    fn get_checkpoint_contents_by_digest(
503        &self,
504        digest: &CheckpointContentsDigest,
505    ) -> Option<sui_types::messages_checkpoint::CheckpointContents> {
506        self.rocks.get_checkpoint_contents_by_digest(digest)
507    }
508
509    fn get_checkpoint_contents_by_sequence_number(
510        &self,
511        sequence_number: CheckpointSequenceNumber,
512    ) -> Option<sui_types::messages_checkpoint::CheckpointContents> {
513        self.rocks
514            .get_checkpoint_contents_by_sequence_number(sequence_number)
515    }
516
517    fn get_transaction(&self, digest: &TransactionDigest) -> Option<Arc<VerifiedTransaction>> {
518        self.rocks.get_transaction(digest)
519    }
520
521    fn multi_get_transactions(
522        &self,
523        digests: &[TransactionDigest],
524    ) -> Vec<Option<Arc<VerifiedTransaction>>> {
525        self.rocks.multi_get_transactions(digests)
526    }
527
528    fn get_transaction_effects(&self, digest: &TransactionDigest) -> Option<TransactionEffects> {
529        self.rocks.get_transaction_effects(digest)
530    }
531
532    fn multi_get_transaction_effects(
533        &self,
534        digests: &[TransactionDigest],
535    ) -> Vec<Option<TransactionEffects>> {
536        self.rocks.multi_get_transaction_effects(digests)
537    }
538
539    fn get_events(&self, digest: &TransactionDigest) -> Option<TransactionEvents> {
540        self.rocks.get_events(digest)
541    }
542
543    fn multi_get_events(&self, digests: &[TransactionDigest]) -> Vec<Option<TransactionEvents>> {
544        self.rocks.multi_get_events(digests)
545    }
546
547    fn get_full_checkpoint_contents(
548        &self,
549        sequence_number: Option<CheckpointSequenceNumber>,
550        digest: &CheckpointContentsDigest,
551    ) -> Option<VersionedFullCheckpointContents> {
552        self.rocks
553            .get_full_checkpoint_contents(sequence_number, digest)
554    }
555
556    fn get_unchanged_loaded_runtime_objects(
557        &self,
558        digest: &TransactionDigest,
559    ) -> Option<Vec<ObjectKey>> {
560        self.rocks.get_unchanged_loaded_runtime_objects(digest)
561    }
562
563    fn multi_get_unchanged_loaded_runtime_objects(
564        &self,
565        digests: &[TransactionDigest],
566    ) -> Vec<Option<Vec<ObjectKey>>> {
567        self.rocks
568            .multi_get_unchanged_loaded_runtime_objects(digests)
569    }
570
571    fn get_transaction_checkpoint(
572        &self,
573        digest: &TransactionDigest,
574    ) -> Option<CheckpointSequenceNumber> {
575        self.rocks.get_transaction_checkpoint(digest)
576    }
577}
578
579impl BackingPackageStore for RestReadStore {
580    fn get_package_object(&self, _package_id: &ObjectID) -> SuiResult<Option<PackageObject>> {
581        Err(SuiErrorKind::UnsupportedFeatureError {
582            error: "RestReadStore does not support loading package objects".to_string(),
583        }
584        .into())
585    }
586}
587
588impl RuntimeObjectResolver for RestReadStore {
589    fn read_child_object(
590        &self,
591        parent: &ObjectID,
592        child: &ObjectID,
593        child_version_upper_bound: SequenceNumber,
594    ) -> SuiResult<Option<Object>> {
595        Ok(self.get_object(child).and_then(|o| {
596            if o.version() <= child_version_upper_bound
597                && o.owner == Owner::ObjectOwner((*parent).into())
598            {
599                Some(o)
600            } else {
601                None
602            }
603        }))
604    }
605
606    fn get_object_received_at_version(
607        &self,
608        _owner: &ObjectID,
609        _receiving_object_id: &ObjectID,
610        _receive_object_at_version: SequenceNumber,
611        _epoch_id: EpochId,
612    ) -> SuiResult<Option<Object>> {
613        Err(SuiErrorKind::UnsupportedFeatureError {
614            error: "RestReadStore does not support receiving objects".to_string(),
615        }
616        .into())
617    }
618}
619
620impl RpcStateReader for RestReadStore {
621    fn get_lowest_available_checkpoint_objects(
622        &self,
623    ) -> sui_types::storage::error::Result<CheckpointSequenceNumber> {
624        Ok(self
625            .state
626            .get_object_cache_reader()
627            .get_highest_pruned_checkpoint()
628            .map(|cp| cp + 1)
629            .unwrap_or(0))
630    }
631
632    fn get_chain_identifier(&self) -> Result<sui_types::digests::ChainIdentifier> {
633        Ok(self.state.get_chain_identifier())
634    }
635
636    fn indexes(&self) -> Option<&dyn RpcIndexes> {
637        // The legacy `rpc-index` backend has been removed; a node serving
638        // reads through `RestReadStore` exposes no index surface. Index
639        // reads are served by the embedded rpc-store via `RpcStoreReadStore`.
640        None
641    }
642
643    fn get_struct_layout_with_overlay(
644        &self,
645        struct_tag: &move_core_types::language_storage::StructTag,
646        overlay: &ObjectSet,
647    ) -> Result<Option<move_core_types::annotated_value::MoveTypeLayout>> {
648        let backing_store = self.state.get_backing_package_store();
649        let overlay_store = OverlayBackingPackageStore::new(overlay, backing_store.as_ref());
650        let epoch_store = self.state.load_epoch_store_one_call_per_task();
651        epoch_store
652            .executor()
653            // TODO(cache) - must read through cache
654            .type_layout_resolver(epoch_store.protocol_config(), Box::new(overlay_store))
655            .get_annotated_layout(struct_tag)
656            .map(|layout| layout.into_layout())
657            .map(Some)
658            .map_err(StorageError::custom)
659    }
660}
661
662/// Read store backed by the embedded [`sui_rpc_store`] indexer.
663///
664/// Like [`RestReadStore`] it serves the `sui-rpc-api` trait stack, but it
665/// additionally exposes the index surface (which [`RestReadStore`] no
666/// longer does). This wrapper composes two backends:
667///
668/// - **Raw chain data** — objects, transactions, effects, events,
669///   checkpoints, committees, and child-object resolution — is served
670///   from the validator's perpetual / checkpoint stores
671///   ([`RocksDbStore`]), exactly like [`RestReadStore`]. The embedded
672///   rpc-store does not duplicate this data.
673/// - **The index surface** ([`RpcIndexes`]) — owner / type / balance /
674///   coin / package-version listings, epoch info, and the
675///   ledger-history bitmaps — is served from the
676///   [`RpcStoreReader`].
677///
678/// The object/state available range is the intersection of the two
679/// backends' ranges (`max` of their lower bounds): a consistent read
680/// at checkpoint `C` needs both the object bytes (perpetual store) and
681/// the index rows (rpc-store) at `C`. Ledger-history-specific
682/// availability (bounded by the history backfill watermark) is exposed
683/// separately.
684pub struct RpcStoreReadStore {
685    state: Arc<AuthorityState>,
686    rocks: RocksDbStore,
687    reader: RpcStoreReader,
688}
689
690impl RpcStoreReadStore {
691    pub fn new(state: Arc<AuthorityState>, rocks: RocksDbStore, reader: RpcStoreReader) -> Self {
692        Self {
693            state,
694            rocks,
695            reader,
696        }
697    }
698}
699
700impl ObjectStore for RpcStoreReadStore {
701    fn get_object(&self, object_id: &ObjectID) -> Option<Object> {
702        self.rocks.get_object(object_id)
703    }
704
705    fn get_object_by_key(&self, object_id: &ObjectID, version: SequenceNumber) -> Option<Object> {
706        self.rocks.get_object_by_key(object_id, version)
707    }
708
709    fn multi_get_objects_by_key(&self, object_keys: &[ObjectKey]) -> Vec<Option<Object>> {
710        self.rocks.multi_get_objects_by_key(object_keys)
711    }
712}
713
714impl ReadStore for RpcStoreReadStore {
715    fn get_committee(&self, epoch: EpochId) -> Option<Arc<Committee>> {
716        self.rocks.get_committee(epoch)
717    }
718
719    fn get_latest_checkpoint(&self) -> Result<VerifiedCheckpoint> {
720        let latest = self.rocks.get_latest_checkpoint()?;
721        // Bound the reported tip to what the live-object index has committed.
722        // The embedded indexer follows the tip asynchronously, so without this
723        // the rpc-api could surface a checkpoint -- and the transactions in it
724        // -- whose indexed state (owned objects, balances, coins) is not yet
725        // readable, breaking read-after-write consistency. The history cohort
726        // backfills independently and bounds the ledger-history APIs
727        // separately, so it does not constrain this tip.
728        match self.reader.highest_live_committed_checkpoint()? {
729            Some(indexed) if indexed < latest.sequence_number => self
730                .rocks
731                .get_checkpoint_by_sequence_number(indexed)
732                .ok_or_else(|| {
733                    StorageError::missing(format!(
734                        "live-indexed checkpoint {indexed} missing from the checkpoint store"
735                    ))
736                }),
737            Some(_) => Ok(latest),
738            // Fail closed: no live watermark means the index surface is
739            // empty -- a fresh node whose indexer has not committed its
740            // first checkpoint yet. Reporting the executed tip here would
741            // advertise checkpoints whose indexed state is not readable,
742            // the exact inconsistency the bound above exists to prevent.
743            // The window closes with the live cohort's first commit,
744            // moments after startup.
745            None => Err(StorageError::missing(
746                "the embedded rpc-store's live index has no committed checkpoint yet",
747            )),
748        }
749    }
750
751    fn get_highest_verified_checkpoint(&self) -> Result<VerifiedCheckpoint> {
752        self.rocks.get_highest_verified_checkpoint()
753    }
754
755    fn get_highest_synced_checkpoint(&self) -> Result<VerifiedCheckpoint> {
756        self.rocks.get_highest_synced_checkpoint()
757    }
758
759    fn get_lowest_available_checkpoint(&self) -> Result<CheckpointSequenceNumber> {
760        // A consistent read needs both the raw chain data (perpetual
761        // store) and the index rows (rpc-store), so the available range
762        // starts at the higher of the two lower bounds.
763        let perpetual = self.rocks.get_lowest_available_checkpoint()?;
764        let rpc_store = self.reader.get_lowest_available_checkpoint()?;
765        Ok(perpetual.max(rpc_store))
766    }
767
768    fn get_checkpoint_by_digest(&self, digest: &CheckpointDigest) -> Option<VerifiedCheckpoint> {
769        self.rocks.get_checkpoint_by_digest(digest)
770    }
771
772    fn get_checkpoint_by_sequence_number(
773        &self,
774        sequence_number: CheckpointSequenceNumber,
775    ) -> Option<VerifiedCheckpoint> {
776        self.rocks
777            .get_checkpoint_by_sequence_number(sequence_number)
778    }
779
780    fn multi_get_checkpoint_by_sequence_number(
781        &self,
782        sequence_numbers: &[CheckpointSequenceNumber],
783    ) -> Vec<Option<VerifiedCheckpoint>> {
784        self.rocks
785            .multi_get_checkpoint_by_sequence_number(sequence_numbers)
786    }
787
788    fn get_checkpoint_contents_by_digest(
789        &self,
790        digest: &CheckpointContentsDigest,
791    ) -> Option<sui_types::messages_checkpoint::CheckpointContents> {
792        self.rocks.get_checkpoint_contents_by_digest(digest)
793    }
794
795    fn get_checkpoint_contents_by_sequence_number(
796        &self,
797        sequence_number: CheckpointSequenceNumber,
798    ) -> Option<sui_types::messages_checkpoint::CheckpointContents> {
799        self.rocks
800            .get_checkpoint_contents_by_sequence_number(sequence_number)
801    }
802
803    fn get_transaction(&self, digest: &TransactionDigest) -> Option<Arc<VerifiedTransaction>> {
804        self.rocks.get_transaction(digest)
805    }
806
807    fn multi_get_transactions(
808        &self,
809        digests: &[TransactionDigest],
810    ) -> Vec<Option<Arc<VerifiedTransaction>>> {
811        self.rocks.multi_get_transactions(digests)
812    }
813
814    fn get_transaction_effects(&self, digest: &TransactionDigest) -> Option<TransactionEffects> {
815        self.rocks.get_transaction_effects(digest)
816    }
817
818    fn multi_get_transaction_effects(
819        &self,
820        digests: &[TransactionDigest],
821    ) -> Vec<Option<TransactionEffects>> {
822        self.rocks.multi_get_transaction_effects(digests)
823    }
824
825    fn get_events(&self, digest: &TransactionDigest) -> Option<TransactionEvents> {
826        self.rocks.get_events(digest)
827    }
828
829    fn multi_get_events(&self, digests: &[TransactionDigest]) -> Vec<Option<TransactionEvents>> {
830        self.rocks.multi_get_events(digests)
831    }
832
833    fn get_full_checkpoint_contents(
834        &self,
835        sequence_number: Option<CheckpointSequenceNumber>,
836        digest: &CheckpointContentsDigest,
837    ) -> Option<VersionedFullCheckpointContents> {
838        self.rocks
839            .get_full_checkpoint_contents(sequence_number, digest)
840    }
841
842    fn get_unchanged_loaded_runtime_objects(
843        &self,
844        digest: &TransactionDigest,
845    ) -> Option<Vec<ObjectKey>> {
846        self.rocks.get_unchanged_loaded_runtime_objects(digest)
847    }
848
849    fn multi_get_unchanged_loaded_runtime_objects(
850        &self,
851        digests: &[TransactionDigest],
852    ) -> Vec<Option<Vec<ObjectKey>>> {
853        self.rocks
854            .multi_get_unchanged_loaded_runtime_objects(digests)
855    }
856
857    fn get_transaction_checkpoint(
858        &self,
859        digest: &TransactionDigest,
860    ) -> Option<CheckpointSequenceNumber> {
861        self.rocks.get_transaction_checkpoint(digest)
862    }
863}
864
865impl BackingPackageStore for RpcStoreReadStore {
866    fn get_package_object(&self, _package_id: &ObjectID) -> SuiResult<Option<PackageObject>> {
867        Err(SuiErrorKind::UnsupportedFeatureError {
868            error: "RpcStoreReadStore does not support loading package objects".to_string(),
869        }
870        .into())
871    }
872}
873
874impl RuntimeObjectResolver for RpcStoreReadStore {
875    fn read_child_object(
876        &self,
877        parent: &ObjectID,
878        child: &ObjectID,
879        child_version_upper_bound: SequenceNumber,
880    ) -> SuiResult<Option<Object>> {
881        Ok(self.get_object(child).and_then(|o| {
882            if o.version() <= child_version_upper_bound
883                && o.owner == Owner::ObjectOwner((*parent).into())
884            {
885                Some(o)
886            } else {
887                None
888            }
889        }))
890    }
891
892    fn get_object_received_at_version(
893        &self,
894        _owner: &ObjectID,
895        _receiving_object_id: &ObjectID,
896        _receive_object_at_version: SequenceNumber,
897        _epoch_id: EpochId,
898    ) -> SuiResult<Option<Object>> {
899        Err(SuiErrorKind::UnsupportedFeatureError {
900            error: "RpcStoreReadStore does not support receiving objects".to_string(),
901        }
902        .into())
903    }
904}
905
906impl RpcStateReader for RpcStoreReadStore {
907    fn get_lowest_available_checkpoint_objects(&self) -> Result<CheckpointSequenceNumber> {
908        let perpetual = self
909            .state
910            .get_object_cache_reader()
911            .get_highest_pruned_checkpoint()
912            .map(|cp| cp + 1)
913            .unwrap_or(0);
914        let rpc_store = self.reader.get_lowest_available_checkpoint_objects()?;
915        Ok(perpetual.max(rpc_store))
916    }
917
918    fn get_chain_identifier(&self) -> Result<sui_types::digests::ChainIdentifier> {
919        Ok(self.state.get_chain_identifier())
920    }
921
922    fn indexes(&self) -> Option<&dyn RpcIndexes> {
923        Some(self)
924    }
925
926    fn get_highest_executed_checkpoint_seq_number(&self) -> Result<CheckpointSequenceNumber> {
927        // The raw executed tip, read straight from the checkpoint store. Unlike
928        // `get_latest_checkpoint`, this is not bounded to the live-object index
929        // frontier -- the health check measures how far that frontier trails
930        // the executed tip, so it must see the unbounded value.
931        Ok(*self.rocks.get_latest_checkpoint()?.sequence_number())
932    }
933
934    fn get_struct_layout_with_overlay(
935        &self,
936        struct_tag: &move_core_types::language_storage::StructTag,
937        overlay: &ObjectSet,
938    ) -> Result<Option<move_core_types::annotated_value::MoveTypeLayout>> {
939        // Resolve through the authority's live executor and backing
940        // package store, matching `RestReadStore`: the perpetual store
941        // backs the package reads and the loaded epoch store carries
942        // the current protocol config.
943        let backing_store = self.state.get_backing_package_store();
944        let overlay_store = OverlayBackingPackageStore::new(overlay, backing_store.as_ref());
945        let epoch_store = self.state.load_epoch_store_one_call_per_task();
946        epoch_store
947            .executor()
948            .type_layout_resolver(epoch_store.protocol_config(), Box::new(overlay_store))
949            .get_annotated_layout(struct_tag)
950            .map(|layout| layout.into_layout())
951            .map(Some)
952            .map_err(StorageError::custom)
953    }
954}
955
956impl RpcIndexes for RpcStoreReadStore {
957    fn get_epoch_info(&self, epoch: EpochId) -> Result<Option<sui_types::storage::EpochInfo>> {
958        self.reader.get_epoch_info(epoch)
959    }
960
961    fn owned_objects_iter(
962        &self,
963        owner: SuiAddress,
964        object_type: Option<StructTag>,
965        cursor: Option<OwnedObjectInfo>,
966    ) -> Result<Box<dyn Iterator<Item = Result<OwnedObjectInfo, TypedStoreError>> + '_>> {
967        self.reader.owned_objects_iter(owner, object_type, cursor)
968    }
969
970    fn dynamic_field_iter(
971        &self,
972        parent: ObjectID,
973        cursor: Option<DynamicFieldKey>,
974    ) -> Result<Box<dyn Iterator<Item = Result<DynamicFieldKey, TypedStoreError>> + '_>> {
975        self.reader.dynamic_field_iter(parent, cursor)
976    }
977
978    fn get_coin_info(&self, coin_type: &StructTag) -> Result<Option<CoinInfo>> {
979        self.reader.get_coin_info(coin_type)
980    }
981
982    fn get_balance(
983        &self,
984        owner: &SuiAddress,
985        coin_type: &StructTag,
986    ) -> Result<Option<BalanceInfo>> {
987        self.reader.get_balance(owner, coin_type)
988    }
989
990    fn balance_iter(
991        &self,
992        owner: &SuiAddress,
993        cursor: Option<(SuiAddress, StructTag)>,
994    ) -> Result<BalanceIterator<'_>> {
995        self.reader.balance_iter(owner, cursor)
996    }
997
998    fn package_versions_iter(
999        &self,
1000        original_id: ObjectID,
1001        cursor: Option<u64>,
1002    ) -> Result<Box<dyn Iterator<Item = Result<(u64, ObjectID), TypedStoreError>> + '_>> {
1003        self.reader.package_versions_iter(original_id, cursor)
1004    }
1005
1006    fn get_highest_indexed_checkpoint_seq_number(
1007        &self,
1008    ) -> Result<Option<CheckpointSequenceNumber>> {
1009        self.reader.get_highest_indexed_checkpoint_seq_number()
1010    }
1011
1012    fn get_highest_live_indexed_checkpoint_seq_number(
1013        &self,
1014    ) -> Result<Option<CheckpointSequenceNumber>> {
1015        self.reader.get_highest_live_indexed_checkpoint_seq_number()
1016    }
1017
1018    fn ledger_tx_seq_digest(&self, tx_seq: u64) -> Result<Option<LedgerTxSeqDigest>> {
1019        self.reader.ledger_tx_seq_digest(tx_seq)
1020    }
1021
1022    fn ledger_tx_seq_digest_multi_get(
1023        &self,
1024        tx_seqs: &[u64],
1025    ) -> Result<Vec<Option<LedgerTxSeqDigest>>> {
1026        self.reader.ledger_tx_seq_digest_multi_get(tx_seqs)
1027    }
1028
1029    fn ledger_tx_seq_digest_iter(
1030        &self,
1031        start: u64,
1032        end_exclusive: u64,
1033        descending: bool,
1034    ) -> Result<LedgerTxSeqDigestIterator<'_>> {
1035        self.reader
1036            .ledger_tx_seq_digest_iter(start, end_exclusive, descending)
1037    }
1038
1039    fn transaction_bitmap_bucket_iter(
1040        &self,
1041        dimension_key: Vec<u8>,
1042        start_bucket: u64,
1043        end_bucket_exclusive: u64,
1044        descending: bool,
1045    ) -> Result<LedgerBitmapBucketIterator<'_>> {
1046        self.reader.transaction_bitmap_bucket_iter(
1047            dimension_key,
1048            start_bucket,
1049            end_bucket_exclusive,
1050            descending,
1051        )
1052    }
1053
1054    fn event_bitmap_bucket_iter(
1055        &self,
1056        dimension_key: Vec<u8>,
1057        start_bucket: u64,
1058        end_bucket_exclusive: u64,
1059        descending: bool,
1060    ) -> Result<LedgerBitmapBucketIterator<'_>> {
1061        self.reader.event_bitmap_bucket_iter(
1062            dimension_key,
1063            start_bucket,
1064            end_bucket_exclusive,
1065            descending,
1066        )
1067    }
1068}