Skip to main content

sui_core/authority/
authority_store.rs

1// Copyright (c) Mysten Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4use std::sync::Arc;
5use std::{iter, mem, thread};
6
7use crate::authority::authority_per_epoch_store::AuthorityPerEpochStore;
8use crate::authority::authority_store_pruner::{
9    AuthorityStorePruner, AuthorityStorePruningMetrics, EPOCH_DURATION_MS_FOR_TESTING,
10};
11use crate::authority::authority_store_types::{StoreObject, StoreObjectWrapper, get_store_object};
12use crate::authority::epoch_marker_key::EpochMarkerKey;
13use crate::authority::epoch_start_configuration::{EpochFlag, EpochStartConfiguration};
14use crate::global_state_hasher::GlobalStateHashStore;
15use crate::transaction_outputs::TransactionOutputs;
16use fastcrypto::hash::{HashFunction, MultisetHash, Sha3_256};
17use futures::stream::FuturesUnordered;
18use move_core_types::account_address::AccountAddress;
19use move_core_types::resolver::{ModuleResolver, SerializedPackage};
20use sui_config::node::AuthorityStorePruningConfig;
21use sui_macros::fail_point_arg;
22use sui_types::execution::TypeLayoutStore;
23use sui_types::global_state_hash::GlobalStateHash;
24use sui_types::message_envelope::Message;
25use sui_types::storage::{
26    BackingPackageStore, FullObjectKey, MarkerValue, ObjectKey, ObjectOrTombstone, ObjectStore,
27    get_module, get_package,
28};
29use sui_types::sui_system_state::get_sui_system_state;
30use sui_types::{base_types::SequenceNumber, fp_ensure};
31use tokio::time::Instant;
32use tracing::{debug, info, trace};
33use typed_store::traits::Map;
34use typed_store::{TypedStoreError, rocks::DBBatch};
35
36use super::authority_store_tables::LiveObject;
37use super::{authority_store_tables::AuthorityPerpetualTables, *};
38use mysten_common::sync::notify_read::NotifyRead;
39use sui_types::effects::{TransactionEffects, TransactionEvents};
40use sui_types::gas_coin::TOTAL_SUPPLY_MIST;
41
42struct AuthorityStoreMetrics {
43    sui_conservation_check_latency: IntGauge,
44    sui_conservation_live_object_count: IntGauge,
45    sui_conservation_live_object_size: IntGauge,
46    sui_conservation_imbalance: IntGauge,
47    sui_conservation_storage_fund: IntGauge,
48    sui_conservation_storage_fund_imbalance: IntGauge,
49    epoch_flags: IntGaugeVec,
50}
51
52impl AuthorityStoreMetrics {
53    pub fn new(registry: &Registry) -> Self {
54        Self {
55            sui_conservation_check_latency: register_int_gauge_with_registry!(
56                "sui_conservation_check_latency",
57                "Number of seconds took to scan all live objects in the store for SUI conservation check",
58                registry,
59            ).unwrap(),
60            sui_conservation_live_object_count: register_int_gauge_with_registry!(
61                "sui_conservation_live_object_count",
62                "Number of live objects in the store",
63                registry,
64            ).unwrap(),
65            sui_conservation_live_object_size: register_int_gauge_with_registry!(
66                "sui_conservation_live_object_size",
67                "Size in bytes of live objects in the store",
68                registry,
69            ).unwrap(),
70            sui_conservation_imbalance: register_int_gauge_with_registry!(
71                "sui_conservation_imbalance",
72                "Total amount of SUI in the network - 10B * 10^9. This delta shows the amount of imbalance",
73                registry,
74            ).unwrap(),
75            sui_conservation_storage_fund: register_int_gauge_with_registry!(
76                "sui_conservation_storage_fund",
77                "Storage Fund pool balance (only includes the storage fund proper that represents object storage)",
78                registry,
79            ).unwrap(),
80            sui_conservation_storage_fund_imbalance: register_int_gauge_with_registry!(
81                "sui_conservation_storage_fund_imbalance",
82                "Imbalance of storage fund, computed with storage_fund_balance - total_object_storage_rebates",
83                registry,
84            ).unwrap(),
85            epoch_flags: register_int_gauge_vec_with_registry!(
86                "epoch_flags",
87                "Local flags of the currently running epoch",
88                &["flag"],
89                registry,
90            ).unwrap(),
91        }
92    }
93}
94
95/// ALL_OBJ_VER determines whether we want to store all past
96/// versions of every object in the store. Authority doesn't store
97/// them, but other entities such as replicas will.
98/// S is a template on Authority signature state. This allows SuiDataStore to be used on either
99/// authorities or non-authorities. Specifically, when storing transactions and effects,
100/// S allows SuiDataStore to either store the authority signed version or unsigned version.
101pub struct AuthorityStore {
102    pub(crate) perpetual_tables: Arc<AuthorityPerpetualTables>,
103
104    pub(crate) root_state_notify_read:
105        NotifyRead<EpochId, (CheckpointSequenceNumber, GlobalStateHash)>,
106
107    /// Whether to enable expensive SUI conservation check at epoch boundaries.
108    enable_epoch_sui_conservation_check: bool,
109
110    metrics: AuthorityStoreMetrics,
111}
112
113pub type ExecutionLockReadGuard<'a> = tokio::sync::RwLockReadGuard<'a, EpochId>;
114pub type ExecutionLockWriteGuard<'a> = tokio::sync::RwLockWriteGuard<'a, EpochId>;
115
116impl AuthorityStore {
117    /// Open an authority store by directory path.
118    /// If the store is empty, initialize it using genesis.
119    pub async fn open(
120        perpetual_tables: Arc<AuthorityPerpetualTables>,
121        genesis: &Genesis,
122        config: &NodeConfig,
123        registry: &Registry,
124    ) -> SuiResult<Arc<Self>> {
125        let enable_epoch_sui_conservation_check = config
126            .expensive_safety_check_config
127            .enable_epoch_sui_conservation_check();
128
129        let epoch_start_configuration = if perpetual_tables.database_is_empty()? {
130            info!("Creating new epoch start config from genesis");
131
132            #[allow(unused_mut)]
133            let mut initial_epoch_flags = EpochFlag::default_flags_for_new_epoch(config);
134            fail_point_arg!("initial_epoch_flags", |flags: Vec<EpochFlag>| {
135                info!("Setting initial epoch flags to {:?}", flags);
136                initial_epoch_flags = flags;
137            });
138
139            let epoch_start_configuration = EpochStartConfiguration::new(
140                genesis.sui_system_object().into_epoch_start_state(),
141                *genesis.checkpoint().digest(),
142                &genesis.objects(),
143                initial_epoch_flags,
144            )?;
145            perpetual_tables.set_epoch_start_configuration(&epoch_start_configuration)?;
146            epoch_start_configuration
147        } else {
148            info!("Loading epoch start config from DB");
149            perpetual_tables
150                .epoch_start_configuration
151                .get(&())?
152                .expect("Epoch start configuration must be set in non-empty DB")
153        };
154        let cur_epoch = perpetual_tables.get_recovery_epoch_at_restart()?;
155        info!("Epoch start config: {:?}", epoch_start_configuration);
156        info!("Cur epoch: {:?}", cur_epoch);
157        let this = Self::open_inner(
158            genesis,
159            perpetual_tables,
160            enable_epoch_sui_conservation_check,
161            registry,
162        )
163        .await?;
164        this.update_epoch_flags_metrics(&[], epoch_start_configuration.flags());
165        Ok(this)
166    }
167
168    pub fn update_epoch_flags_metrics(&self, old: &[EpochFlag], new: &[EpochFlag]) {
169        for flag in old {
170            self.metrics
171                .epoch_flags
172                .with_label_values(&[&flag.to_string()])
173                .set(0);
174        }
175        for flag in new {
176            self.metrics
177                .epoch_flags
178                .with_label_values(&[&flag.to_string()])
179                .set(1);
180        }
181    }
182
183    // NB: This must only be called at time of reconfiguration. We take the execution lock write
184    // guard as an argument to ensure that this is the case.
185    pub fn clear_object_per_epoch_marker_table(
186        &self,
187        _execution_guard: &ExecutionLockWriteGuard<'_>,
188    ) -> SuiResult<()> {
189        // We can safely delete all entries in the per epoch marker table since this is only called
190        // at epoch boundaries (during reconfiguration). Therefore any entries that currently
191        // exist can be removed. Because of this we can use the `schedule_delete_all` method.
192        self.perpetual_tables
193            .object_per_epoch_marker_table
194            .schedule_delete_all()?;
195        #[cfg(not(tidehunter))]
196        {
197            self.perpetual_tables
198                .object_per_epoch_marker_table_v2
199                .schedule_delete_all()?;
200        }
201        #[cfg(tidehunter)]
202        {
203            self.perpetual_tables
204                .object_per_epoch_marker_table_v2
205                .drop_cells_in_range_raw(&EpochMarkerKey::MIN_KEY, &EpochMarkerKey::MAX_KEY)?;
206        }
207        Ok(())
208    }
209
210    pub async fn open_with_committee_for_testing(
211        perpetual_tables: Arc<AuthorityPerpetualTables>,
212        committee: &Committee,
213        genesis: &Genesis,
214    ) -> SuiResult<Arc<Self>> {
215        // TODO: Since we always start at genesis, the committee should be technically the same
216        // as the genesis committee.
217        assert_eq!(committee.epoch, 0);
218        Self::open_inner(genesis, perpetual_tables, true, &Registry::new()).await
219    }
220
221    async fn open_inner(
222        genesis: &Genesis,
223        perpetual_tables: Arc<AuthorityPerpetualTables>,
224        enable_epoch_sui_conservation_check: bool,
225        registry: &Registry,
226    ) -> SuiResult<Arc<Self>> {
227        let store = Arc::new(Self {
228            perpetual_tables,
229            root_state_notify_read: NotifyRead::<
230                EpochId,
231                (CheckpointSequenceNumber, GlobalStateHash),
232            >::new(),
233            enable_epoch_sui_conservation_check,
234            metrics: AuthorityStoreMetrics::new(registry),
235        });
236        // Only initialize an empty database.
237        if store
238            .database_is_empty()
239            .expect("Database read should not fail at init.")
240        {
241            store
242                .bulk_insert_genesis_objects(genesis.objects())
243                .expect("Cannot bulk insert genesis objects");
244
245            // insert txn and effects of genesis
246            let transaction = VerifiedTransaction::new_unchecked(genesis.transaction().clone());
247
248            store
249                .perpetual_tables
250                .transactions
251                .insert(transaction.digest(), transaction.serializable_ref())
252                .unwrap();
253
254            store
255                .perpetual_tables
256                .effects
257                .insert(&genesis.effects().digest(), genesis.effects())
258                .unwrap();
259            // We don't insert the effects to executed_effects yet because the genesis tx hasn't but will be executed.
260            // This is important for fullnodes to be able to generate indexing data right now.
261
262            if genesis.effects().events_digest().is_some() {
263                store
264                    .perpetual_tables
265                    .events_2
266                    .insert(transaction.digest(), genesis.events())
267                    .unwrap();
268            }
269        }
270
271        Ok(store)
272    }
273
274    /// Open authority store without any operations that require
275    /// genesis, such as constructing EpochStartConfiguration
276    /// or inserting genesis objects.
277    pub fn open_no_genesis(
278        perpetual_tables: Arc<AuthorityPerpetualTables>,
279        enable_epoch_sui_conservation_check: bool,
280        registry: &Registry,
281    ) -> SuiResult<Arc<Self>> {
282        let store = Arc::new(Self {
283            perpetual_tables,
284            root_state_notify_read: NotifyRead::<
285                EpochId,
286                (CheckpointSequenceNumber, GlobalStateHash),
287            >::new(),
288            enable_epoch_sui_conservation_check,
289            metrics: AuthorityStoreMetrics::new(registry),
290        });
291        Ok(store)
292    }
293
294    pub fn get_recovery_epoch_at_restart(&self) -> SuiResult<EpochId> {
295        self.perpetual_tables.get_recovery_epoch_at_restart()
296    }
297
298    pub fn get_effects(
299        &self,
300        effects_digest: &TransactionEffectsDigest,
301    ) -> SuiResult<Option<TransactionEffects>> {
302        Ok(self.perpetual_tables.effects.get(effects_digest)?)
303    }
304
305    pub fn get_events(
306        &self,
307        digest: &TransactionDigest,
308    ) -> Result<Option<TransactionEvents>, TypedStoreError> {
309        self.perpetual_tables.events_2.get(digest)
310    }
311
312    pub fn multi_get_events(
313        &self,
314        event_digests: &[TransactionDigest],
315    ) -> SuiResult<Vec<Option<TransactionEvents>>> {
316        Ok(event_digests
317            .iter()
318            .map(|digest| self.get_events(digest))
319            .collect::<Result<Vec<_>, _>>()?)
320    }
321
322    pub fn get_unchanged_loaded_runtime_objects(
323        &self,
324        digest: &TransactionDigest,
325    ) -> Result<Option<Vec<ObjectKey>>, TypedStoreError> {
326        self.perpetual_tables
327            .unchanged_loaded_runtime_objects
328            .get(digest)
329    }
330
331    pub fn multi_get_unchanged_loaded_runtime_objects(
332        &self,
333        digests: &[TransactionDigest],
334    ) -> Result<Vec<Option<Vec<ObjectKey>>>, TypedStoreError> {
335        self.perpetual_tables
336            .unchanged_loaded_runtime_objects
337            .multi_get(digests)
338    }
339
340    pub fn multi_get_effects<'a>(
341        &self,
342        effects_digests: impl Iterator<Item = &'a TransactionEffectsDigest>,
343    ) -> Result<Vec<Option<TransactionEffects>>, TypedStoreError> {
344        self.perpetual_tables.effects.multi_get(effects_digests)
345    }
346
347    pub fn get_executed_effects(
348        &self,
349        tx_digest: &TransactionDigest,
350    ) -> Result<Option<TransactionEffects>, TypedStoreError> {
351        let effects_digest = self.perpetual_tables.executed_effects.get(tx_digest)?;
352        match effects_digest {
353            Some(digest) => Ok(self.perpetual_tables.effects.get(&digest)?),
354            None => Ok(None),
355        }
356    }
357
358    /// Given a list of transaction digests, returns a list of the corresponding effects only if they have been
359    /// executed. For transactions that have not been executed, None is returned.
360    pub fn multi_get_executed_effects_digests(
361        &self,
362        digests: &[TransactionDigest],
363    ) -> Result<Vec<Option<TransactionEffectsDigest>>, TypedStoreError> {
364        self.perpetual_tables.executed_effects.multi_get(digests)
365    }
366
367    /// Given a list of transaction digests, returns a list of the corresponding effects only if they have been
368    /// executed. For transactions that have not been executed, None is returned.
369    pub fn multi_get_executed_effects(
370        &self,
371        digests: &[TransactionDigest],
372    ) -> Result<Vec<Option<TransactionEffects>>, TypedStoreError> {
373        let executed_effects_digests = self.perpetual_tables.executed_effects.multi_get(digests)?;
374        let effects = self.multi_get_effects(executed_effects_digests.iter().flatten())?;
375        let mut tx_to_effects_map = effects
376            .into_iter()
377            .flatten()
378            .map(|effects| (*effects.transaction_digest(), effects))
379            .collect::<HashMap<_, _>>();
380        Ok(digests
381            .iter()
382            .map(|digest| tx_to_effects_map.remove(digest))
383            .collect())
384    }
385
386    pub fn is_tx_already_executed(&self, digest: &TransactionDigest) -> SuiResult<bool> {
387        Ok(self
388            .perpetual_tables
389            .executed_effects
390            .contains_key(digest)?)
391    }
392
393    pub fn get_marker_value(
394        &self,
395        object_key: FullObjectKey,
396        epoch_id: EpochId,
397    ) -> SuiResult<Option<MarkerValue>> {
398        Ok(self
399            .perpetual_tables
400            .object_per_epoch_marker_table_v2
401            .get(&EpochMarkerKey(epoch_id, object_key))?)
402    }
403
404    pub fn get_latest_marker(
405        &self,
406        object_id: FullObjectID,
407        epoch_id: EpochId,
408    ) -> SuiResult<Option<(SequenceNumber, MarkerValue)>> {
409        let min_key = EpochMarkerKey(epoch_id, FullObjectKey::min_for_id(&object_id));
410        let max_key = EpochMarkerKey(epoch_id, FullObjectKey::max_for_id(&object_id));
411
412        let marker_entry = self
413            .perpetual_tables
414            .object_per_epoch_marker_table_v2
415            .reversed_safe_iter_with_bounds(Some(min_key), Some(max_key))?
416            .next();
417        match marker_entry {
418            Some(Ok((EpochMarkerKey(epoch, key), marker))) => {
419                // because of the iterator bounds these cannot fail
420                assert_eq!(epoch, epoch_id);
421                assert_eq!(key.id(), object_id);
422                Ok(Some((key.version(), marker)))
423            }
424            Some(Err(e)) => Err(e.into()),
425            None => Ok(None),
426        }
427    }
428
429    // DEPRECATED -- use function of same name in AuthorityPerEpochStore
430    pub fn deprecated_insert_finalized_transactions(
431        &self,
432        digests: &[TransactionDigest],
433        epoch: EpochId,
434        sequence: CheckpointSequenceNumber,
435    ) -> SuiResult {
436        let mut batch = self
437            .perpetual_tables
438            .executed_transactions_to_checkpoint
439            .batch();
440        batch.insert_batch(
441            &self.perpetual_tables.executed_transactions_to_checkpoint,
442            digests.iter().map(|d| (*d, (epoch, sequence))),
443        )?;
444        batch.write()?;
445        trace!("Transactions {digests:?} finalized at checkpoint {sequence} epoch {epoch}");
446        Ok(())
447    }
448
449    // DEPRECATED -- use function of same name in AuthorityPerEpochStore
450    pub fn deprecated_get_transaction_checkpoint(
451        &self,
452        digest: &TransactionDigest,
453    ) -> SuiResult<Option<(EpochId, CheckpointSequenceNumber)>> {
454        Ok(self
455            .perpetual_tables
456            .executed_transactions_to_checkpoint
457            .get(digest)?)
458    }
459
460    // DEPRECATED -- use function of same name in AuthorityPerEpochStore
461    pub fn deprecated_multi_get_transaction_checkpoint(
462        &self,
463        digests: &[TransactionDigest],
464    ) -> SuiResult<Vec<Option<(EpochId, CheckpointSequenceNumber)>>> {
465        Ok(self
466            .perpetual_tables
467            .executed_transactions_to_checkpoint
468            .multi_get(digests)?
469            .into_iter()
470            .collect())
471    }
472
473    /// Returns true if there are no objects in the database
474    pub fn database_is_empty(&self) -> SuiResult<bool> {
475        self.perpetual_tables.database_is_empty()
476    }
477
478    pub fn object_exists_by_key(
479        &self,
480        object_id: &ObjectID,
481        version: VersionNumber,
482    ) -> SuiResult<bool> {
483        Ok(self
484            .perpetual_tables
485            .objects
486            .contains_key(&ObjectKey(*object_id, version))?)
487    }
488
489    pub fn multi_object_exists_by_key(&self, object_keys: &[ObjectKey]) -> SuiResult<Vec<bool>> {
490        Ok(self
491            .perpetual_tables
492            .objects
493            .multi_contains_keys(object_keys.to_vec())?
494            .into_iter()
495            .collect())
496    }
497
498    fn get_object_ref_prior_to_key(
499        &self,
500        object_id: &ObjectID,
501        version: VersionNumber,
502    ) -> Result<Option<ObjectRef>, SuiError> {
503        let Some(prior_version) = version.one_before() else {
504            return Ok(None);
505        };
506        let mut iterator = self
507            .perpetual_tables
508            .objects
509            .reversed_safe_iter_with_bounds(
510                Some(ObjectKey::min_for_id(object_id)),
511                Some(ObjectKey(*object_id, prior_version)),
512            )?;
513
514        if let Some((object_key, value)) = iterator.next().transpose()?
515            && object_key.0 == *object_id
516        {
517            return Ok(Some(
518                self.perpetual_tables.object_reference(&object_key, value)?,
519            ));
520        }
521        Ok(None)
522    }
523
524    pub fn multi_get_objects_by_key(
525        &self,
526        object_keys: &[ObjectKey],
527    ) -> Result<Vec<Option<Object>>, SuiError> {
528        let wrappers = self
529            .perpetual_tables
530            .objects
531            .multi_get(object_keys.to_vec())?;
532        let mut ret = vec![];
533
534        for (idx, w) in wrappers.into_iter().enumerate() {
535            ret.push(
536                w.map(|object| self.perpetual_tables.object(&object_keys[idx], object))
537                    .transpose()?
538                    .flatten(),
539            );
540        }
541        Ok(ret)
542    }
543
544    /// Get many objects
545    pub fn get_objects(&self, objects: &[ObjectID]) -> Result<Vec<Option<Object>>, SuiError> {
546        let mut result = Vec::new();
547        for id in objects {
548            result.push(self.get_object(id));
549        }
550        Ok(result)
551    }
552
553    // Methods to mutate the store
554
555    /// Insert a genesis object.
556    /// TODO: delete this method entirely (still used by authority_tests.rs)
557    pub(crate) fn insert_genesis_object(&self, object: Object) -> SuiResult {
558        // We only side load objects with a genesis parent transaction.
559        debug_assert!(object.previous_transaction == TransactionDigest::genesis_marker());
560        let object_ref = object.compute_object_reference();
561        self.insert_object_direct(object_ref, &object)
562    }
563
564    /// Insert an object directly into the store, and also update relevant tables
565    /// NOTE: does not handle transaction lock.
566    /// This is used to insert genesis objects
567    fn insert_object_direct(&self, object_ref: ObjectRef, object: &Object) -> SuiResult {
568        let mut write_batch = self.perpetual_tables.objects.batch();
569
570        // Insert object
571        let store_object = get_store_object(object.clone());
572        write_batch.insert_batch(
573            &self.perpetual_tables.objects,
574            std::iter::once((ObjectKey::from(object_ref), store_object)),
575        )?;
576
577        write_batch.write()?;
578
579        Ok(())
580    }
581
582    /// This function should only be used for initializing genesis and should remain private.
583    #[instrument(level = "debug", skip_all)]
584    pub(crate) fn bulk_insert_genesis_objects(&self, objects: &[Object]) -> SuiResult<()> {
585        let mut batch = self.perpetual_tables.objects.batch();
586        let ref_and_objects: Vec<_> = objects
587            .iter()
588            .map(|o| (o.compute_object_reference(), o))
589            .collect();
590
591        batch.insert_batch(
592            &self.perpetual_tables.objects,
593            ref_and_objects
594                .iter()
595                .map(|(oref, o)| (ObjectKey::from(oref), get_store_object((*o).clone()))),
596        )?;
597
598        batch.write()?;
599
600        Ok(())
601    }
602
603    pub async fn bulk_insert_live_objects(
604        perpetual_db: Arc<AuthorityPerpetualTables>,
605        objects: Vec<LiveObject>,
606        expected_sha3_digest: &[u8; 32],
607        num_parallel_chunks: usize,
608    ) -> SuiResult<()> {
609        // Verify SHA3 over the full object set before inserting.
610        let mut hasher = Sha3_256::default();
611        for object in &objects {
612            hasher.update(object.object_reference().2.inner());
613        }
614        let sha3_digest = hasher.finalize().digest;
615        if *expected_sha3_digest != sha3_digest {
616            error!(
617                "Sha does not match! expected: {:?}, actual: {:?}",
618                expected_sha3_digest, sha3_digest
619            );
620            return Err(SuiError::from("Sha does not match"));
621        }
622
623        let chunk_size = objects.len().div_ceil(num_parallel_chunks).max(1);
624        let mut remaining = objects;
625        let mut handles = Vec::new();
626        while !remaining.is_empty() {
627            let take = chunk_size.min(remaining.len());
628            let chunk: Vec<LiveObject> = remaining.drain(..take).collect();
629            let db = perpetual_db.clone();
630            handles.push(tokio::task::spawn_blocking(move || {
631                Self::insert_objects_chunk(db, chunk)
632            }));
633        }
634        for handle in handles {
635            handle.await.expect("insert task panicked")?;
636        }
637        Ok(())
638    }
639
640    fn insert_objects_chunk(
641        perpetual_db: Arc<AuthorityPerpetualTables>,
642        objects: Vec<LiveObject>,
643    ) -> SuiResult<()> {
644        let mut batch = perpetual_db.objects.batch();
645        let mut written = 0usize;
646        const MAX_BATCH_SIZE: usize = 100_000;
647        for object in objects {
648            match object {
649                LiveObject::Normal(object) => {
650                    let object_key = ObjectKey::from(object.compute_object_reference());
651                    let store_object_wrapper = get_store_object(object);
652                    batch.insert_batch(
653                        &perpetual_db.objects,
654                        std::iter::once((object_key, store_object_wrapper)),
655                    )?;
656                }
657                LiveObject::Wrapped(object_key) => {
658                    batch.insert_batch(
659                        &perpetual_db.objects,
660                        std::iter::once::<(ObjectKey, StoreObjectWrapper)>((
661                            object_key,
662                            StoreObject::Wrapped.into(),
663                        )),
664                    )?;
665                }
666            }
667            written += 1;
668            if written > MAX_BATCH_SIZE {
669                batch.write()?;
670                batch = perpetual_db.objects.batch();
671                written = 0;
672            }
673        }
674        batch.write()?;
675        Ok(())
676    }
677
678    pub fn set_epoch_start_configuration(
679        &self,
680        epoch_start_configuration: &EpochStartConfiguration,
681    ) -> SuiResult {
682        self.perpetual_tables
683            .set_epoch_start_configuration(epoch_start_configuration)?;
684        Ok(())
685    }
686
687    pub fn get_epoch_start_configuration(&self) -> SuiResult<Option<EpochStartConfiguration>> {
688        Ok(self.perpetual_tables.epoch_start_configuration.get(&())?)
689    }
690
691    /// Updates the state resulting from the execution of a certificate.
692    ///
693    /// Internally it checks that all locks for active inputs are at the correct
694    /// version, and then writes objects, certificates, parents and clean up locks atomically.
695    #[instrument(level = "debug", skip_all)]
696    pub fn build_db_batch(
697        &self,
698        epoch_id: EpochId,
699        tx_outputs: &[Arc<TransactionOutputs>],
700    ) -> SuiResult<DBBatch> {
701        let mut write_batch = self.perpetual_tables.transactions.batch();
702        for outputs in tx_outputs {
703            self.write_one_transaction_outputs(&mut write_batch, epoch_id, outputs)?;
704        }
705        // test crashing before writing the batch
706        fail_point!("crash");
707
708        trace!(
709            "built batch for committed transactions: {:?}",
710            tx_outputs
711                .iter()
712                .map(|tx| tx.transaction.digest())
713                .collect::<Vec<_>>()
714        );
715
716        // test crashing before notifying
717        fail_point!("crash");
718
719        Ok(write_batch)
720    }
721
722    fn write_one_transaction_outputs(
723        &self,
724        write_batch: &mut DBBatch,
725        epoch_id: EpochId,
726        tx_outputs: &TransactionOutputs,
727    ) -> SuiResult {
728        let TransactionOutputs {
729            transaction,
730            effects,
731            markers,
732            wrapped,
733            deleted,
734            written,
735            events,
736            unchanged_loaded_runtime_objects,
737            ..
738        } = tx_outputs;
739
740        let effects_digest = effects.digest();
741        let transaction_digest = transaction.digest();
742        // effects must be inserted before the corresponding dependent entries
743        // because they carry epoch information necessary for correct pruning via relocation filters
744        write_batch
745            .insert_batch(
746                &self.perpetual_tables.effects,
747                [(effects_digest, effects.clone())],
748            )?
749            .insert_batch(
750                &self.perpetual_tables.executed_effects,
751                [(transaction_digest, effects_digest)],
752            )?;
753
754        // Store the certificate indexed by transaction digest
755        write_batch.insert_batch(
756            &self.perpetual_tables.transactions,
757            iter::once((transaction_digest, transaction.serializable_ref())),
758        )?;
759
760        write_batch.insert_batch(
761            &self.perpetual_tables.executed_transaction_digests,
762            [((epoch_id, *transaction_digest), ())],
763        )?;
764
765        // Add batched writes for objects and locks.
766        write_batch.insert_batch(
767            &self.perpetual_tables.object_per_epoch_marker_table_v2,
768            markers
769                .iter()
770                .map(|(key, marker_value)| (EpochMarkerKey(epoch_id, *key), *marker_value)),
771        )?;
772        write_batch.insert_batch(
773            &self.perpetual_tables.objects,
774            deleted
775                .iter()
776                .map(|key| (key, StoreObject::Deleted))
777                .chain(wrapped.iter().map(|key| (key, StoreObject::Wrapped)))
778                .map(|(key, store_object)| (key, StoreObjectWrapper::from(store_object))),
779        )?;
780
781        // Insert each output object into the stores
782        let new_objects = written.iter().map(|(id, new_object)| {
783            let version = new_object.version();
784            trace!(?id, ?version, "writing object");
785            let store_object = get_store_object(new_object.clone());
786            (ObjectKey(*id, version), store_object)
787        });
788
789        write_batch.insert_batch(&self.perpetual_tables.objects, new_objects)?;
790
791        // Write events into the new table keyed off of transaction_digest
792        if effects.events_digest().is_some() {
793            write_batch.insert_batch(
794                &self.perpetual_tables.events_2,
795                [(transaction_digest, events)],
796            )?;
797        }
798
799        // Write unchanged_loaded_runtime_objects
800        if !unchanged_loaded_runtime_objects.is_empty() {
801            write_batch.insert_batch(
802                &self.perpetual_tables.unchanged_loaded_runtime_objects,
803                [(transaction_digest, unchanged_loaded_runtime_objects)],
804            )?;
805        }
806
807        debug!(effects_digest = ?effects.digest(), "commit_certificate finished");
808
809        Ok(())
810    }
811
812    /// Commits transactions only (not effects or other transaction outputs) to the db.
813    /// See ExecutionCache::persist_transaction for more info
814    pub(crate) fn persist_transaction(&self, tx: &VerifiedExecutableTransaction) -> SuiResult {
815        let mut batch = self.perpetual_tables.transactions.batch();
816        batch.insert_batch(
817            &self.perpetual_tables.transactions,
818            [(tx.digest(), tx.clone().into_unsigned().serializable_ref())],
819        )?;
820        batch.write()?;
821        Ok(())
822    }
823
824    /// Return the object with version less then or eq to the provided seq number.
825    /// This is used by indexer to find the correct version of dynamic field child object.
826    /// We do not store the version of the child object, but because of lamport timestamp,
827    /// we know the child must have version number less then or eq to the parent.
828    pub fn find_object_lt_or_eq_version(
829        &self,
830        object_id: ObjectID,
831        version: SequenceNumber,
832    ) -> SuiResult<Option<Object>> {
833        self.perpetual_tables
834            .find_object_lt_or_eq_version(object_id, version)
835    }
836
837    /// Returns the latest object reference we have for this object_id in the objects table.
838    ///
839    /// The method may also return the reference to a deleted object with a digest of
840    /// ObjectDigest::deleted() or ObjectDigest::wrapped() and lamport version
841    /// of a transaction that deleted the object.
842    /// Note that a deleted object may re-appear if the deletion was the result of the object
843    /// being wrapped in another object.
844    ///
845    /// If no entry for the object_id is found, return None.
846    pub fn get_latest_object_ref_or_tombstone(
847        &self,
848        object_id: ObjectID,
849    ) -> Result<Option<ObjectRef>, SuiError> {
850        self.perpetual_tables
851            .get_latest_object_ref_or_tombstone(object_id)
852    }
853
854    /// Returns the latest object we have for this object_id in the objects table.
855    ///
856    /// If no entry for the object_id is found, return None.
857    pub fn get_latest_object_or_tombstone(
858        &self,
859        object_id: ObjectID,
860    ) -> Result<Option<(ObjectKey, ObjectOrTombstone)>, SuiError> {
861        let Some((object_key, store_object)) = self
862            .perpetual_tables
863            .get_latest_object_or_tombstone(object_id)?
864        else {
865            return Ok(None);
866        };
867
868        if let Some(object_ref) = self
869            .perpetual_tables
870            .tombstone_reference(&object_key, &store_object)?
871        {
872            return Ok(Some((object_key, ObjectOrTombstone::Tombstone(object_ref))));
873        }
874
875        let object = self
876            .perpetual_tables
877            .object(&object_key, store_object)?
878            .expect("Non tombstone store object could not be converted to object");
879
880        Ok(Some((object_key, ObjectOrTombstone::Object(object))))
881    }
882
883    pub fn insert_transaction_and_effects(
884        &self,
885        transaction: &VerifiedTransaction,
886        transaction_effects: &TransactionEffects,
887    ) -> Result<(), TypedStoreError> {
888        let mut write_batch = self.perpetual_tables.transactions.batch();
889        // effects must be inserted before the corresponding transaction entry
890        // because they carry epoch information necessary for correct pruning via relocation filters
891        write_batch
892            .insert_batch(
893                &self.perpetual_tables.effects,
894                [(transaction_effects.digest(), transaction_effects)],
895            )?
896            .insert_batch(
897                &self.perpetual_tables.transactions,
898                [(transaction.digest(), transaction.serializable_ref())],
899            )?;
900
901        write_batch.write()?;
902        Ok(())
903    }
904
905    pub fn multi_insert_transaction_and_effects<'a>(
906        &self,
907        transactions: impl Iterator<Item = &'a VerifiedExecutionData>,
908    ) -> Result<(), TypedStoreError> {
909        let mut write_batch = self.perpetual_tables.transactions.batch();
910        for tx in transactions {
911            write_batch
912                .insert_batch(
913                    &self.perpetual_tables.effects,
914                    [(tx.effects.digest(), &tx.effects)],
915                )?
916                .insert_batch(
917                    &self.perpetual_tables.transactions,
918                    [(tx.transaction.digest(), tx.transaction.serializable_ref())],
919                )?;
920        }
921
922        write_batch.write()?;
923        Ok(())
924    }
925
926    pub fn multi_get_transaction_blocks(
927        &self,
928        tx_digests: &[TransactionDigest],
929    ) -> Result<Vec<Option<VerifiedTransaction>>, TypedStoreError> {
930        self.perpetual_tables
931            .transactions
932            .multi_get(tx_digests)
933            .map(|v| v.into_iter().map(|v| v.map(|v| v.into())).collect())
934    }
935
936    pub fn get_transaction_block(
937        &self,
938        tx_digest: &TransactionDigest,
939    ) -> Result<Option<VerifiedTransaction>, TypedStoreError> {
940        self.perpetual_tables
941            .transactions
942            .get(tx_digest)
943            .map(|v| v.map(|v| v.into()))
944    }
945
946    pub fn list_transactions_from(
947        &self,
948        start: Option<TransactionDigest>,
949        limit: usize,
950    ) -> Result<Vec<TransactionDigest>, TypedStoreError> {
951        self.perpetual_tables.list_transactions_from(start, limit)
952    }
953
954    pub fn get_executed_effects_digest_for_tx(
955        &self,
956        tx_digest: &TransactionDigest,
957    ) -> Result<Option<TransactionEffectsDigest>, TypedStoreError> {
958        self.perpetual_tables.get_executed_effects_digest(tx_digest)
959    }
960
961    /// This function reads the DB directly to get the system state object.
962    /// If reconfiguration is happening at the same time, there is no guarantee whether we would be getting
963    /// the old or the new system state object.
964    /// Hence this function should only be called during RPC reads where data race is not a major concern.
965    /// In general we should avoid this as much as possible.
966    /// If the intent is for testing, you can use AuthorityState:: get_sui_system_state_object_for_testing.
967    pub fn get_sui_system_state_object_unsafe(&self) -> SuiResult<SuiSystemState> {
968        get_sui_system_state(self.perpetual_tables.as_ref())
969    }
970
971    pub fn expensive_check_sui_conservation<T>(
972        self: &Arc<Self>,
973        type_layout_store: T,
974        old_epoch_store: &AuthorityPerEpochStore,
975    ) -> SuiResult
976    where
977        T: TypeLayoutStore + Send + Copy,
978    {
979        if !self.enable_epoch_sui_conservation_check {
980            return Ok(());
981        }
982
983        let executor = old_epoch_store.executor();
984        info!("Starting SUI conservation check. This may take a while..");
985        let cur_time = Instant::now();
986        let mut pending_objects = vec![];
987        let mut count = 0;
988        let mut size = 0;
989        let (mut total_sui, mut total_storage_rebate) = thread::scope(|s| {
990            let pending_tasks = FuturesUnordered::new();
991            for o in self.iter_live_object_set(false) {
992                match o {
993                    LiveObject::Normal(object) => {
994                        size += object.object_size_for_gas_metering();
995                        count += 1;
996                        pending_objects.push(object);
997                        if count % 1_000_000 == 0 {
998                            let mut task_objects = vec![];
999                            mem::swap(&mut pending_objects, &mut task_objects);
1000                            pending_tasks.push(s.spawn(move || {
1001                                let mut layout_resolver = executor.type_layout_resolver(
1002                                    old_epoch_store.protocol_config(),
1003                                    Box::new(type_layout_store),
1004                                );
1005                                let mut total_storage_rebate = 0;
1006                                let mut total_sui = 0;
1007                                for object in task_objects {
1008                                    total_storage_rebate += object.storage_rebate;
1009                                    // get_total_sui includes storage rebate, however all storage rebate is
1010                                    // also stored in the storage fund, so we need to subtract it here.
1011                                    let object_contained_sui = match object
1012                                        .get_total_sui(layout_resolver.as_mut())
1013                                    {
1014                                        Ok(sui) => sui,
1015                                        Err(e)
1016                                            if old_epoch_store.get_chain()
1017                                                == sui_protocol_config::Chain::Testnet =>
1018                                        {
1019                                            error!(
1020                                                "Error calculating total SUI for object {:?}: {:?}",
1021                                                object.compute_object_reference(),
1022                                                e
1023                                            );
1024                                            0
1025                                        }
1026                                        Err(e) => panic!(
1027                                            "Error calculating total SUI for object {:?}: {:?}",
1028                                            object.compute_object_reference(),
1029                                            e
1030                                        ),
1031                                    };
1032                                    total_sui += object_contained_sui - object.storage_rebate;
1033                                }
1034                                if count % 50_000_000 == 0 {
1035                                    info!("Processed {} objects", count);
1036                                }
1037                                (total_sui, total_storage_rebate)
1038                            }));
1039                        }
1040                    }
1041                    LiveObject::Wrapped(_) => {
1042                        unreachable!("Explicitly asked to not include wrapped tombstones")
1043                    }
1044                }
1045            }
1046            pending_tasks.into_iter().fold((0, 0), |init, result| {
1047                let result = result.join().unwrap();
1048                (init.0 + result.0, init.1 + result.1)
1049            })
1050        });
1051        let mut layout_resolver = executor.type_layout_resolver(
1052            old_epoch_store.protocol_config(),
1053            Box::new(type_layout_store),
1054        );
1055        for object in pending_objects {
1056            total_storage_rebate += object.storage_rebate;
1057            total_sui +=
1058                object.get_total_sui(layout_resolver.as_mut()).unwrap() - object.storage_rebate;
1059        }
1060        info!(
1061            "Scanned {} live objects, took {:?}",
1062            count,
1063            cur_time.elapsed()
1064        );
1065        self.metrics
1066            .sui_conservation_live_object_count
1067            .set(count as i64);
1068        self.metrics
1069            .sui_conservation_live_object_size
1070            .set(size as i64);
1071        self.metrics
1072            .sui_conservation_check_latency
1073            .set(cur_time.elapsed().as_secs() as i64);
1074
1075        // It is safe to call this function because we are in the middle of reconfiguration.
1076        let system_state = self
1077            .get_sui_system_state_object_unsafe()
1078            .expect("Reading sui system state object cannot fail")
1079            .into_sui_system_state_summary();
1080        let storage_fund_balance = system_state.storage_fund_total_object_storage_rebates;
1081        info!(
1082            "Total SUI amount in the network: {}, storage fund balance: {}, total storage rebate: {} at beginning of epoch {}",
1083            total_sui, storage_fund_balance, total_storage_rebate, system_state.epoch
1084        );
1085
1086        let imbalance = (storage_fund_balance as i64) - (total_storage_rebate as i64);
1087        self.metrics
1088            .sui_conservation_storage_fund
1089            .set(storage_fund_balance as i64);
1090        self.metrics
1091            .sui_conservation_storage_fund_imbalance
1092            .set(imbalance);
1093        self.metrics
1094            .sui_conservation_imbalance
1095            .set((total_sui as i128 - TOTAL_SUPPLY_MIST as i128) as i64);
1096
1097        if let Some(expected_imbalance) = self
1098            .perpetual_tables
1099            .expected_storage_fund_imbalance
1100            .get(&())
1101            .expect("DB read cannot fail")
1102        {
1103            fp_ensure!(
1104                imbalance == expected_imbalance,
1105                SuiError::from(
1106                    format!(
1107                        "Inconsistent state detected at epoch {}: total storage rebate: {}, storage fund balance: {}, expected imbalance: {}",
1108                        system_state.epoch, total_storage_rebate, storage_fund_balance, expected_imbalance
1109                    ).as_str()
1110                )
1111            );
1112        } else {
1113            self.perpetual_tables
1114                .expected_storage_fund_imbalance
1115                .insert(&(), &imbalance)
1116                .expect("DB write cannot fail");
1117        }
1118
1119        if let Some(expected_sui) = self
1120            .perpetual_tables
1121            .expected_network_sui_amount
1122            .get(&())
1123            .expect("DB read cannot fail")
1124        {
1125            fp_ensure!(
1126                total_sui == expected_sui,
1127                SuiError::from(
1128                    format!(
1129                        "Inconsistent state detected at epoch {}: total sui: {}, expecting {}",
1130                        system_state.epoch, total_sui, expected_sui
1131                    )
1132                    .as_str()
1133                )
1134            );
1135        } else {
1136            self.perpetual_tables
1137                .expected_network_sui_amount
1138                .insert(&(), &total_sui)
1139                .expect("DB write cannot fail");
1140        }
1141
1142        Ok(())
1143    }
1144
1145    /// This is a temporary method to be used when we enable simplified_unwrap_then_delete.
1146    /// It re-accumulates state hash for the new epoch if simplified_unwrap_then_delete is enabled.
1147    #[instrument(level = "error", skip_all)]
1148    pub fn maybe_reaccumulate_state_hash(
1149        &self,
1150        cur_epoch_store: &AuthorityPerEpochStore,
1151        new_protocol_version: ProtocolVersion,
1152    ) {
1153        let old_simplified_unwrap_then_delete = cur_epoch_store
1154            .protocol_config()
1155            .simplified_unwrap_then_delete();
1156        let new_simplified_unwrap_then_delete =
1157            ProtocolConfig::get_for_version(new_protocol_version, cur_epoch_store.get_chain())
1158                .simplified_unwrap_then_delete();
1159        // If in the new epoch the simplified_unwrap_then_delete is enabled for the first time,
1160        // we re-accumulate state root.
1161        let should_reaccumulate =
1162            !old_simplified_unwrap_then_delete && new_simplified_unwrap_then_delete;
1163        if !should_reaccumulate {
1164            return;
1165        }
1166        info!(
1167            "[Re-accumulate] simplified_unwrap_then_delete is enabled in the new protocol version, re-accumulating state hash"
1168        );
1169        let cur_time = Instant::now();
1170        std::thread::scope(|s| {
1171            let pending_tasks = FuturesUnordered::new();
1172            // Shard the object IDs into different ranges so that we can process them in parallel.
1173            // We divide the range into 2^BITS number of ranges. To do so we use the highest BITS bits
1174            // to mark the starting/ending point of the range. For example, when BITS = 5, we
1175            // divide the range into 32 ranges, and the first few ranges are:
1176            // 00000000_.... to 00000111_....
1177            // 00001000_.... to 00001111_....
1178            // 00010000_.... to 00010111_....
1179            // and etc.
1180            const BITS: u8 = 5;
1181            for index in 0u8..(1 << BITS) {
1182                pending_tasks.push(s.spawn(move || {
1183                    let mut id_bytes = [0; ObjectID::LENGTH];
1184                    id_bytes[0] = index << (8 - BITS);
1185                    let start_id = ObjectID::new(id_bytes);
1186
1187                    id_bytes[0] |= (1 << (8 - BITS)) - 1;
1188                    for element in id_bytes.iter_mut().skip(1) {
1189                        *element = u8::MAX;
1190                    }
1191                    let end_id = ObjectID::new(id_bytes);
1192
1193                    info!(
1194                        "[Re-accumulate] Scanning object ID range {:?}..{:?}",
1195                        start_id, end_id
1196                    );
1197                    let mut prev = (
1198                        ObjectKey::min_for_id(&ObjectID::ZERO),
1199                        StoreObjectWrapper::V1(StoreObject::Deleted),
1200                    );
1201                    let mut object_scanned: u64 = 0;
1202                    let mut wrapped_objects_to_remove = vec![];
1203                    for db_result in self.perpetual_tables.objects.safe_range_iter(
1204                        ObjectKey::min_for_id(&start_id)..=ObjectKey::max_for_id(&end_id),
1205                    ) {
1206                        match db_result {
1207                            Ok((object_key, object)) => {
1208                                object_scanned += 1;
1209                                if object_scanned.is_multiple_of(100000) {
1210                                    info!(
1211                                        "[Re-accumulate] Task {}: object scanned: {}",
1212                                        index, object_scanned,
1213                                    );
1214                                }
1215                                if matches!(prev.1.inner(), StoreObject::Wrapped)
1216                                    && object_key.0 != prev.0.0
1217                                {
1218                                    wrapped_objects_to_remove
1219                                        .push(WrappedObject::new(prev.0.0, prev.0.1));
1220                                }
1221
1222                                prev = (object_key, object);
1223                            }
1224                            Err(err) => {
1225                                warn!("Object iterator encounter RocksDB error {:?}", err);
1226                                return Err(err);
1227                            }
1228                        }
1229                    }
1230                    if matches!(prev.1.inner(), StoreObject::Wrapped) {
1231                        wrapped_objects_to_remove.push(WrappedObject::new(prev.0.0, prev.0.1));
1232                    }
1233                    info!(
1234                        "[Re-accumulate] Task {}: object scanned: {}, wrapped objects: {}",
1235                        index,
1236                        object_scanned,
1237                        wrapped_objects_to_remove.len(),
1238                    );
1239                    Ok((wrapped_objects_to_remove, object_scanned))
1240                }));
1241            }
1242            let (last_checkpoint_of_epoch, cur_accumulator) = self
1243                .get_root_state_hash_for_epoch(cur_epoch_store.epoch())
1244                .expect("read cannot fail")
1245                .expect("accumulator must exist");
1246            let (accumulator, total_objects_scanned, total_wrapped_objects) =
1247                pending_tasks.into_iter().fold(
1248                    (cur_accumulator, 0u64, 0usize),
1249                    |(mut accumulator, total_objects_scanned, total_wrapped_objects), task| {
1250                        let (wrapped_objects_to_remove, object_scanned) =
1251                            task.join().unwrap().unwrap();
1252                        accumulator.remove_all(
1253                            wrapped_objects_to_remove
1254                                .iter()
1255                                .map(|wrapped| bcs::to_bytes(wrapped).unwrap().to_vec())
1256                                .collect::<Vec<Vec<u8>>>(),
1257                        );
1258                        (
1259                            accumulator,
1260                            total_objects_scanned + object_scanned,
1261                            total_wrapped_objects + wrapped_objects_to_remove.len(),
1262                        )
1263                    },
1264                );
1265            info!(
1266                "[Re-accumulate] Total objects scanned: {}, total wrapped objects: {}",
1267                total_objects_scanned, total_wrapped_objects,
1268            );
1269            info!(
1270                "[Re-accumulate] New accumulator value: {:?}",
1271                accumulator.digest()
1272            );
1273            self.insert_state_hash_for_epoch(
1274                cur_epoch_store.epoch(),
1275                &last_checkpoint_of_epoch,
1276                &accumulator,
1277            )
1278            .unwrap();
1279        });
1280        info!(
1281            "[Re-accumulate] Re-accumulating took {}seconds",
1282            cur_time.elapsed().as_secs()
1283        );
1284    }
1285
1286    pub async fn prune_objects_and_compact_for_testing(
1287        &self,
1288        checkpoint_store: &Arc<CheckpointStore>,
1289    ) {
1290        let pruning_config = AuthorityStorePruningConfig {
1291            num_epochs_to_retain: 0,
1292            ..Default::default()
1293        };
1294        let _ = AuthorityStorePruner::prune_objects_for_eligible_epochs(
1295            &self.perpetual_tables,
1296            checkpoint_store,
1297            None,
1298            pruning_config,
1299            AuthorityStorePruningMetrics::new_for_test(),
1300            EPOCH_DURATION_MS_FOR_TESTING,
1301        )
1302        .await;
1303        let _ = AuthorityStorePruner::compact(&self.perpetual_tables);
1304    }
1305
1306    pub fn remove_executed_effects_for_testing(
1307        &self,
1308        tx_digest: &TransactionDigest,
1309    ) -> anyhow::Result<()> {
1310        let effects_digest = self.perpetual_tables.executed_effects.get(tx_digest)?;
1311        if let Some(effects_digest) = effects_digest {
1312            self.perpetual_tables.executed_effects.remove(tx_digest)?;
1313            self.perpetual_tables.effects.remove(&effects_digest)?;
1314        }
1315        Ok(())
1316    }
1317
1318    #[cfg(test)]
1319    pub async fn prune_objects_immediately_for_testing(
1320        &self,
1321        transaction_effects: Vec<TransactionEffects>,
1322    ) -> anyhow::Result<()> {
1323        let mut wb = self.perpetual_tables.objects.batch();
1324
1325        let mut object_keys_to_prune = vec![];
1326        for effects in &transaction_effects {
1327            for (object_id, seq_number) in effects.modified_at_versions() {
1328                info!("Pruning object {:?} version {:?}", object_id, seq_number);
1329                object_keys_to_prune.push(ObjectKey(object_id, seq_number));
1330            }
1331        }
1332
1333        wb.delete_batch(&self.perpetual_tables.objects, object_keys_to_prune)?;
1334        wb.write()?;
1335        Ok(())
1336    }
1337
1338    // Counts the number of versions exist in object store for `object_id`. This includes tombstone.
1339    #[cfg(msim)]
1340    pub fn count_object_versions(&self, object_id: ObjectID) -> usize {
1341        self.perpetual_tables
1342            .objects
1343            .safe_iter_with_bounds(
1344                Some(ObjectKey(object_id, VersionNumber::MIN)),
1345                Some(ObjectKey(object_id, VersionNumber::MAX)),
1346            )
1347            .collect::<Result<Vec<_>, _>>()
1348            .unwrap()
1349            .len()
1350    }
1351}
1352
1353impl GlobalStateHashStore for AuthorityStore {
1354    fn get_object_ref_prior_to_key_deprecated(
1355        &self,
1356        object_id: &ObjectID,
1357        version: VersionNumber,
1358    ) -> SuiResult<Option<ObjectRef>> {
1359        self.get_object_ref_prior_to_key(object_id, version)
1360    }
1361
1362    fn get_root_state_hash_for_epoch(
1363        &self,
1364        epoch: EpochId,
1365    ) -> SuiResult<Option<(CheckpointSequenceNumber, GlobalStateHash)>> {
1366        self.perpetual_tables
1367            .root_state_hash_by_epoch
1368            .get(&epoch)
1369            .map_err(Into::into)
1370    }
1371
1372    fn get_root_state_hash_for_highest_epoch(
1373        &self,
1374    ) -> SuiResult<Option<(EpochId, (CheckpointSequenceNumber, GlobalStateHash))>> {
1375        Ok(self
1376            .perpetual_tables
1377            .root_state_hash_by_epoch
1378            .reversed_safe_iter_with_bounds(None, None)?
1379            .next()
1380            .transpose()?)
1381    }
1382
1383    fn insert_state_hash_for_epoch(
1384        &self,
1385        epoch: EpochId,
1386        last_checkpoint_of_epoch: &CheckpointSequenceNumber,
1387        acc: &GlobalStateHash,
1388    ) -> SuiResult {
1389        self.perpetual_tables
1390            .root_state_hash_by_epoch
1391            .insert(&epoch, &(*last_checkpoint_of_epoch, acc.clone()))?;
1392        self.root_state_notify_read
1393            .notify(&epoch, &(*last_checkpoint_of_epoch, acc.clone()));
1394
1395        Ok(())
1396    }
1397
1398    fn iter_live_object_set(
1399        &self,
1400        include_wrapped_object: bool,
1401    ) -> Box<dyn Iterator<Item = LiveObject> + '_> {
1402        Box::new(
1403            self.perpetual_tables
1404                .iter_live_object_set(include_wrapped_object),
1405        )
1406    }
1407}
1408
1409impl ObjectStore for AuthorityStore {
1410    /// Read an object and return it, or Ok(None) if the object was not found.
1411    fn get_object(&self, object_id: &ObjectID) -> Option<Object> {
1412        self.perpetual_tables.as_ref().get_object(object_id)
1413    }
1414
1415    fn get_object_by_key(&self, object_id: &ObjectID, version: VersionNumber) -> Option<Object> {
1416        self.perpetual_tables.get_object_by_key(object_id, version)
1417    }
1418}
1419
1420/// A wrapper to make Orphan Rule happy
1421pub struct ResolverWrapper {
1422    pub resolver: Arc<dyn BackingPackageStore + Send + Sync>,
1423    pub metrics: Arc<ResolverMetrics>,
1424}
1425
1426impl ResolverWrapper {
1427    pub fn new(
1428        resolver: Arc<dyn BackingPackageStore + Send + Sync>,
1429        metrics: Arc<ResolverMetrics>,
1430    ) -> Self {
1431        metrics.module_cache_size.set(0);
1432        ResolverWrapper { resolver, metrics }
1433    }
1434
1435    fn inc_cache_size_gauge(&self) {
1436        // reset the gauge after a restart of the cache
1437        let current = self.metrics.module_cache_size.get();
1438        self.metrics.module_cache_size.set(current + 1);
1439    }
1440}
1441
1442impl ModuleResolver for ResolverWrapper {
1443    type Error = SuiError;
1444    fn get_module(&self, module_id: &ModuleId) -> Result<Option<Vec<u8>>, Self::Error> {
1445        self.inc_cache_size_gauge();
1446        get_module(&*self.resolver, module_id)
1447    }
1448
1449    fn get_packages_static<const N: usize>(
1450        &self,
1451        ids: [AccountAddress; N],
1452    ) -> Result<[Option<SerializedPackage>; N], Self::Error> {
1453        let mut packages = [const { None }; N];
1454        for (i, id) in ids.iter().enumerate() {
1455            packages[i] = get_package(&*self.resolver, &ObjectID::from(*id))?;
1456        }
1457        Ok(packages)
1458    }
1459
1460    fn get_packages<'a>(
1461        &self,
1462        ids: impl ExactSizeIterator<Item = &'a AccountAddress>,
1463    ) -> Result<Vec<Option<SerializedPackage>>, Self::Error> {
1464        ids.map(|id| get_package(&*self.resolver, &ObjectID::from(*id)))
1465            .collect()
1466    }
1467}
1468
1469#[cfg(test)]
1470pub type SuiLockResult = SuiResult<ObjectLockStatus>;
1471
1472#[cfg(test)]
1473#[derive(Debug, PartialEq, Eq)]
1474pub enum ObjectLockStatus {
1475    Initialized,
1476    LockedToTx { locked_by_tx: LockDetailsDeprecated },
1477    LockedAtDifferentVersion { locked_ref: ObjectRef },
1478}
1479
1480#[cfg(test)]
1481#[derive(Clone, Debug, PartialEq, Eq)]
1482pub struct LockDetailsV1Deprecated {
1483    pub epoch: EpochId,
1484    pub tx_digest: TransactionDigest,
1485}
1486
1487#[cfg(test)]
1488pub type LockDetailsDeprecated = LockDetailsV1Deprecated;