Skip to main content

sui_core/checkpoints/
mod.rs

1// Copyright (c) Mysten Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4pub(crate) mod causal_order;
5pub mod checkpoint_executor;
6mod checkpoint_output;
7mod metrics;
8
9use crate::accumulators::{self, AccumulatorSettlementTxBuilder};
10use crate::authority::AuthorityState;
11use crate::authority_client::{AuthorityAPI, make_network_authority_clients_with_network_config};
12use crate::checkpoints::causal_order::CausalOrder;
13use crate::checkpoints::checkpoint_output::CertifiedCheckpointOutput;
14pub use crate::checkpoints::checkpoint_output::{
15    CheckpointOutput, LogCheckpointOutput, SendCheckpointToStateSync, SubmitCheckpointToConsensus,
16};
17pub use crate::checkpoints::metrics::CheckpointMetrics;
18use crate::consensus_manager::ReplayWaiter;
19use crate::execution_cache::TransactionCacheRead;
20
21use crate::global_state_hasher::GlobalStateHasher;
22use crate::stake_aggregator::{InsertResult, MultiStakeAggregator};
23use consensus_core::CommitRef;
24use diffy::create_patch;
25use itertools::Itertools;
26use mysten_common::ZipDebugEqIteratorExt;
27use mysten_common::random::get_rng;
28use mysten_common::sync::notify_read::{CHECKPOINT_BUILDER_NOTIFY_READ_TASK_NAME, NotifyRead};
29use mysten_common::{assert_reachable, debug_fatal, fatal, in_antithesis};
30use mysten_metrics::{MonitoredFutureExt, monitored_scope, spawn_monitored_task};
31use parking_lot::Mutex;
32use pin_project_lite::pin_project;
33use serde::{Deserialize, Serialize};
34use sui_macros::fail_point_arg;
35use sui_network::default_mysten_network_config;
36use sui_types::accumulator_metadata;
37use sui_types::base_types::{ConciseableName, SequenceNumber};
38use sui_types::execution::ExecutionTimeObservationKey;
39use sui_types::messages_checkpoint::{
40    CheckpointArtifacts, CheckpointCommitment, VersionedFullCheckpointContents,
41};
42use sui_types::sui_system_state::epoch_start_sui_system_state::EpochStartSystemStateTrait;
43use tokio::sync::{mpsc, watch};
44#[cfg(not(tidehunter))]
45use typed_store::rocks::{DBOptions, ReadWriteOptions, default_db_options};
46
47use crate::authority::authority_per_epoch_store::AuthorityPerEpochStore;
48use crate::authority::authority_store_pruner::PrunerWatermarks;
49use crate::consensus_handler::SequencedConsensusTransactionKey;
50use rand::seq::SliceRandom;
51use std::collections::{BTreeMap, HashMap, HashSet};
52use std::fs::File;
53use std::future::Future;
54use std::io::Write;
55use std::path::Path;
56use std::pin::Pin;
57use std::sync::Arc;
58use std::sync::OnceLock;
59use std::sync::Weak;
60use std::task::{Context, Poll};
61use std::time::{Duration, SystemTime};
62use sui_protocol_config::ProtocolVersion;
63use sui_types::base_types::{AuthorityName, EpochId, TransactionDigest};
64use sui_types::committee::StakeUnit;
65use sui_types::crypto::AuthorityStrongQuorumSignInfo;
66use sui_types::digests::{
67    CheckpointContentsDigest, CheckpointDigest, Digest, TransactionEffectsDigest,
68};
69use sui_types::effects::{TransactionEffects, TransactionEffectsAPI};
70use sui_types::error::{SuiErrorKind, SuiResult};
71use sui_types::gas::GasCostSummary;
72use sui_types::message_envelope::Message;
73use sui_types::messages_checkpoint::{
74    CertifiedCheckpointSummary, CheckpointContents, CheckpointResponseV2, CheckpointSequenceNumber,
75    CheckpointSignatureMessage, CheckpointSummary, CheckpointSummaryResponse, CheckpointTimestamp,
76    EndOfEpochData, FullCheckpointContents, TrustedCheckpoint, VerifiedCheckpoint,
77    VerifiedCheckpointContents,
78};
79use sui_types::messages_checkpoint::{CheckpointRequestV2, SignedCheckpointSummary};
80use sui_types::messages_consensus::ConsensusTransactionKey;
81use sui_types::signature::GenericSignature;
82use sui_types::sui_system_state::{SuiSystemState, SuiSystemStateTrait};
83use sui_types::transaction::{
84    TransactionDataAPI, TransactionKey, TransactionKind, VerifiedTransaction,
85};
86use tokio::sync::Notify;
87use tracing::{debug, error, info, instrument, trace, warn};
88use typed_store::DBMapUtils;
89use typed_store::Map;
90use typed_store::{
91    TypedStoreError,
92    rocks::{DBMap, MetricConf},
93};
94
95const TRANSACTION_FORK_DETECTED_KEY: u8 = 0;
96const CHECKPOINT_FORK_DETECTED_KEY: u8 = 0;
97
98pub type CheckpointHeight = u64;
99
100#[derive(Clone, Debug, Serialize, Deserialize)]
101pub struct TransactionForkInfo {
102    pub tx_digest: TransactionDigest,
103    pub expected_effects_digest: TransactionEffectsDigest,
104    pub actual_effects_digest: TransactionEffectsDigest,
105    /// Sequence of the certified checkpoint whose contents supplied `expected_effects_digest`,
106    /// when the fork was detected while executing a synced checkpoint. None when the
107    /// expectation came from a non-certified source (e.g. this validator's own previously
108    /// signed effects), in which case automatic recovery is never allowed: nothing proves the
109    /// network certified the outcome.
110    pub certified_checkpoint_seq: Option<CheckpointSequenceNumber>,
111    /// The binary version that produced the fork (the node sets the store's version at startup,
112    /// before any fork can be recorded). Automatic recovery refuses to clear a marker carrying
113    /// the currently running version: re-deriving with the binary that forked would
114    /// deterministically fork again, so the node hangs until a corrected binary is deployed.
115    pub binary_version: String,
116}
117
118#[derive(Clone, Debug, Serialize, Deserialize)]
119pub struct CheckpointForkInfo {
120    pub checkpoint_seq: CheckpointSequenceNumber,
121    /// Digest of the locally computed checkpoint that diverged.
122    pub checkpoint_digest: CheckpointDigest,
123    /// Digest of the certified checkpoint the local one diverged from, when the fork was
124    /// detected against a certified checkpoint. None when the builder re-derived a previously
125    /// computed sequence with a different result (e.g. after a binary upgrade): nothing proves
126    /// the network certified this sequence, so automatic recovery is never allowed.
127    pub certified_checkpoint_digest: Option<CheckpointDigest>,
128    /// See [`TransactionForkInfo::binary_version`].
129    pub binary_version: String,
130}
131
132pub struct EpochStats {
133    pub checkpoint_count: u64,
134    pub transaction_count: u64,
135    pub total_gas_reward: u64,
136}
137
138#[derive(Clone, Debug)]
139pub struct PendingCheckpointInfo {
140    pub timestamp_ms: CheckpointTimestamp,
141    pub last_of_epoch: bool,
142    pub checkpoint_height: CheckpointHeight,
143    // Consensus commit ref and rejected transactions digest which corresponds to this checkpoint.
144    pub consensus_commit_ref: CommitRef,
145    pub rejected_transactions_digest: Digest,
146    // Pre-assigned checkpoint sequence number from consensus handler.
147    pub checkpoint_seq: CheckpointSequenceNumber,
148}
149
150#[derive(Clone, Debug, Default)]
151pub struct CheckpointRoots {
152    pub tx_roots: Vec<TransactionKey>,
153    pub settlement_root: Option<TransactionKey>,
154    pub height: CheckpointHeight,
155}
156
157/// Consensus commits are merged and split into PendingCheckpoints in ConsensusHandler.
158/// Each CheckpointRoots represents a group of transactions settled together.
159#[derive(Clone, Debug)]
160pub struct PendingCheckpoint {
161    pub roots: Vec<CheckpointRoots>,
162    pub details: PendingCheckpointInfo,
163}
164
165#[derive(Clone, Debug, Serialize, Deserialize)]
166pub struct BuilderCheckpointSummary {
167    pub summary: CheckpointSummary,
168    // Height at which this checkpoint summary was built. None for genesis checkpoint
169    pub checkpoint_height: Option<CheckpointHeight>,
170    // Always 0: each height now maps to exactly one checkpoint. Kept for DB format
171    // compatibility; old rows may contain nonzero values from builder-side splitting.
172    pub position_in_commit: usize,
173}
174
175#[derive(DBMapUtils)]
176#[cfg_attr(tidehunter, tidehunter)]
177pub struct CheckpointStoreTables {
178    /// Maps checkpoint contents digest to checkpoint contents
179    pub(crate) checkpoint_content: DBMap<CheckpointContentsDigest, CheckpointContents>,
180
181    /// Maps checkpoint contents digest to checkpoint sequence number
182    pub(crate) checkpoint_sequence_by_contents_digest:
183        DBMap<CheckpointContentsDigest, CheckpointSequenceNumber>,
184
185    /// Stores entire checkpoint contents from state sync, indexed by sequence number, for
186    /// efficient reads of full checkpoints. Entries from this table are deleted after state
187    /// accumulation has completed.
188    #[default_options_override_fn = "full_checkpoint_content_table_default_config"]
189    // TODO: Once the switch to `full_checkpoint_content_v2` is fully active on mainnet,
190    // deprecate this table (and remove when possible).
191    full_checkpoint_content: DBMap<CheckpointSequenceNumber, FullCheckpointContents>,
192
193    /// Stores certified checkpoints
194    pub(crate) certified_checkpoints: DBMap<CheckpointSequenceNumber, TrustedCheckpoint>,
195    /// Map from checkpoint digest to certified checkpoint
196    pub(crate) checkpoint_by_digest: DBMap<CheckpointDigest, TrustedCheckpoint>,
197
198    /// Store locally computed checkpoint summaries so that we can detect forks and log useful
199    /// information. Can be pruned as soon as we verify that we are in agreement with the latest
200    /// certified checkpoint.
201    pub(crate) locally_computed_checkpoints: DBMap<CheckpointSequenceNumber, CheckpointSummary>,
202
203    /// A map from epoch ID to the sequence number of the last checkpoint in that epoch.
204    epoch_last_checkpoint_map: DBMap<EpochId, CheckpointSequenceNumber>,
205
206    /// Watermarks used to determine the highest verified, fully synced, and
207    /// fully executed checkpoints
208    pub(crate) watermarks: DBMap<CheckpointWatermark, (CheckpointSequenceNumber, CheckpointDigest)>,
209
210    /// Stores transaction fork detection information, including the certified checkpoint (if
211    /// any) that supplied the expected effects digest and the binary version that forked.
212    /// Automatic fork recovery is gated on both: it may only proceed once the network has
213    /// certified the forked outcome, and only under a different binary version.
214    pub(crate) transaction_fork_detected_v2: DBMap<u8, TransactionForkInfo>,
215    #[default_options_override_fn = "full_checkpoint_content_table_default_config"]
216    full_checkpoint_content_v2: DBMap<CheckpointSequenceNumber, VersionedFullCheckpointContents>,
217
218    /// Stores checkpoint fork detection information, including the binary version that forked.
219    pub(crate) checkpoint_fork_detected: DBMap<u8, CheckpointForkInfo>,
220}
221
222#[cfg(not(tidehunter))]
223fn full_checkpoint_content_table_default_config() -> DBOptions {
224    DBOptions {
225        options: default_db_options().options,
226        // We have seen potential data corruption issues in this table after forced shutdowns
227        // so we enable value hash logging to help with debugging.
228        // TODO: remove this once we have a better understanding of the root cause.
229        rw_options: ReadWriteOptions::default().set_log_value_hash(true),
230    }
231}
232
233impl CheckpointStoreTables {
234    #[cfg(not(tidehunter))]
235    pub fn new(path: &Path, metric_name: &'static str, _: Arc<PrunerWatermarks>) -> Self {
236        Self::open_tables_read_write(path.to_path_buf(), MetricConf::new(metric_name), None, None)
237    }
238
239    #[cfg(tidehunter)]
240    pub fn new(
241        path: &Path,
242        metric_name: &'static str,
243        pruner_watermarks: Arc<PrunerWatermarks>,
244    ) -> Self {
245        tracing::warn!("Checkpoint DB using tidehunter");
246        use crate::authority::authority_store_pruner::apply_relocation_filter;
247        use typed_store::tidehunter_util::{
248            Decision, KeySpaceConfig, KeyType, ThConfig, default_cells_per_mutex,
249            default_max_dirty_keys, default_mutex_count, default_value_cache_size,
250        };
251        let mutexes = default_mutex_count();
252        let u64_sequence_key = KeyType::from_prefix_bits(6 * 8);
253        let override_dirty_keys_config = KeySpaceConfig::new()
254            .with_max_dirty_keys(16 * default_max_dirty_keys())
255            .with_value_cache_size(default_value_cache_size());
256        let config_u64 = ThConfig::new_with_config(
257            8,
258            mutexes,
259            u64_sequence_key,
260            override_dirty_keys_config.clone(),
261        );
262        let digest_config = ThConfig::new_with_rm_prefix(
263            32,
264            mutexes,
265            KeyType::uniform(default_cells_per_mutex()),
266            KeySpaceConfig::default(),
267            vec![0, 0, 0, 0, 0, 0, 0, 32],
268        );
269        let watermarks_config = KeySpaceConfig::new()
270            .with_value_cache_size(10)
271            .disable_unload();
272        let lru_config = KeySpaceConfig::new().with_value_cache_size(100);
273        let configs = vec![
274            (
275                "checkpoint_content",
276                digest_config.clone().with_config(
277                    KeySpaceConfig::new().with_relocation_filter(|_, _| Decision::Remove),
278                ),
279            ),
280            (
281                "checkpoint_sequence_by_contents_digest",
282                digest_config.clone().with_config(apply_relocation_filter(
283                    KeySpaceConfig::default(),
284                    pruner_watermarks.checkpoint_id.clone(),
285                    |sequence_number: CheckpointSequenceNumber| sequence_number,
286                    false,
287                )),
288            ),
289            (
290                "full_checkpoint_content",
291                config_u64.clone().with_config(apply_relocation_filter(
292                    override_dirty_keys_config.clone(),
293                    pruner_watermarks.checkpoint_id.clone(),
294                    |sequence_number: CheckpointSequenceNumber| sequence_number,
295                    true,
296                )),
297            ),
298            ("certified_checkpoints", config_u64.clone()),
299            (
300                "checkpoint_by_digest",
301                digest_config.clone().with_config(apply_relocation_filter(
302                    lru_config,
303                    pruner_watermarks.epoch_id.clone(),
304                    |checkpoint: TrustedCheckpoint| checkpoint.inner().epoch,
305                    false,
306                )),
307            ),
308            (
309                "locally_computed_checkpoints",
310                config_u64.clone().with_config(apply_relocation_filter(
311                    override_dirty_keys_config.clone(),
312                    pruner_watermarks.checkpoint_id.clone(),
313                    |checkpoint_id: CheckpointSequenceNumber| checkpoint_id,
314                    true,
315                )),
316            ),
317            ("epoch_last_checkpoint_map", config_u64.clone()),
318            (
319                "watermarks",
320                ThConfig::new_with_config(4, 1, KeyType::uniform(1), watermarks_config.clone()),
321            ),
322            (
323                "transaction_fork_detected_v2",
324                ThConfig::new_with_config(
325                    1,
326                    1,
327                    KeyType::uniform(1),
328                    watermarks_config
329                        .clone()
330                        .with_relocation_filter(|_, _| Decision::Remove),
331                ),
332            ),
333            (
334                "checkpoint_fork_detected",
335                ThConfig::new_with_config(
336                    1,
337                    1,
338                    KeyType::uniform(1),
339                    watermarks_config.with_relocation_filter(|_, _| Decision::Remove),
340                ),
341            ),
342            (
343                "full_checkpoint_content_v2",
344                config_u64.clone().with_config(apply_relocation_filter(
345                    override_dirty_keys_config.clone(),
346                    pruner_watermarks.checkpoint_id.clone(),
347                    |sequence_number: CheckpointSequenceNumber| sequence_number,
348                    true,
349                )),
350            ),
351        ];
352        Self::open_tables_read_write(
353            path.to_path_buf(),
354            MetricConf::new(metric_name),
355            configs
356                .into_iter()
357                .map(|(cf, config)| (cf.to_string(), config))
358                .collect(),
359        )
360    }
361
362    #[cfg(not(tidehunter))]
363    pub fn open_readonly(path: &Path) -> CheckpointStoreTablesReadOnly {
364        Self::get_read_only_handle(
365            path.to_path_buf(),
366            None,
367            None,
368            MetricConf::new("checkpoint_readonly"),
369        )
370    }
371
372    #[cfg(tidehunter)]
373    pub fn open_readonly(path: &Path) -> Self {
374        Self::new(path, "checkpoint", Arc::new(PrunerWatermarks::default()))
375    }
376}
377
378pub struct CheckpointStore {
379    pub(crate) tables: CheckpointStoreTables,
380    synced_checkpoint_notify_read: NotifyRead<CheckpointSequenceNumber, VerifiedCheckpoint>,
381    executed_checkpoint_notify_read: NotifyRead<CheckpointSequenceNumber, VerifiedCheckpoint>,
382    /// The running binary's version, set once at node startup. Embedded in fork markers so
383    /// recovery can refuse to clear a fork with the same binary version that produced it.
384    binary_version: OnceLock<String>,
385}
386
387impl CheckpointStore {
388    pub fn new(path: &Path, pruner_watermarks: Arc<PrunerWatermarks>) -> Arc<Self> {
389        let tables = CheckpointStoreTables::new(path, "checkpoint", pruner_watermarks);
390        Arc::new(Self {
391            tables,
392            synced_checkpoint_notify_read: NotifyRead::new(),
393            executed_checkpoint_notify_read: NotifyRead::new(),
394            binary_version: OnceLock::new(),
395        })
396    }
397
398    pub fn new_for_tests() -> Arc<Self> {
399        let ckpt_dir = mysten_common::tempdir().unwrap();
400        CheckpointStore::new(ckpt_dir.path(), Arc::new(PrunerWatermarks::default()))
401    }
402
403    pub fn new_for_db_checkpoint_handler(path: &Path) -> Arc<Self> {
404        let tables = CheckpointStoreTables::new(
405            path,
406            "db_checkpoint",
407            Arc::new(PrunerWatermarks::default()),
408        );
409        Arc::new(Self {
410            tables,
411            synced_checkpoint_notify_read: NotifyRead::new(),
412            executed_checkpoint_notify_read: NotifyRead::new(),
413            binary_version: OnceLock::new(),
414        })
415    }
416
417    pub fn set_binary_version(&self, version: &str) {
418        let _ = self.binary_version.set(version.to_string());
419    }
420
421    #[cfg(not(tidehunter))]
422    pub fn open_readonly(path: &Path) -> CheckpointStoreTablesReadOnly {
423        CheckpointStoreTables::open_readonly(path)
424    }
425
426    #[cfg(tidehunter)]
427    pub fn open_readonly(path: &Path) -> CheckpointStoreTables {
428        CheckpointStoreTables::open_readonly(path)
429    }
430
431    #[instrument(level = "info", skip_all)]
432    pub fn insert_genesis_checkpoint(
433        &self,
434        checkpoint: VerifiedCheckpoint,
435        contents: CheckpointContents,
436        epoch_store: &AuthorityPerEpochStore,
437    ) {
438        assert_eq!(
439            checkpoint.epoch(),
440            0,
441            "can't call insert_genesis_checkpoint with a checkpoint not in epoch 0"
442        );
443        assert_eq!(
444            *checkpoint.sequence_number(),
445            0,
446            "can't call insert_genesis_checkpoint with a checkpoint that doesn't have a sequence number of 0"
447        );
448
449        // Only insert the genesis checkpoint if the DB is empty and doesn't have it already
450        match self.get_checkpoint_by_sequence_number(0).unwrap() {
451            Some(existing_checkpoint) => {
452                assert_eq!(existing_checkpoint.digest(), checkpoint.digest())
453            }
454            None => {
455                if epoch_store.epoch() == checkpoint.epoch {
456                    epoch_store
457                        .put_genesis_checkpoint_in_builder(checkpoint.data())
458                        .unwrap();
459                } else {
460                    debug!(
461                        validator_epoch =% epoch_store.epoch(),
462                        genesis_epoch =% checkpoint.epoch(),
463                        "Not inserting checkpoint builder data for genesis checkpoint",
464                    );
465                }
466                self.insert_checkpoint_contents(contents).unwrap();
467                self.insert_verified_checkpoint(&checkpoint).unwrap();
468                self.update_highest_synced_checkpoint(&checkpoint).unwrap();
469            }
470        }
471    }
472
473    pub fn get_checkpoint_by_digest(
474        &self,
475        digest: &CheckpointDigest,
476    ) -> Result<Option<VerifiedCheckpoint>, TypedStoreError> {
477        self.tables
478            .checkpoint_by_digest
479            .get(digest)
480            .map(|maybe_checkpoint| maybe_checkpoint.map(|c| c.into()))
481    }
482
483    pub fn get_checkpoint_by_sequence_number(
484        &self,
485        sequence_number: CheckpointSequenceNumber,
486    ) -> Result<Option<VerifiedCheckpoint>, TypedStoreError> {
487        self.tables
488            .certified_checkpoints
489            .get(&sequence_number)
490            .map(|maybe_checkpoint| maybe_checkpoint.map(|c| c.into()))
491    }
492
493    pub fn get_locally_computed_checkpoint(
494        &self,
495        sequence_number: CheckpointSequenceNumber,
496    ) -> Result<Option<CheckpointSummary>, TypedStoreError> {
497        self.tables
498            .locally_computed_checkpoints
499            .get(&sequence_number)
500    }
501
502    pub fn multi_get_locally_computed_checkpoints(
503        &self,
504        sequence_numbers: &[CheckpointSequenceNumber],
505    ) -> Result<Vec<Option<CheckpointSummary>>, TypedStoreError> {
506        let checkpoints = self
507            .tables
508            .locally_computed_checkpoints
509            .multi_get(sequence_numbers)?;
510
511        Ok(checkpoints)
512    }
513
514    pub fn get_sequence_number_by_contents_digest(
515        &self,
516        digest: &CheckpointContentsDigest,
517    ) -> Result<Option<CheckpointSequenceNumber>, TypedStoreError> {
518        self.tables
519            .checkpoint_sequence_by_contents_digest
520            .get(digest)
521    }
522
523    pub fn delete_contents_digest_sequence_number_mapping(
524        &self,
525        digest: &CheckpointContentsDigest,
526    ) -> Result<(), TypedStoreError> {
527        self.tables
528            .checkpoint_sequence_by_contents_digest
529            .remove(digest)
530    }
531
532    pub fn get_latest_certified_checkpoint(
533        &self,
534    ) -> Result<Option<VerifiedCheckpoint>, TypedStoreError> {
535        Ok(self
536            .tables
537            .certified_checkpoints
538            .reversed_safe_iter_with_bounds(None, None)?
539            .next()
540            .transpose()?
541            .map(|(_, v)| v.into()))
542    }
543
544    pub fn get_latest_locally_computed_checkpoint(
545        &self,
546    ) -> Result<Option<CheckpointSummary>, TypedStoreError> {
547        Ok(self
548            .tables
549            .locally_computed_checkpoints
550            .reversed_safe_iter_with_bounds(None, None)?
551            .next()
552            .transpose()?
553            .map(|(_, v)| v))
554    }
555
556    pub fn multi_get_checkpoint_by_sequence_number(
557        &self,
558        sequence_numbers: &[CheckpointSequenceNumber],
559    ) -> Result<Vec<Option<VerifiedCheckpoint>>, TypedStoreError> {
560        let checkpoints = self
561            .tables
562            .certified_checkpoints
563            .multi_get(sequence_numbers)?
564            .into_iter()
565            .map(|maybe_checkpoint| maybe_checkpoint.map(|c| c.into()))
566            .collect();
567
568        Ok(checkpoints)
569    }
570
571    pub fn multi_get_checkpoint_content(
572        &self,
573        contents_digest: &[CheckpointContentsDigest],
574    ) -> Result<Vec<Option<CheckpointContents>>, TypedStoreError> {
575        self.tables.checkpoint_content.multi_get(contents_digest)
576    }
577
578    pub fn get_highest_verified_checkpoint(
579        &self,
580    ) -> Result<Option<VerifiedCheckpoint>, TypedStoreError> {
581        let highest_verified = if let Some(highest_verified) = self
582            .tables
583            .watermarks
584            .get(&CheckpointWatermark::HighestVerified)?
585        {
586            highest_verified
587        } else {
588            return Ok(None);
589        };
590        self.get_checkpoint_by_digest(&highest_verified.1)
591    }
592
593    pub fn get_highest_synced_checkpoint(
594        &self,
595    ) -> Result<Option<VerifiedCheckpoint>, TypedStoreError> {
596        let highest_synced = if let Some(highest_synced) = self
597            .tables
598            .watermarks
599            .get(&CheckpointWatermark::HighestSynced)?
600        {
601            highest_synced
602        } else {
603            return Ok(None);
604        };
605        self.get_checkpoint_by_digest(&highest_synced.1)
606    }
607
608    pub fn get_highest_synced_checkpoint_seq_number(
609        &self,
610    ) -> Result<Option<CheckpointSequenceNumber>, TypedStoreError> {
611        if let Some(highest_synced) = self
612            .tables
613            .watermarks
614            .get(&CheckpointWatermark::HighestSynced)?
615        {
616            Ok(Some(highest_synced.0))
617        } else {
618            Ok(None)
619        }
620    }
621
622    pub fn get_highest_executed_checkpoint_seq_number(
623        &self,
624    ) -> Result<Option<CheckpointSequenceNumber>, TypedStoreError> {
625        if let Some(highest_executed) = self
626            .tables
627            .watermarks
628            .get(&CheckpointWatermark::HighestExecuted)?
629        {
630            Ok(Some(highest_executed.0))
631        } else {
632            Ok(None)
633        }
634    }
635
636    pub fn get_highest_executed_checkpoint(
637        &self,
638    ) -> Result<Option<VerifiedCheckpoint>, TypedStoreError> {
639        let highest_executed = if let Some(highest_executed) = self
640            .tables
641            .watermarks
642            .get(&CheckpointWatermark::HighestExecuted)?
643        {
644            highest_executed
645        } else {
646            return Ok(None);
647        };
648        self.get_checkpoint_by_digest(&highest_executed.1)
649    }
650
651    pub fn get_highest_pruned_checkpoint_seq_number(
652        &self,
653    ) -> Result<Option<CheckpointSequenceNumber>, TypedStoreError> {
654        self.tables
655            .watermarks
656            .get(&CheckpointWatermark::HighestPruned)
657            .map(|watermark| watermark.map(|w| w.0))
658    }
659
660    pub fn get_checkpoint_contents(
661        &self,
662        digest: &CheckpointContentsDigest,
663    ) -> Result<Option<CheckpointContents>, TypedStoreError> {
664        self.tables.checkpoint_content.get(digest)
665    }
666
667    pub fn get_full_checkpoint_contents_by_sequence_number(
668        &self,
669        seq: CheckpointSequenceNumber,
670    ) -> Result<Option<VersionedFullCheckpointContents>, TypedStoreError> {
671        self.tables.full_checkpoint_content_v2.get(&seq)
672    }
673
674    fn prune_local_summaries(&self) -> SuiResult {
675        if let Some((last_local_summary, _)) = self
676            .tables
677            .locally_computed_checkpoints
678            .reversed_safe_iter_with_bounds(None, None)?
679            .next()
680            .transpose()?
681        {
682            let mut batch = self.tables.locally_computed_checkpoints.batch();
683            batch.schedule_delete_range(
684                &self.tables.locally_computed_checkpoints,
685                &0,
686                &last_local_summary,
687            )?;
688            batch.write()?;
689            info!("Pruned local summaries up to {:?}", last_local_summary);
690        }
691        Ok(())
692    }
693
694    pub fn clear_locally_computed_checkpoints_from(
695        &self,
696        from_seq: CheckpointSequenceNumber,
697    ) -> SuiResult {
698        let keys: Vec<_> = self
699            .tables
700            .locally_computed_checkpoints
701            .safe_iter_with_bounds(Some(from_seq), None)
702            .map(|r| r.map(|(k, _)| k))
703            .collect::<Result<_, _>>()?;
704        if let Some(&last_local_summary) = keys.last() {
705            let mut batch = self.tables.locally_computed_checkpoints.batch();
706            batch
707                .delete_batch(&self.tables.locally_computed_checkpoints, keys.iter())
708                .expect("Failed to delete locally computed checkpoints");
709            batch
710                .write()
711                .expect("Failed to delete locally computed checkpoints");
712            warn!(
713                from_seq,
714                last_local_summary,
715                "Cleared locally_computed_checkpoints from {} (inclusive) through {} (inclusive)",
716                from_seq,
717                last_local_summary
718            );
719        }
720        Ok(())
721    }
722
723    fn check_for_checkpoint_fork(
724        &self,
725        local_checkpoint: &CheckpointSummary,
726        verified_checkpoint: &VerifiedCheckpoint,
727    ) {
728        if local_checkpoint != verified_checkpoint.data() {
729            let verified_contents = self
730                .get_checkpoint_contents(&verified_checkpoint.content_digest)
731                .map(|opt_contents| {
732                    opt_contents
733                        .map(|contents| format!("{:?}", contents))
734                        .unwrap_or_else(|| {
735                            format!(
736                                "Verified checkpoint contents not found, digest: {:?}",
737                                verified_checkpoint.content_digest,
738                            )
739                        })
740                })
741                .map_err(|e| {
742                    format!(
743                        "Failed to get verified checkpoint contents, digest: {:?} error: {:?}",
744                        verified_checkpoint.content_digest, e
745                    )
746                })
747                .unwrap_or_else(|err_msg| err_msg);
748
749            let local_contents = self
750                .get_checkpoint_contents(&local_checkpoint.content_digest)
751                .map(|opt_contents| {
752                    opt_contents
753                        .map(|contents| format!("{:?}", contents))
754                        .unwrap_or_else(|| {
755                            format!(
756                                "Local checkpoint contents not found, digest: {:?}",
757                                local_checkpoint.content_digest
758                            )
759                        })
760                })
761                .map_err(|e| {
762                    format!(
763                        "Failed to get local checkpoint contents, digest: {:?} error: {:?}",
764                        local_checkpoint.content_digest, e
765                    )
766                })
767                .unwrap_or_else(|err_msg| err_msg);
768
769            // checkpoint contents may be too large for panic message.
770            error!(
771                verified_checkpoint = ?verified_checkpoint.data(),
772                ?verified_contents,
773                ?local_checkpoint,
774                ?local_contents,
775                "Local checkpoint fork detected!",
776            );
777
778            // Record the fork in the database before crashing
779            if let Err(e) = self.record_checkpoint_fork_detected(
780                *local_checkpoint.sequence_number(),
781                local_checkpoint.digest(),
782                Some(*verified_checkpoint.digest()),
783            ) {
784                error!("Failed to record checkpoint fork in database: {:?}", e);
785            }
786
787            fail_point_arg!(
788                "kill_checkpoint_fork_node",
789                |checkpoint_overrides: std::sync::Arc<
790                    std::sync::Mutex<std::collections::BTreeMap<u64, String>>,
791                >| {
792                    #[cfg(msim)]
793                    {
794                        if let Ok(mut overrides) = checkpoint_overrides.lock() {
795                            overrides.insert(
796                                local_checkpoint.sequence_number,
797                                verified_checkpoint.digest().to_string(),
798                            );
799                        }
800                        tracing::error!(
801                            fatal = true,
802                            "Fork recovery test: killing node due to checkpoint fork for sequence number: {}, using verified digest: {}",
803                            local_checkpoint.sequence_number(),
804                            verified_checkpoint.digest()
805                        );
806                        sui_simulator::task::shutdown_current_node();
807                    }
808                }
809            );
810
811            fatal!(
812                "Local checkpoint fork detected for sequence number: {}",
813                local_checkpoint.sequence_number()
814            );
815        }
816    }
817
818    // Called by consensus (ConsensusAggregator).
819    // Different from `insert_verified_checkpoint`, it does not touch
820    // the highest_verified_checkpoint watermark such that state sync
821    // will have a chance to process this checkpoint and perform some
822    // state-sync only things.
823    pub fn insert_certified_checkpoint(
824        &self,
825        checkpoint: &VerifiedCheckpoint,
826    ) -> Result<(), TypedStoreError> {
827        debug!(
828            checkpoint_seq = checkpoint.sequence_number(),
829            "Inserting certified checkpoint",
830        );
831        let mut batch = self.tables.certified_checkpoints.batch();
832        batch
833            .insert_batch(
834                &self.tables.certified_checkpoints,
835                [(checkpoint.sequence_number(), checkpoint.serializable_ref())],
836            )?
837            .insert_batch(
838                &self.tables.checkpoint_by_digest,
839                [(checkpoint.digest(), checkpoint.serializable_ref())],
840            )?;
841        if checkpoint.next_epoch_committee().is_some() {
842            batch.insert_batch(
843                &self.tables.epoch_last_checkpoint_map,
844                [(&checkpoint.epoch(), checkpoint.sequence_number())],
845            )?;
846        }
847        batch.write()?;
848
849        if let Some(local_checkpoint) = self
850            .tables
851            .locally_computed_checkpoints
852            .get(checkpoint.sequence_number())?
853        {
854            self.check_for_checkpoint_fork(&local_checkpoint, checkpoint);
855        }
856
857        Ok(())
858    }
859
860    // Called by state sync, apart from inserting the checkpoint and updating
861    // related tables, it also bumps the highest_verified_checkpoint watermark.
862    #[instrument(level = "debug", skip_all)]
863    pub fn insert_verified_checkpoint(
864        &self,
865        checkpoint: &VerifiedCheckpoint,
866    ) -> Result<(), TypedStoreError> {
867        self.insert_certified_checkpoint(checkpoint)?;
868        self.update_highest_verified_checkpoint(checkpoint)
869    }
870
871    pub fn update_highest_verified_checkpoint(
872        &self,
873        checkpoint: &VerifiedCheckpoint,
874    ) -> Result<(), TypedStoreError> {
875        if Some(*checkpoint.sequence_number())
876            > self
877                .get_highest_verified_checkpoint()?
878                .map(|x| *x.sequence_number())
879        {
880            debug!(
881                checkpoint_seq = checkpoint.sequence_number(),
882                "Updating highest verified checkpoint",
883            );
884            self.tables.watermarks.insert(
885                &CheckpointWatermark::HighestVerified,
886                &(*checkpoint.sequence_number(), *checkpoint.digest()),
887            )?;
888        }
889
890        Ok(())
891    }
892
893    pub fn update_highest_synced_checkpoint(
894        &self,
895        checkpoint: &VerifiedCheckpoint,
896    ) -> Result<(), TypedStoreError> {
897        let seq = *checkpoint.sequence_number();
898        debug!(checkpoint_seq = seq, "Updating highest synced checkpoint",);
899        self.tables.watermarks.insert(
900            &CheckpointWatermark::HighestSynced,
901            &(seq, *checkpoint.digest()),
902        )?;
903        self.synced_checkpoint_notify_read.notify(&seq, checkpoint);
904        Ok(())
905    }
906
907    async fn notify_read_checkpoint_watermark<F>(
908        &self,
909        notify_read: &NotifyRead<CheckpointSequenceNumber, VerifiedCheckpoint>,
910        seq: CheckpointSequenceNumber,
911        get_watermark: F,
912    ) -> VerifiedCheckpoint
913    where
914        F: Fn() -> Option<CheckpointSequenceNumber>,
915    {
916        notify_read
917            .read("notify_read_checkpoint_watermark", &[seq], |seqs| {
918                let seq = seqs[0];
919                let Some(highest) = get_watermark() else {
920                    return vec![None];
921                };
922                if highest < seq {
923                    return vec![None];
924                }
925                let checkpoint = self
926                    .get_checkpoint_by_sequence_number(seq)
927                    .expect("db error")
928                    .expect("checkpoint not found");
929                vec![Some(checkpoint)]
930            })
931            .await
932            .into_iter()
933            .next()
934            .unwrap()
935    }
936
937    pub async fn notify_read_synced_checkpoint(
938        &self,
939        seq: CheckpointSequenceNumber,
940    ) -> VerifiedCheckpoint {
941        self.notify_read_checkpoint_watermark(&self.synced_checkpoint_notify_read, seq, || {
942            self.get_highest_synced_checkpoint_seq_number()
943                .expect("db error")
944        })
945        .await
946    }
947
948    pub async fn notify_read_executed_checkpoint(
949        &self,
950        seq: CheckpointSequenceNumber,
951    ) -> VerifiedCheckpoint {
952        self.notify_read_checkpoint_watermark(&self.executed_checkpoint_notify_read, seq, || {
953            self.get_highest_executed_checkpoint_seq_number()
954                .expect("db error")
955        })
956        .await
957    }
958
959    pub fn update_highest_executed_checkpoint(
960        &self,
961        checkpoint: &VerifiedCheckpoint,
962    ) -> Result<(), TypedStoreError> {
963        if let Some(seq_number) = self.get_highest_executed_checkpoint_seq_number()? {
964            if seq_number >= *checkpoint.sequence_number() {
965                return Ok(());
966            }
967            assert_eq!(
968                seq_number + 1,
969                *checkpoint.sequence_number(),
970                "Cannot update highest executed checkpoint to {} when current highest executed checkpoint is {}",
971                checkpoint.sequence_number(),
972                seq_number
973            );
974        }
975        let seq = *checkpoint.sequence_number();
976        debug!(checkpoint_seq = seq, "Updating highest executed checkpoint",);
977        self.tables.watermarks.insert(
978            &CheckpointWatermark::HighestExecuted,
979            &(seq, *checkpoint.digest()),
980        )?;
981        self.executed_checkpoint_notify_read
982            .notify(&seq, checkpoint);
983        Ok(())
984    }
985
986    pub fn update_highest_pruned_checkpoint(
987        &self,
988        checkpoint: &VerifiedCheckpoint,
989    ) -> Result<(), TypedStoreError> {
990        self.tables.watermarks.insert(
991            &CheckpointWatermark::HighestPruned,
992            &(*checkpoint.sequence_number(), *checkpoint.digest()),
993        )
994    }
995
996    /// Sets highest executed checkpoint to any value.
997    ///
998    /// WARNING: This method is very subtle and can corrupt the database if used incorrectly.
999    /// It should only be used in one-off cases or tests after fully understanding the risk.
1000    pub fn set_highest_executed_checkpoint_subtle(
1001        &self,
1002        checkpoint: &VerifiedCheckpoint,
1003    ) -> Result<(), TypedStoreError> {
1004        self.tables.watermarks.insert(
1005            &CheckpointWatermark::HighestExecuted,
1006            &(*checkpoint.sequence_number(), *checkpoint.digest()),
1007        )
1008    }
1009
1010    pub fn insert_checkpoint_contents(
1011        &self,
1012        contents: CheckpointContents,
1013    ) -> Result<(), TypedStoreError> {
1014        debug!(
1015            checkpoint_seq = ?contents.digest(),
1016            "Inserting checkpoint contents",
1017        );
1018        self.tables
1019            .checkpoint_content
1020            .insert(contents.digest(), &contents)
1021    }
1022
1023    pub fn insert_verified_checkpoint_contents(
1024        &self,
1025        checkpoint: &VerifiedCheckpoint,
1026        full_contents: VerifiedCheckpointContents,
1027    ) -> Result<(), TypedStoreError> {
1028        let mut batch = self.tables.full_checkpoint_content_v2.batch();
1029        batch.insert_batch(
1030            &self.tables.checkpoint_sequence_by_contents_digest,
1031            [(&checkpoint.content_digest, checkpoint.sequence_number())],
1032        )?;
1033        let full_contents = full_contents.into_inner();
1034        batch.insert_batch(
1035            &self.tables.full_checkpoint_content_v2,
1036            [(checkpoint.sequence_number(), &full_contents)],
1037        )?;
1038
1039        let contents = full_contents.into_checkpoint_contents();
1040        assert_eq!(&checkpoint.content_digest, contents.digest());
1041
1042        batch.insert_batch(
1043            &self.tables.checkpoint_content,
1044            [(contents.digest(), &contents)],
1045        )?;
1046
1047        batch.write()
1048    }
1049
1050    pub fn delete_full_checkpoint_contents(
1051        &self,
1052        seq: CheckpointSequenceNumber,
1053    ) -> Result<(), TypedStoreError> {
1054        self.tables.full_checkpoint_content.remove(&seq)?;
1055        self.tables.full_checkpoint_content_v2.remove(&seq)
1056    }
1057
1058    pub fn get_epoch_last_checkpoint(
1059        &self,
1060        epoch_id: EpochId,
1061    ) -> SuiResult<Option<VerifiedCheckpoint>> {
1062        let seq = self.get_epoch_last_checkpoint_seq_number(epoch_id)?;
1063        let checkpoint = match seq {
1064            Some(seq) => self.get_checkpoint_by_sequence_number(seq)?,
1065            None => None,
1066        };
1067        Ok(checkpoint)
1068    }
1069
1070    pub fn get_epoch_last_checkpoint_seq_number(
1071        &self,
1072        epoch_id: EpochId,
1073    ) -> SuiResult<Option<CheckpointSequenceNumber>> {
1074        let seq = self.tables.epoch_last_checkpoint_map.get(&epoch_id)?;
1075        Ok(seq)
1076    }
1077
1078    /// Returns the sequence number of the first checkpoint in the given epoch.
1079    /// For epoch 0 this is always 0; for epoch N > 0 it is last_checkpoint(N-1) + 1.
1080    pub fn get_epoch_first_checkpoint_seq(
1081        &self,
1082        epoch: EpochId,
1083    ) -> SuiResult<Option<CheckpointSequenceNumber>> {
1084        if epoch == 0 {
1085            return Ok(Some(0));
1086        }
1087        Ok(self
1088            .tables
1089            .epoch_last_checkpoint_map
1090            .get(&(epoch - 1))?
1091            .map(|s| s + 1))
1092    }
1093
1094    /// Iterate certified checkpoints starting at `start` (inclusive), up to `limit`.
1095    pub fn list_checkpoints_from_seq(
1096        &self,
1097        start: Option<CheckpointSequenceNumber>,
1098        limit: usize,
1099    ) -> Result<Vec<(CheckpointSequenceNumber, VerifiedCheckpoint)>, TypedStoreError> {
1100        self.tables
1101            .certified_checkpoints
1102            .safe_iter_with_bounds(start, None)
1103            .take(limit)
1104            .map(|r| r.map(|(seq, cp)| (seq, cp.into())))
1105            .collect()
1106    }
1107
1108    /// Iterate epoch→last-checkpoint-seq entries starting at `start` epoch, up to `limit`.
1109    pub fn list_epoch_last_checkpoints(
1110        &self,
1111        start: Option<EpochId>,
1112        limit: usize,
1113    ) -> Result<Vec<(EpochId, CheckpointSequenceNumber)>, TypedStoreError> {
1114        self.tables
1115            .epoch_last_checkpoint_map
1116            .safe_iter_with_bounds(start, None)
1117            .take(limit)
1118            .collect()
1119    }
1120
1121    /// Iterate checkpoint digests from `checkpoint_by_digest`, starting at `start`, up to `limit`.
1122    pub fn list_checkpoint_digests(
1123        &self,
1124        start: Option<CheckpointDigest>,
1125        limit: usize,
1126    ) -> Result<Vec<CheckpointDigest>, TypedStoreError> {
1127        self.tables
1128            .checkpoint_by_digest
1129            .safe_iter_with_bounds(start, None)
1130            .take(limit)
1131            .map(|r| r.map(|(d, _)| d))
1132            .collect()
1133    }
1134
1135    /// Iterate checkpoint contents digests, starting at `start`, up to `limit`.
1136    pub fn list_checkpoint_contents_digests(
1137        &self,
1138        start: Option<CheckpointContentsDigest>,
1139        limit: usize,
1140    ) -> Result<Vec<CheckpointContentsDigest>, TypedStoreError> {
1141        self.tables
1142            .checkpoint_content
1143            .safe_iter_with_bounds(start, None)
1144            .take(limit)
1145            .map(|r| r.map(|(d, _)| d))
1146            .collect()
1147    }
1148
1149    /// Iterate certified checkpoints belonging to `epoch`, starting at `start_seq`, up to `limit`.
1150    pub fn list_epoch_checkpoints(
1151        &self,
1152        epoch: EpochId,
1153        start_seq: Option<CheckpointSequenceNumber>,
1154        limit: usize,
1155    ) -> Result<Vec<(CheckpointSequenceNumber, VerifiedCheckpoint)>, TypedStoreError> {
1156        let Some(last_seq) = self.tables.epoch_last_checkpoint_map.get(&epoch)? else {
1157            return Ok(vec![]);
1158        };
1159        // Compute first seq directly to stay within TypedStoreError.
1160        let first_seq = if epoch == 0 {
1161            0
1162        } else {
1163            self.tables
1164                .epoch_last_checkpoint_map
1165                .get(&(epoch - 1))?
1166                .map(|s| s + 1)
1167                .unwrap_or(0)
1168        };
1169        let start = start_seq.map(|s| s.max(first_seq)).unwrap_or(first_seq);
1170        self.tables
1171            .certified_checkpoints
1172            .safe_iter_with_bounds(Some(start), Some(last_seq + 1))
1173            .take(limit)
1174            .map(|r| r.map(|(seq, cp)| (seq, cp.into())))
1175            .collect()
1176    }
1177
1178    pub fn insert_epoch_last_checkpoint(
1179        &self,
1180        epoch_id: EpochId,
1181        checkpoint: &VerifiedCheckpoint,
1182    ) -> SuiResult {
1183        self.tables
1184            .epoch_last_checkpoint_map
1185            .insert(&epoch_id, checkpoint.sequence_number())?;
1186        Ok(())
1187    }
1188
1189    pub fn get_epoch_state_commitments(
1190        &self,
1191        epoch: EpochId,
1192    ) -> SuiResult<Option<Vec<CheckpointCommitment>>> {
1193        let commitments = self.get_epoch_last_checkpoint(epoch)?.map(|checkpoint| {
1194            checkpoint
1195                .end_of_epoch_data
1196                .as_ref()
1197                .expect("Last checkpoint of epoch expected to have EndOfEpochData")
1198                .epoch_commitments
1199                .clone()
1200        });
1201        Ok(commitments)
1202    }
1203
1204    /// Given the epoch ID, and the last checkpoint of the epoch, derive a few statistics of the epoch.
1205    pub fn get_epoch_stats(
1206        &self,
1207        epoch: EpochId,
1208        last_checkpoint: &CheckpointSummary,
1209    ) -> Option<EpochStats> {
1210        let (first_checkpoint, prev_epoch_network_transactions) = if epoch == 0 {
1211            (0, 0)
1212        } else if let Ok(Some(checkpoint)) = self.get_epoch_last_checkpoint(epoch - 1) {
1213            (
1214                checkpoint.sequence_number + 1,
1215                checkpoint.network_total_transactions,
1216            )
1217        } else {
1218            return None;
1219        };
1220        Some(EpochStats {
1221            checkpoint_count: last_checkpoint.sequence_number - first_checkpoint + 1,
1222            transaction_count: last_checkpoint.network_total_transactions
1223                - prev_epoch_network_transactions,
1224            total_gas_reward: last_checkpoint
1225                .epoch_rolling_gas_cost_summary
1226                .computation_cost,
1227        })
1228    }
1229
1230    pub fn checkpoint_db(&self, path: &Path) -> SuiResult {
1231        // This checkpoints the entire db and not one column family
1232        self.tables
1233            .checkpoint_content
1234            .checkpoint_db(path)
1235            .map_err(Into::into)
1236    }
1237
1238    pub fn delete_highest_executed_checkpoint_test_only(&self) -> Result<(), TypedStoreError> {
1239        let mut wb = self.tables.watermarks.batch();
1240        wb.delete_batch(
1241            &self.tables.watermarks,
1242            std::iter::once(CheckpointWatermark::HighestExecuted),
1243        )?;
1244        wb.write()?;
1245        Ok(())
1246    }
1247
1248    pub fn reset_db_for_execution_since_genesis(&self) -> SuiResult {
1249        self.delete_highest_executed_checkpoint_test_only()?;
1250        Ok(())
1251    }
1252
1253    pub fn record_checkpoint_fork_detected(
1254        &self,
1255        checkpoint_seq: CheckpointSequenceNumber,
1256        checkpoint_digest: CheckpointDigest,
1257        certified_checkpoint_digest: Option<CheckpointDigest>,
1258    ) -> Result<(), TypedStoreError> {
1259        let binary_version = self.binary_version.get().cloned().unwrap_or_default();
1260        info!(
1261            checkpoint_seq = checkpoint_seq,
1262            checkpoint_digest = ?checkpoint_digest,
1263            certified_checkpoint_digest = ?certified_checkpoint_digest,
1264            binary_version,
1265            "Recording checkpoint fork detection in database"
1266        );
1267        self.tables.checkpoint_fork_detected.insert(
1268            &CHECKPOINT_FORK_DETECTED_KEY,
1269            &CheckpointForkInfo {
1270                checkpoint_seq,
1271                checkpoint_digest,
1272                certified_checkpoint_digest,
1273                binary_version,
1274            },
1275        )
1276    }
1277
1278    pub fn get_checkpoint_fork_detected(
1279        &self,
1280    ) -> Result<Option<CheckpointForkInfo>, TypedStoreError> {
1281        self.tables
1282            .checkpoint_fork_detected
1283            .get(&CHECKPOINT_FORK_DETECTED_KEY)
1284    }
1285
1286    pub fn clear_checkpoint_fork_detected(&self) -> Result<(), TypedStoreError> {
1287        self.tables
1288            .checkpoint_fork_detected
1289            .remove(&CHECKPOINT_FORK_DETECTED_KEY)
1290    }
1291
1292    pub fn record_transaction_fork_detected(
1293        &self,
1294        tx_digest: TransactionDigest,
1295        expected_effects_digest: TransactionEffectsDigest,
1296        actual_effects_digest: TransactionEffectsDigest,
1297        certified_checkpoint_seq: Option<CheckpointSequenceNumber>,
1298    ) -> Result<(), TypedStoreError> {
1299        let binary_version = self.binary_version.get().cloned().unwrap_or_default();
1300        info!(
1301            tx_digest = ?tx_digest,
1302            expected_effects_digest = ?expected_effects_digest,
1303            actual_effects_digest = ?actual_effects_digest,
1304            certified_checkpoint_seq = ?certified_checkpoint_seq,
1305            binary_version,
1306            "Recording transaction fork detection in database"
1307        );
1308        self.tables.transaction_fork_detected_v2.insert(
1309            &TRANSACTION_FORK_DETECTED_KEY,
1310            &TransactionForkInfo {
1311                tx_digest,
1312                expected_effects_digest,
1313                actual_effects_digest,
1314                certified_checkpoint_seq,
1315                binary_version,
1316            },
1317        )
1318    }
1319
1320    pub fn get_transaction_fork_detected(
1321        &self,
1322    ) -> Result<Option<TransactionForkInfo>, TypedStoreError> {
1323        self.tables
1324            .transaction_fork_detected_v2
1325            .get(&TRANSACTION_FORK_DETECTED_KEY)
1326    }
1327
1328    pub fn clear_transaction_fork_detected(&self) -> Result<(), TypedStoreError> {
1329        self.tables
1330            .transaction_fork_detected_v2
1331            .remove(&TRANSACTION_FORK_DETECTED_KEY)
1332    }
1333}
1334
1335#[derive(Copy, Clone, Debug, Serialize, Deserialize)]
1336pub enum CheckpointWatermark {
1337    HighestVerified,
1338    HighestSynced,
1339    HighestExecuted,
1340    HighestPruned,
1341    /// Retired: checkpoint fork markers now live in the `checkpoint_fork_detected` table. Kept
1342    /// because existing databases may hold this key.
1343    CheckpointForkDetected,
1344}
1345
1346struct CheckpointStateHasher {
1347    epoch_store: Arc<AuthorityPerEpochStore>,
1348    hasher: Weak<GlobalStateHasher>,
1349    receive_from_builder: mpsc::Receiver<(CheckpointSequenceNumber, Vec<TransactionEffects>)>,
1350}
1351
1352impl CheckpointStateHasher {
1353    fn new(
1354        epoch_store: Arc<AuthorityPerEpochStore>,
1355        hasher: Weak<GlobalStateHasher>,
1356        receive_from_builder: mpsc::Receiver<(CheckpointSequenceNumber, Vec<TransactionEffects>)>,
1357    ) -> Self {
1358        Self {
1359            epoch_store,
1360            hasher,
1361            receive_from_builder,
1362        }
1363    }
1364
1365    async fn run(self) {
1366        let Self {
1367            epoch_store,
1368            hasher,
1369            mut receive_from_builder,
1370        } = self;
1371        while let Some((seq, effects)) = receive_from_builder.recv().await {
1372            let Some(hasher) = hasher.upgrade() else {
1373                info!("Object state hasher was dropped, stopping checkpoint accumulation");
1374                break;
1375            };
1376            hasher.accumulate_checkpoint(&effects, seq, &epoch_store);
1377        }
1378    }
1379}
1380
1381#[derive(Debug)]
1382pub enum CheckpointBuilderError {
1383    ChangeEpochTxAlreadyExecuted,
1384    SystemPackagesMissing,
1385    Retry(anyhow::Error),
1386}
1387
1388impl<SuiError: std::error::Error + Send + Sync + 'static> From<SuiError>
1389    for CheckpointBuilderError
1390{
1391    fn from(e: SuiError) -> Self {
1392        Self::Retry(e.into())
1393    }
1394}
1395
1396pub type CheckpointBuilderResult<T = ()> = Result<T, CheckpointBuilderError>;
1397
1398pub struct CheckpointBuilder {
1399    state: Arc<AuthorityState>,
1400    store: Arc<CheckpointStore>,
1401    epoch_store: Arc<AuthorityPerEpochStore>,
1402    notify: Arc<Notify>,
1403    notify_aggregator: Arc<Notify>,
1404    last_built: watch::Sender<CheckpointSequenceNumber>,
1405    effects_store: Arc<dyn TransactionCacheRead>,
1406    global_state_hasher: Weak<GlobalStateHasher>,
1407    send_to_hasher: mpsc::Sender<(CheckpointSequenceNumber, Vec<TransactionEffects>)>,
1408    output: Box<dyn CheckpointOutput>,
1409    metrics: Arc<CheckpointMetrics>,
1410}
1411
1412pub struct CheckpointAggregator {
1413    store: Arc<CheckpointStore>,
1414    epoch_store: Arc<AuthorityPerEpochStore>,
1415    notify: Arc<Notify>,
1416    receiver: mpsc::UnboundedReceiver<CheckpointSignatureMessage>,
1417    pending: BTreeMap<CheckpointSequenceNumber, Vec<CheckpointSignatureMessage>>,
1418    current: Option<CheckpointSignatureAggregator>,
1419    output: Box<dyn CertifiedCheckpointOutput>,
1420    state: Arc<AuthorityState>,
1421    metrics: Arc<CheckpointMetrics>,
1422}
1423
1424// This holds information to aggregate signatures for one checkpoint
1425pub struct CheckpointSignatureAggregator {
1426    summary: CheckpointSummary,
1427    digest: CheckpointDigest,
1428    /// Aggregates voting stake for each signed checkpoint proposal by authority
1429    signatures_by_digest: MultiStakeAggregator<CheckpointDigest, CheckpointSummary, true>,
1430    store: Arc<CheckpointStore>,
1431    state: Arc<AuthorityState>,
1432    metrics: Arc<CheckpointMetrics>,
1433}
1434
1435impl CheckpointBuilder {
1436    fn new(
1437        state: Arc<AuthorityState>,
1438        store: Arc<CheckpointStore>,
1439        epoch_store: Arc<AuthorityPerEpochStore>,
1440        notify: Arc<Notify>,
1441        effects_store: Arc<dyn TransactionCacheRead>,
1442        // for synchronous accumulation of end-of-epoch checkpoint
1443        global_state_hasher: Weak<GlobalStateHasher>,
1444        // for asynchronous/concurrent accumulation of regular checkpoints
1445        send_to_hasher: mpsc::Sender<(CheckpointSequenceNumber, Vec<TransactionEffects>)>,
1446        output: Box<dyn CheckpointOutput>,
1447        notify_aggregator: Arc<Notify>,
1448        last_built: watch::Sender<CheckpointSequenceNumber>,
1449        metrics: Arc<CheckpointMetrics>,
1450    ) -> Self {
1451        Self {
1452            state,
1453            store,
1454            epoch_store,
1455            notify,
1456            effects_store,
1457            global_state_hasher,
1458            send_to_hasher,
1459            output,
1460            notify_aggregator,
1461            last_built,
1462            metrics,
1463        }
1464    }
1465
1466    /// This function first waits for ConsensusCommitHandler to finish reprocessing
1467    /// commits that have been processed before the last restart, if consensus_replay_waiter
1468    /// is supplied. Then it starts building checkpoints in a loop.
1469    ///
1470    /// It is optional to pass in consensus_replay_waiter, to make it easier to attribute
1471    /// if slow recovery of previously built checkpoints is due to consensus replay or
1472    /// checkpoint building.
1473    async fn run(mut self, consensus_replay_waiter: Option<ReplayWaiter>) {
1474        if let Some(replay_waiter) = consensus_replay_waiter {
1475            info!("Waiting for consensus commits to replay ...");
1476            replay_waiter.wait_for_replay().await;
1477            info!("Consensus commits finished replaying");
1478        }
1479        info!("Starting CheckpointBuilder");
1480        loop {
1481            match self.maybe_build_checkpoints().await {
1482                Ok(()) => {}
1483                err @ Err(
1484                    CheckpointBuilderError::ChangeEpochTxAlreadyExecuted
1485                    | CheckpointBuilderError::SystemPackagesMissing,
1486                ) => {
1487                    info!("CheckpointBuilder stopping: {:?}", err);
1488                    return;
1489                }
1490                Err(CheckpointBuilderError::Retry(inner)) => {
1491                    let msg = format!("{:?}", inner);
1492                    debug_fatal!("Error while making checkpoint, will retry in 1s: {}", msg);
1493                    tokio::time::sleep(Duration::from_secs(1)).await;
1494                    self.metrics.checkpoint_errors.inc();
1495                    continue;
1496                }
1497            }
1498
1499            self.notify.notified().await;
1500        }
1501    }
1502
1503    async fn maybe_build_checkpoints(&mut self) -> CheckpointBuilderResult {
1504        let _scope = monitored_scope("BuildCheckpoints");
1505
1506        // Collect info about the most recently built checkpoint.
1507        let last_height = self
1508            .epoch_store
1509            .last_built_checkpoint_builder_summary()
1510            .expect("epoch should not have ended")
1511            .and_then(|s| s.checkpoint_height);
1512
1513        for (height, pending) in self.epoch_store.get_pending_checkpoints(last_height) {
1514            debug!(checkpoint_commit_height = height, "Making checkpoint");
1515
1516            let seq = self.make_checkpoint(pending).await?;
1517
1518            self.last_built.send_if_modified(|cur| {
1519                // when rebuilding checkpoints at startup, seq can be for an old checkpoint
1520                if seq > *cur {
1521                    *cur = seq;
1522                    true
1523                } else {
1524                    false
1525                }
1526            });
1527
1528            // ensure that the task can be cancelled at end of epoch, even if no other await yields
1529            // execution.
1530            tokio::task::yield_now().await;
1531        }
1532
1533        Ok(())
1534    }
1535
1536    #[instrument(level = "debug", skip_all, fields(height = pending.details.checkpoint_height))]
1537    async fn make_checkpoint(
1538        &mut self,
1539        pending: PendingCheckpoint,
1540    ) -> CheckpointBuilderResult<CheckpointSequenceNumber> {
1541        let _scope = monitored_scope("CheckpointBuilder::make_checkpoint");
1542
1543        let details = pending.details.clone();
1544
1545        let highest_executed_sequence = self
1546            .store
1547            .get_highest_executed_checkpoint_seq_number()
1548            .expect("db error")
1549            .unwrap_or(0);
1550
1551        let (poll_count, result) = poll_count(self.resolve_checkpoint_transactions(pending)).await;
1552        let (sorted_tx_effects_included_in_checkpoint, all_roots) = result?;
1553
1554        let new_checkpoint = self
1555            .create_checkpoint(
1556                sorted_tx_effects_included_in_checkpoint,
1557                &details,
1558                &all_roots,
1559            )
1560            .await?;
1561        let sequence = *new_checkpoint.0.sequence_number();
1562        let digest = new_checkpoint.0.digest();
1563        if sequence <= highest_executed_sequence && poll_count > 1 {
1564            debug_fatal!(
1565                "resolve_checkpoint_transactions should be instantaneous when executed checkpoint is ahead of checkpoint builder"
1566            );
1567        }
1568
1569        self.write_checkpoint(details.checkpoint_height, new_checkpoint)
1570            .await?;
1571        info!(
1572            seq = sequence,
1573            %digest,
1574            height = details.checkpoint_height,
1575            commit = %details.consensus_commit_ref,
1576            "Made new checkpoint"
1577        );
1578
1579        Ok(sequence)
1580    }
1581
1582    // Given the root transactions of a pending checkpoint, resolve the transactions should be included in
1583    // the checkpoint, and return them in the order they should be included in the checkpoint.
1584    #[instrument(level = "debug", skip_all)]
1585    async fn resolve_checkpoint_transactions(
1586        &self,
1587        pending: PendingCheckpoint,
1588    ) -> SuiResult<(Vec<TransactionEffects>, HashSet<TransactionDigest>)> {
1589        let _scope = monitored_scope("CheckpointBuilder::resolve_checkpoint_transactions");
1590
1591        debug!(
1592            checkpoint_commit_height = pending.details.checkpoint_height,
1593            "Resolving checkpoint transactions for pending checkpoint.",
1594        );
1595
1596        trace!(
1597            "roots for pending checkpoint {:?}: {:?}",
1598            pending.details.checkpoint_height, pending.roots,
1599        );
1600
1601        assert!(
1602            self.epoch_store
1603                .protocol_config()
1604                .prepend_prologue_tx_in_consensus_commit_in_checkpoints()
1605        );
1606
1607        let mut all_effects: Vec<TransactionEffects> = Vec::new();
1608        let mut all_root_digests: Vec<TransactionDigest> = Vec::new();
1609
1610        for checkpoint_roots in &pending.roots {
1611            let tx_roots = &checkpoint_roots.tx_roots;
1612
1613            self.metrics
1614                .checkpoint_roots_count
1615                .inc_by(tx_roots.len() as u64);
1616
1617            let root_digests = self
1618                .epoch_store
1619                .notify_read_tx_key_to_digest(tx_roots)
1620                .in_monitored_scope("CheckpointNotifyDigests")
1621                .await?;
1622
1623            all_root_digests.extend(root_digests.iter().cloned());
1624
1625            let root_effects = self
1626                .effects_store
1627                .notify_read_executed_effects(
1628                    CHECKPOINT_BUILDER_NOTIFY_READ_TASK_NAME,
1629                    &root_digests,
1630                )
1631                .in_monitored_scope("CheckpointNotifyRead")
1632                .await;
1633            let consensus_commit_prologue =
1634                self.extract_consensus_commit_prologue(&root_digests, &root_effects)?;
1635
1636            let _scope = monitored_scope("CheckpointBuilder::causal_sort");
1637            let ccp_digest = consensus_commit_prologue.map(|(d, _)| d);
1638            let mut sorted = CausalOrder::order_for_checkpoint(
1639                root_effects,
1640                ccp_digest,
1641                self.epoch_store.protocol_config(),
1642            );
1643
1644            if let Some(settlement_key) = &checkpoint_roots.settlement_root {
1645                let checkpoint_seq = pending.details.checkpoint_seq;
1646                let tx_index_offset = all_effects.len() as u64;
1647                let effects = self
1648                    .resolve_settlement_effects(
1649                        *settlement_key,
1650                        &sorted,
1651                        checkpoint_roots.height,
1652                        checkpoint_seq,
1653                        tx_index_offset,
1654                    )
1655                    .await;
1656                sorted.extend(effects);
1657            }
1658
1659            #[cfg(msim)]
1660            {
1661                self.expensive_consensus_commit_prologue_invariants_check(&root_digests, &sorted);
1662            }
1663
1664            all_effects.extend(sorted);
1665        }
1666        Ok((all_effects, all_root_digests.into_iter().collect()))
1667    }
1668
1669    /// Constructs settlement transactions to compute their digests, then reads effects
1670    /// directly from the cache. If execution is ahead of the checkpoint builder, the
1671    /// effects are already cached and this returns instantly. Otherwise it waits for
1672    /// the execution scheduler's queue worker to execute them.
1673    async fn resolve_settlement_effects(
1674        &self,
1675        settlement_key: TransactionKey,
1676        sorted_root_effects: &[TransactionEffects],
1677        checkpoint_height: CheckpointHeight,
1678        checkpoint_seq: CheckpointSequenceNumber,
1679        tx_index_offset: u64,
1680    ) -> Vec<TransactionEffects> {
1681        let epoch = self.epoch_store.epoch();
1682        let accumulator_root_obj_initial_shared_version = self
1683            .epoch_store
1684            .epoch_start_config()
1685            .accumulator_root_obj_initial_shared_version()
1686            .expect("accumulator root object must exist");
1687
1688        let builder = AccumulatorSettlementTxBuilder::new(
1689            None,
1690            sorted_root_effects,
1691            checkpoint_seq,
1692            tx_index_offset,
1693        );
1694
1695        let settlement_digests: Vec<_> = builder
1696            .build_tx(
1697                self.epoch_store.protocol_config(),
1698                epoch,
1699                accumulator_root_obj_initial_shared_version,
1700                checkpoint_height,
1701                checkpoint_seq,
1702            )
1703            .into_iter()
1704            .map(|tx| *VerifiedTransaction::new_system_transaction(tx).digest())
1705            .collect();
1706
1707        debug!(
1708            ?settlement_digests,
1709            ?settlement_key,
1710            "reading settlement effects from cache"
1711        );
1712
1713        let settlement_effects = wait_for_effects_with_retry(
1714            self.effects_store.as_ref(),
1715            "CheckpointBuilder::settlement_effects",
1716            &settlement_digests,
1717            settlement_key,
1718        )
1719        .await;
1720        let (accounts_created, accounts_deleted) =
1721            accumulators::count_accumulator_object_changes(&settlement_effects);
1722        self.metrics
1723            .report_accumulator_account_changes(accounts_created, accounts_deleted);
1724
1725        let barrier_digest = *VerifiedTransaction::new_system_transaction(
1726            accumulators::build_accumulator_barrier_tx(
1727                epoch,
1728                accumulator_root_obj_initial_shared_version,
1729                checkpoint_height,
1730                &settlement_effects,
1731            ),
1732        )
1733        .digest();
1734
1735        let barrier_effects = wait_for_effects_with_retry(
1736            self.effects_store.as_ref(),
1737            "CheckpointBuilder::barrier_effects",
1738            &[barrier_digest],
1739            settlement_key,
1740        )
1741        .await;
1742
1743        // Assert success here, in the builder task, before these effects are included in
1744        // the checkpoint. The settlement scheduler also asserts this, but it runs in a
1745        // separate task, so its assertion does not order against checkpoint persistence -
1746        // checking it here is what prevents a checkpoint from being built over the effects
1747        // of a failed settlement transaction.
1748        for fx in settlement_effects.iter().chain(barrier_effects.iter()) {
1749            assert!(
1750                fx.status().is_ok(),
1751                "settlement transaction cannot fail (digest: {:?}) {:#?}",
1752                fx.transaction_digest(),
1753                fx
1754            );
1755        }
1756
1757        settlement_effects
1758            .into_iter()
1759            .chain(barrier_effects)
1760            .collect()
1761    }
1762
1763    // Extracts the consensus commit prologue digest and effects from the root transactions.
1764    // The consensus commit prologue is expected to be the first transaction in the roots.
1765    fn extract_consensus_commit_prologue(
1766        &self,
1767        root_digests: &[TransactionDigest],
1768        root_effects: &[TransactionEffects],
1769    ) -> SuiResult<Option<(TransactionDigest, TransactionEffects)>> {
1770        let _scope = monitored_scope("CheckpointBuilder::extract_consensus_commit_prologue");
1771        if root_digests.is_empty() {
1772            return Ok(None);
1773        }
1774
1775        // Reads the first transaction in the roots, and checks whether it is a consensus commit
1776        // prologue transaction. The consensus commit prologue transaction should be the first
1777        // transaction in the roots written by the consensus handler.
1778        let first_tx = self
1779            .state
1780            .get_transaction_cache_reader()
1781            .get_transaction_block(&root_digests[0])
1782            .expect("Transaction block must exist");
1783
1784        Ok(first_tx
1785            .transaction_data()
1786            .is_consensus_commit_prologue()
1787            .then(|| {
1788                assert_eq!(first_tx.digest(), root_effects[0].transaction_digest());
1789                (*first_tx.digest(), root_effects[0].clone())
1790            }))
1791    }
1792
1793    #[instrument(level = "debug", skip_all)]
1794    async fn write_checkpoint(
1795        &mut self,
1796        height: CheckpointHeight,
1797        new_checkpoint: (CheckpointSummary, CheckpointContents),
1798    ) -> SuiResult {
1799        let _scope = monitored_scope("CheckpointBuilder::write_checkpoint");
1800        let mut batch = self.store.tables.checkpoint_content.batch();
1801
1802        let (summary, contents) = &new_checkpoint;
1803        debug!(
1804            checkpoint_commit_height = height,
1805            checkpoint_seq = summary.sequence_number,
1806            contents_digest = ?contents.digest(),
1807            "writing checkpoint",
1808        );
1809
1810        if let Some(previously_computed_summary) = self
1811            .store
1812            .tables
1813            .locally_computed_checkpoints
1814            .get(&summary.sequence_number)?
1815            && previously_computed_summary.digest() != summary.digest()
1816        {
1817            // The builder re-derived this sequence with a different result than a previous run
1818            // (e.g. after a binary upgrade), which may already have signed and sent its result
1819            // to consensus. No certified checkpoint proves which result is canonical, so the
1820            // marker carries no certified digest and is never cleared automatically.
1821            if let Err(e) = self.store.record_checkpoint_fork_detected(
1822                summary.sequence_number,
1823                previously_computed_summary.digest(),
1824                None,
1825            ) {
1826                error!("Failed to record checkpoint fork in database: {:?}", e);
1827            }
1828            fatal!(
1829                "Checkpoint {} was previously built with a different result: previously_computed_summary {:?} vs current_summary {:?}",
1830                summary.sequence_number,
1831                previously_computed_summary.digest(),
1832                summary.digest()
1833            );
1834        }
1835
1836        self.metrics
1837            .transactions_included_in_checkpoint
1838            .inc_by(contents.size() as u64);
1839        let sequence_number = summary.sequence_number;
1840        self.metrics
1841            .last_constructed_checkpoint
1842            .set(sequence_number as i64);
1843
1844        batch.insert_batch(
1845            &self.store.tables.checkpoint_content,
1846            [(contents.digest(), contents)],
1847        )?;
1848
1849        batch.insert_batch(
1850            &self.store.tables.locally_computed_checkpoints,
1851            [(sequence_number, summary)],
1852        )?;
1853
1854        batch.write()?;
1855
1856        // Send checkpoint sigs to consensus.
1857        self.output
1858            .checkpoint_created(summary, contents, &self.epoch_store, &self.store)
1859            .await?;
1860
1861        if let Some(certified_checkpoint) = self
1862            .store
1863            .tables
1864            .certified_checkpoints
1865            .get(summary.sequence_number())?
1866        {
1867            self.store
1868                .check_for_checkpoint_fork(summary, &certified_checkpoint.into());
1869        }
1870
1871        self.notify_aggregator.notify_one();
1872        self.epoch_store
1873            .process_constructed_checkpoint(height, new_checkpoint.0);
1874        Ok(())
1875    }
1876
1877    fn load_last_built_checkpoint_summary(
1878        epoch_store: &AuthorityPerEpochStore,
1879        store: &CheckpointStore,
1880    ) -> SuiResult<Option<(CheckpointSequenceNumber, CheckpointSummary)>> {
1881        let mut last_checkpoint = epoch_store.last_built_checkpoint_summary()?;
1882        if last_checkpoint.is_none() {
1883            let epoch = epoch_store.epoch();
1884            if epoch > 0 {
1885                let previous_epoch = epoch - 1;
1886                let last_verified = store.get_epoch_last_checkpoint(previous_epoch)?;
1887                last_checkpoint = last_verified.map(VerifiedCheckpoint::into_summary_and_sequence);
1888                if let Some((ref seq, _)) = last_checkpoint {
1889                    debug!(
1890                        "No checkpoints in builder DB, taking checkpoint from previous epoch with sequence {seq}"
1891                    );
1892                } else {
1893                    // This is some serious bug with when CheckpointBuilder started so surfacing it via panic
1894                    panic!("Can not find last checkpoint for previous epoch {previous_epoch}");
1895                }
1896            }
1897        }
1898        Ok(last_checkpoint)
1899    }
1900
1901    #[instrument(level = "debug", skip_all)]
1902    async fn create_checkpoint(
1903        &self,
1904        all_effects: Vec<TransactionEffects>,
1905        details: &PendingCheckpointInfo,
1906        all_roots: &HashSet<TransactionDigest>,
1907    ) -> CheckpointBuilderResult<(CheckpointSummary, CheckpointContents)> {
1908        let _scope = monitored_scope("CheckpointBuilder::create_checkpoint");
1909
1910        let last_checkpoint =
1911            Self::load_last_built_checkpoint_summary(&self.epoch_store, &self.store)?;
1912        let last_checkpoint_seq = last_checkpoint.as_ref().map(|(seq, _)| *seq);
1913        debug!(
1914            checkpoint_commit_height = details.checkpoint_height,
1915            next_checkpoint_seq = last_checkpoint_seq.unwrap_or_default() + 1,
1916            checkpoint_timestamp = details.timestamp_ms,
1917            "Creating checkpoint for {} transactions",
1918            all_effects.len(),
1919        );
1920
1921        let all_digests: Vec<_> = all_effects
1922            .iter()
1923            .map(|effect| *effect.transaction_digest())
1924            .collect();
1925        let transaction_blocks = self
1926            .state
1927            .get_transaction_cache_reader()
1928            .multi_get_transaction_blocks(&all_digests);
1929        let mut transactions = Vec::with_capacity(all_effects.len());
1930        let mut transaction_keys = Vec::with_capacity(all_effects.len());
1931        let mut randomness_rounds = BTreeMap::new();
1932        {
1933            let _guard = monitored_scope("CheckpointBuilder::wait_for_transactions_sequenced");
1934            debug!(
1935                ?last_checkpoint_seq,
1936                "Waiting for {:?} certificates to appear in consensus",
1937                all_effects.len()
1938            );
1939
1940            for (effects, transaction) in all_effects
1941                .iter()
1942                .zip_debug_eq(transaction_blocks.into_iter())
1943            {
1944                let transaction = transaction
1945                    .unwrap_or_else(|| panic!("Could not find executed transaction {:?}", effects));
1946                match transaction.inner().transaction_data().kind() {
1947                    TransactionKind::ConsensusCommitPrologue(_)
1948                    | TransactionKind::ConsensusCommitPrologueV2(_)
1949                    | TransactionKind::ConsensusCommitPrologueV3(_)
1950                    | TransactionKind::ConsensusCommitPrologueV4(_)
1951                    | TransactionKind::AuthenticatorStateUpdate(_) => {
1952                        // ConsensusCommitPrologue and AuthenticatorStateUpdate are guaranteed to be
1953                        // processed before we reach here.
1954                    }
1955                    TransactionKind::ProgrammableSystemTransaction(_) => {
1956                        // settlement transactions are added by checkpoint builder
1957                    }
1958                    TransactionKind::ChangeEpoch(_)
1959                    | TransactionKind::Genesis(_)
1960                    | TransactionKind::EndOfEpochTransaction(_) => {
1961                        fatal!(
1962                            "unexpected transaction in checkpoint effects: {:?}",
1963                            transaction
1964                        );
1965                    }
1966                    TransactionKind::RandomnessStateUpdate(rsu) => {
1967                        randomness_rounds
1968                            .insert(*effects.transaction_digest(), rsu.randomness_round);
1969                    }
1970                    TransactionKind::ProgrammableTransaction(_) => {
1971                        // Only transactions that are not roots should be included in the call to
1972                        // `consensus_messages_processed_notify`. roots come directly from the consensus
1973                        // commit and so are known to be processed already.
1974                        let digest = *effects.transaction_digest();
1975                        if !all_roots.contains(&digest) {
1976                            transaction_keys.push(SequencedConsensusTransactionKey::External(
1977                                ConsensusTransactionKey::Certificate(digest),
1978                            ));
1979                        }
1980                    }
1981                }
1982                transactions.push((*transaction).clone());
1983            }
1984
1985            self.epoch_store
1986                .consensus_messages_processed_notify(transaction_keys)
1987                .await;
1988        }
1989
1990        let signatures = self
1991            .epoch_store
1992            .user_signatures_for_checkpoint(&transactions, &all_digests);
1993        debug!(
1994            ?last_checkpoint_seq,
1995            "Received {} checkpoint user signatures from consensus",
1996            signatures.len()
1997        );
1998
1999        let end_of_epoch_observation_keys: Option<Vec<_>> = if details.last_of_epoch {
2000            Some(
2001                transactions
2002                    .iter()
2003                    .flat_map(|tx| {
2004                        if let TransactionKind::ProgrammableTransaction(ptb) =
2005                            tx.transaction_data().kind()
2006                        {
2007                            itertools::Either::Left(
2008                                ptb.commands
2009                                    .iter()
2010                                    .map(ExecutionTimeObservationKey::from_command),
2011                            )
2012                        } else {
2013                            itertools::Either::Right(std::iter::empty())
2014                        }
2015                    })
2016                    .collect(),
2017            )
2018        } else {
2019            None
2020        };
2021
2022        let epoch = self.epoch_store.epoch();
2023        let first_checkpoint_of_epoch = last_checkpoint
2024            .as_ref()
2025            .map(|(_, c)| c.epoch != epoch)
2026            .unwrap_or(true);
2027        if first_checkpoint_of_epoch {
2028            self.epoch_store
2029                .record_epoch_first_checkpoint_creation_time_metric();
2030        }
2031        let last_checkpoint_of_epoch = details.last_of_epoch;
2032
2033        let sequence_number = details.checkpoint_seq;
2034        let timestamp_ms = details.timestamp_ms;
2035        if let Some((_, last_checkpoint)) = &last_checkpoint
2036            && last_checkpoint.timestamp_ms > timestamp_ms
2037        {
2038            debug_fatal!(
2039                "Decrease of checkpoint timestamp. Sequence: {}, previous: {}, current: {}",
2040                sequence_number,
2041                last_checkpoint.timestamp_ms,
2042                timestamp_ms
2043            );
2044        }
2045
2046        let mut effects = all_effects;
2047        let mut signatures = signatures;
2048        let epoch_rolling_gas_cost_summary =
2049            self.get_epoch_total_gas_cost(last_checkpoint.as_ref().map(|(_, c)| c), &effects);
2050
2051        let end_of_epoch_data = if last_checkpoint_of_epoch {
2052            let system_state_obj = self
2053                .augment_epoch_last_checkpoint(
2054                    &epoch_rolling_gas_cost_summary,
2055                    timestamp_ms,
2056                    &mut effects,
2057                    &mut signatures,
2058                    sequence_number,
2059                    end_of_epoch_observation_keys.expect(
2060                        "end_of_epoch_observation_keys must be populated for the last checkpoint",
2061                    ),
2062                    last_checkpoint_seq.unwrap_or_default(),
2063                )
2064                .await?;
2065
2066            let committee = system_state_obj
2067                .get_current_epoch_committee()
2068                .committee()
2069                .clone();
2070
2071            // This must happen after the call to augment_epoch_last_checkpoint,
2072            // otherwise we will not capture the change_epoch tx.
2073            let root_state_digest = {
2074                let state_acc = self
2075                    .global_state_hasher
2076                    .upgrade()
2077                    .expect("No checkpoints should be getting built after local configuration");
2078                let acc =
2079                    state_acc.accumulate_checkpoint(&effects, sequence_number, &self.epoch_store);
2080
2081                state_acc
2082                    .wait_for_previous_running_root(&self.epoch_store, sequence_number)
2083                    .await?;
2084
2085                state_acc.accumulate_running_root(&self.epoch_store, sequence_number, Some(acc))?;
2086                state_acc
2087                    .digest_epoch(self.epoch_store.clone(), sequence_number)
2088                    .await?
2089            };
2090            self.metrics.highest_accumulated_epoch.set(epoch as i64);
2091            info!("Epoch {epoch} root state hash digest: {root_state_digest:?}");
2092
2093            let epoch_commitments = if self
2094                .epoch_store
2095                .protocol_config()
2096                .commit_root_state_digest()
2097            {
2098                vec![root_state_digest.into()]
2099            } else {
2100                vec![]
2101            };
2102
2103            Some(EndOfEpochData {
2104                next_epoch_committee: committee.voting_rights,
2105                next_epoch_protocol_version: ProtocolVersion::new(
2106                    system_state_obj.protocol_version(),
2107                ),
2108                epoch_commitments,
2109            })
2110        } else {
2111            self.send_to_hasher
2112                .send((sequence_number, effects.clone()))
2113                .await?;
2114
2115            None
2116        };
2117        let contents = if self.epoch_store.protocol_config().address_aliases() {
2118            CheckpointContents::new_v2(&effects, signatures)
2119        } else {
2120            CheckpointContents::new_with_digests_and_signatures(
2121                effects.iter().map(TransactionEffects::execution_digests),
2122                signatures
2123                    .into_iter()
2124                    .map(|sigs| sigs.into_iter().map(|(s, _)| s).collect())
2125                    .collect(),
2126            )
2127        };
2128
2129        let num_txns = contents.size() as u64;
2130
2131        let network_total_transactions = last_checkpoint
2132            .as_ref()
2133            .map(|(_, c)| c.network_total_transactions + num_txns)
2134            .unwrap_or(num_txns);
2135
2136        let previous_digest = last_checkpoint.as_ref().map(|(_, c)| c.digest());
2137
2138        let matching_randomness_rounds: Vec<_> = effects
2139            .iter()
2140            .filter_map(|e| randomness_rounds.get(e.transaction_digest()))
2141            .copied()
2142            .collect();
2143
2144        let checkpoint_commitments = if self
2145            .epoch_store
2146            .protocol_config()
2147            .include_checkpoint_artifacts_digest_in_summary()
2148        {
2149            let artifacts = CheckpointArtifacts::from(&effects[..]);
2150            let artifacts_digest = artifacts.digest()?;
2151            vec![artifacts_digest.into()]
2152        } else {
2153            Default::default()
2154        };
2155
2156        let summary = CheckpointSummary::new(
2157            self.epoch_store.protocol_config(),
2158            epoch,
2159            sequence_number,
2160            network_total_transactions,
2161            &contents,
2162            previous_digest,
2163            epoch_rolling_gas_cost_summary,
2164            end_of_epoch_data,
2165            timestamp_ms,
2166            matching_randomness_rounds,
2167            checkpoint_commitments,
2168        );
2169        summary.report_checkpoint_age(
2170            &self.metrics.last_created_checkpoint_age,
2171            &self.metrics.last_created_checkpoint_age_ms,
2172        );
2173        if last_checkpoint_of_epoch {
2174            info!(
2175                checkpoint_seq = sequence_number,
2176                "creating last checkpoint of epoch {}", epoch
2177            );
2178            if let Some(stats) = self.store.get_epoch_stats(epoch, &summary) {
2179                self.epoch_store
2180                    .report_epoch_metrics_at_last_checkpoint(stats);
2181            }
2182        }
2183
2184        Ok((summary, contents))
2185    }
2186
2187    fn get_epoch_total_gas_cost(
2188        &self,
2189        last_checkpoint: Option<&CheckpointSummary>,
2190        cur_checkpoint_effects: &[TransactionEffects],
2191    ) -> GasCostSummary {
2192        let (previous_epoch, previous_gas_costs) = last_checkpoint
2193            .map(|c| (c.epoch, c.epoch_rolling_gas_cost_summary.clone()))
2194            .unwrap_or_default();
2195        let current_gas_costs = GasCostSummary::new_from_txn_effects(cur_checkpoint_effects.iter());
2196        if previous_epoch == self.epoch_store.epoch() {
2197            // sum only when we are within the same epoch
2198            GasCostSummary::new(
2199                previous_gas_costs.computation_cost + current_gas_costs.computation_cost,
2200                previous_gas_costs.storage_cost + current_gas_costs.storage_cost,
2201                previous_gas_costs.storage_rebate + current_gas_costs.storage_rebate,
2202                previous_gas_costs.non_refundable_storage_fee
2203                    + current_gas_costs.non_refundable_storage_fee,
2204            )
2205        } else {
2206            current_gas_costs
2207        }
2208    }
2209
2210    #[instrument(level = "error", skip_all)]
2211    async fn augment_epoch_last_checkpoint(
2212        &self,
2213        epoch_total_gas_cost: &GasCostSummary,
2214        epoch_start_timestamp_ms: CheckpointTimestamp,
2215        checkpoint_effects: &mut Vec<TransactionEffects>,
2216        signatures: &mut Vec<Vec<(GenericSignature, Option<SequenceNumber>)>>,
2217        checkpoint: CheckpointSequenceNumber,
2218        end_of_epoch_observation_keys: Vec<ExecutionTimeObservationKey>,
2219        // This may be less than `checkpoint - 1` if the end-of-epoch PendingCheckpoint produced
2220        // >1 checkpoint.
2221        last_checkpoint: CheckpointSequenceNumber,
2222    ) -> CheckpointBuilderResult<SuiSystemState> {
2223        let (system_state, effects) = self
2224            .state
2225            .create_and_execute_advance_epoch_tx(
2226                &self.epoch_store,
2227                epoch_total_gas_cost,
2228                checkpoint,
2229                epoch_start_timestamp_ms,
2230                end_of_epoch_observation_keys,
2231                last_checkpoint,
2232            )
2233            .await?;
2234        checkpoint_effects.push(effects);
2235        signatures.push(vec![]);
2236        Ok(system_state)
2237    }
2238
2239    // Checks the invariants of the consensus commit prologue transactions in the checkpoint
2240    // in simtest.
2241    #[cfg(msim)]
2242    fn expensive_consensus_commit_prologue_invariants_check(
2243        &self,
2244        root_digests: &[TransactionDigest],
2245        sorted: &[TransactionEffects],
2246    ) {
2247        // Gets all the consensus commit prologue transactions from the roots.
2248        let root_txs = self
2249            .state
2250            .get_transaction_cache_reader()
2251            .multi_get_transaction_blocks(root_digests);
2252        let ccps = root_txs
2253            .iter()
2254            .filter_map(|tx| {
2255                if let Some(tx) = tx {
2256                    if tx.transaction_data().is_consensus_commit_prologue() {
2257                        Some(tx)
2258                    } else {
2259                        None
2260                    }
2261                } else {
2262                    None
2263                }
2264            })
2265            .collect::<Vec<_>>();
2266
2267        // There should be at most one consensus commit prologue transaction in the roots.
2268        assert!(ccps.len() <= 1);
2269
2270        // Get all the transactions in the checkpoint.
2271        let txs = self
2272            .state
2273            .get_transaction_cache_reader()
2274            .multi_get_transaction_blocks(
2275                &sorted
2276                    .iter()
2277                    .map(|tx| tx.transaction_digest().clone())
2278                    .collect::<Vec<_>>(),
2279            );
2280
2281        if ccps.len() == 0 {
2282            // If there is no consensus commit prologue transaction in the roots, then there should be no
2283            // consensus commit prologue transaction in the checkpoint.
2284            for tx in txs.iter() {
2285                if let Some(tx) = tx {
2286                    assert!(!tx.transaction_data().is_consensus_commit_prologue());
2287                }
2288            }
2289        } else {
2290            // If there is one consensus commit prologue, it must be the first one in the checkpoint.
2291            assert!(
2292                txs[0]
2293                    .as_ref()
2294                    .unwrap()
2295                    .transaction_data()
2296                    .is_consensus_commit_prologue()
2297            );
2298
2299            assert_eq!(ccps[0].digest(), txs[0].as_ref().unwrap().digest());
2300
2301            for tx in txs.iter().skip(1) {
2302                if let Some(tx) = tx {
2303                    assert!(!tx.transaction_data().is_consensus_commit_prologue());
2304                }
2305            }
2306        }
2307    }
2308}
2309
2310async fn wait_for_effects_with_retry(
2311    effects_store: &dyn TransactionCacheRead,
2312    task_name: &'static str,
2313    digests: &[TransactionDigest],
2314    tx_key: TransactionKey,
2315) -> Vec<TransactionEffects> {
2316    let delay = if in_antithesis() {
2317        // antithesis pauses containers and threads for tens of seconds, so shorter
2318        // timeouts produce false positives
2319        60
2320    } else {
2321        5
2322    };
2323    loop {
2324        match tokio::time::timeout(Duration::from_secs(delay), async {
2325            effects_store
2326                .notify_read_executed_effects(task_name, digests)
2327                .await
2328        })
2329        .await
2330        {
2331            Ok(effects) => break effects,
2332            Err(_) => {
2333                debug_fatal!(
2334                    "Timeout waiting for transactions to be executed {:?}, retrying...",
2335                    tx_key
2336                );
2337            }
2338        }
2339    }
2340}
2341
2342impl CheckpointAggregator {
2343    fn new(
2344        tables: Arc<CheckpointStore>,
2345        epoch_store: Arc<AuthorityPerEpochStore>,
2346        notify: Arc<Notify>,
2347        receiver: mpsc::UnboundedReceiver<CheckpointSignatureMessage>,
2348        output: Box<dyn CertifiedCheckpointOutput>,
2349        state: Arc<AuthorityState>,
2350        metrics: Arc<CheckpointMetrics>,
2351    ) -> Self {
2352        Self {
2353            store: tables,
2354            epoch_store,
2355            notify,
2356            receiver,
2357            pending: BTreeMap::new(),
2358            current: None,
2359            output,
2360            state,
2361            metrics,
2362        }
2363    }
2364
2365    async fn run(mut self) {
2366        info!("Starting CheckpointAggregator");
2367        loop {
2368            // Drain all signatures that arrived since the last iteration into the pending buffer
2369            while let Ok(sig) = self.receiver.try_recv() {
2370                self.pending
2371                    .entry(sig.summary.sequence_number)
2372                    .or_default()
2373                    .push(sig);
2374            }
2375
2376            if let Err(e) = self.run_and_notify().await {
2377                error!(
2378                    "Error while aggregating checkpoint, will retry in 1s: {:?}",
2379                    e
2380                );
2381                self.metrics.checkpoint_errors.inc();
2382                tokio::time::sleep(Duration::from_secs(1)).await;
2383                continue;
2384            }
2385
2386            tokio::select! {
2387                Some(sig) = self.receiver.recv() => {
2388                    self.pending
2389                        .entry(sig.summary.sequence_number)
2390                        .or_default()
2391                        .push(sig);
2392                }
2393                _ = self.notify.notified() => {}
2394                _ = tokio::time::sleep(Duration::from_secs(1)) => {}
2395            }
2396        }
2397    }
2398
2399    async fn run_and_notify(&mut self) -> SuiResult {
2400        let summaries = self.run_inner()?;
2401        for summary in summaries {
2402            self.output.certified_checkpoint_created(&summary).await?;
2403        }
2404        Ok(())
2405    }
2406
2407    fn run_inner(&mut self) -> SuiResult<Vec<CertifiedCheckpointSummary>> {
2408        let _scope = monitored_scope("CheckpointAggregator");
2409        let mut result = vec![];
2410        'outer: loop {
2411            let next_to_certify = self.next_checkpoint_to_certify()?;
2412            // Discard buffered signatures for checkpoints already certified
2413            // (e.g. certified via StateSync before local aggregation completed).
2414            self.pending.retain(|&seq, _| seq >= next_to_certify);
2415            let current = if let Some(current) = &mut self.current {
2416                // It's possible that the checkpoint was already certified by
2417                // the rest of the network and we've already received the
2418                // certified checkpoint via StateSync. In this case, we reset
2419                // the current signature aggregator to the next checkpoint to
2420                // be certified
2421                if current.summary.sequence_number < next_to_certify {
2422                    assert_reachable!("skip checkpoint certification");
2423                    self.current = None;
2424                    continue;
2425                }
2426                current
2427            } else {
2428                let Some(summary) = self
2429                    .epoch_store
2430                    .get_built_checkpoint_summary(next_to_certify)?
2431                else {
2432                    return Ok(result);
2433                };
2434                self.current = Some(CheckpointSignatureAggregator {
2435                    digest: summary.digest(),
2436                    summary,
2437                    signatures_by_digest: MultiStakeAggregator::new(
2438                        self.epoch_store.committee().clone(),
2439                    ),
2440                    store: self.store.clone(),
2441                    state: self.state.clone(),
2442                    metrics: self.metrics.clone(),
2443                });
2444                self.current.as_mut().unwrap()
2445            };
2446
2447            let seq = current.summary.sequence_number;
2448            let sigs = self.pending.remove(&seq).unwrap_or_default();
2449            if sigs.is_empty() {
2450                trace!(
2451                    checkpoint_seq =? seq,
2452                    "Not enough checkpoint signatures",
2453                );
2454                return Ok(result);
2455            }
2456            for data in sigs {
2457                trace!(
2458                    checkpoint_seq = seq,
2459                    "Processing signature for checkpoint (digest: {:?}) from {:?}",
2460                    current.summary.digest(),
2461                    data.summary.auth_sig().authority.concise()
2462                );
2463                self.metrics
2464                    .checkpoint_participation
2465                    .with_label_values(&[&format!(
2466                        "{:?}",
2467                        data.summary.auth_sig().authority.concise()
2468                    )])
2469                    .inc();
2470                if let Ok(auth_signature) = current.try_aggregate(data) {
2471                    debug!(
2472                        checkpoint_seq = seq,
2473                        "Successfully aggregated signatures for checkpoint (digest: {:?})",
2474                        current.summary.digest(),
2475                    );
2476                    let summary = VerifiedCheckpoint::new_unchecked(
2477                        CertifiedCheckpointSummary::new_from_data_and_sig(
2478                            current.summary.clone(),
2479                            auth_signature,
2480                        ),
2481                    );
2482
2483                    self.store.insert_certified_checkpoint(&summary)?;
2484                    self.metrics.last_certified_checkpoint.set(seq as i64);
2485                    current.summary.report_checkpoint_age(
2486                        &self.metrics.last_certified_checkpoint_age,
2487                        &self.metrics.last_certified_checkpoint_age_ms,
2488                    );
2489                    result.push(summary.into_inner());
2490                    self.current = None;
2491                    continue 'outer;
2492                }
2493            }
2494            break;
2495        }
2496        Ok(result)
2497    }
2498
2499    fn next_checkpoint_to_certify(&self) -> SuiResult<CheckpointSequenceNumber> {
2500        Ok(self
2501            .store
2502            .tables
2503            .certified_checkpoints
2504            .reversed_safe_iter_with_bounds(None, None)?
2505            .next()
2506            .transpose()?
2507            .map(|(seq, _)| seq + 1)
2508            .unwrap_or_default())
2509    }
2510}
2511
2512impl CheckpointSignatureAggregator {
2513    #[allow(clippy::result_unit_err)]
2514    pub fn try_aggregate(
2515        &mut self,
2516        data: CheckpointSignatureMessage,
2517    ) -> Result<AuthorityStrongQuorumSignInfo, ()> {
2518        let their_digest = *data.summary.digest();
2519        let (_, signature) = data.summary.into_data_and_sig();
2520        let author = signature.authority;
2521        let envelope =
2522            SignedCheckpointSummary::new_from_data_and_sig(self.summary.clone(), signature);
2523        match self.signatures_by_digest.insert(their_digest, envelope) {
2524            // ignore repeated signatures
2525            InsertResult::Failed { error }
2526                if matches!(
2527                    error.as_inner(),
2528                    SuiErrorKind::StakeAggregatorRepeatedSigner {
2529                        conflicting_sig: false,
2530                        ..
2531                    },
2532                ) =>
2533            {
2534                Err(())
2535            }
2536            InsertResult::Failed { error } => {
2537                warn!(
2538                    checkpoint_seq = self.summary.sequence_number,
2539                    "Failed to aggregate new signature from validator {:?}: {:?}",
2540                    author.concise(),
2541                    error
2542                );
2543                self.check_for_split_brain();
2544                Err(())
2545            }
2546            InsertResult::QuorumReached(cert) => {
2547                // It is not guaranteed that signature.authority == narwhal_cert.author, but we do verify
2548                // the signature so we know that the author signed the message at some point.
2549                if their_digest != self.digest {
2550                    self.metrics.remote_checkpoint_forks.inc();
2551                    warn!(
2552                        checkpoint_seq = self.summary.sequence_number,
2553                        "Validator {:?} has mismatching checkpoint digest {}, we have digest {}",
2554                        author.concise(),
2555                        their_digest,
2556                        self.digest
2557                    );
2558                    return Err(());
2559                }
2560                Ok(cert)
2561            }
2562            InsertResult::NotEnoughVotes {
2563                bad_votes: _,
2564                bad_authorities: _,
2565            } => {
2566                self.check_for_split_brain();
2567                Err(())
2568            }
2569        }
2570    }
2571
2572    /// Check if there is a split brain condition in checkpoint signature aggregation, defined
2573    /// as any state wherein it is no longer possible to achieve quorum on a checkpoint proposal,
2574    /// irrespective of the outcome of any outstanding votes.
2575    fn check_for_split_brain(&self) {
2576        debug!(
2577            checkpoint_seq = self.summary.sequence_number,
2578            "Checking for split brain condition"
2579        );
2580        if self.signatures_by_digest.quorum_unreachable() {
2581            // TODO: at this point we should immediately halt processing
2582            // of new transaction certificates to avoid building on top of
2583            // forked output
2584            // self.halt_all_execution();
2585
2586            let all_unique_values = self.signatures_by_digest.get_all_unique_values();
2587            let digests_by_stake_messages = all_unique_values
2588                .iter()
2589                .sorted_by_key(|(_, (_, stake))| -(*stake as i64))
2590                .map(|(digest, (_authorities, total_stake))| {
2591                    format!("{:?} (total stake: {})", digest, total_stake)
2592                })
2593                .collect::<Vec<String>>();
2594            fail_point_arg!("kill_split_brain_node", |(
2595                checkpoint_overrides,
2596                forked_authorities,
2597            ): (
2598                std::sync::Arc<std::sync::Mutex<std::collections::BTreeMap<u64, String>>>,
2599                std::sync::Arc<std::sync::Mutex<std::collections::HashSet<AuthorityName>>>,
2600            )| {
2601                #[cfg(msim)]
2602                {
2603                    if let (Ok(mut overrides), Ok(forked_authorities_set)) =
2604                        (checkpoint_overrides.lock(), forked_authorities.lock())
2605                    {
2606                        // Find the digest produced by non-forked authorities
2607                        let correct_digest = all_unique_values
2608                            .iter()
2609                            .find(|(_, (authorities, _))| {
2610                                // Check if any authority that produced this digest is NOT in the forked set
2611                                authorities
2612                                    .iter()
2613                                    .any(|auth| !forked_authorities_set.contains(auth))
2614                            })
2615                            .map(|(digest, _)| digest.to_string())
2616                            .unwrap_or_else(|| {
2617                                // Fallback: use the digest with the highest stake
2618                                all_unique_values
2619                                    .iter()
2620                                    .max_by_key(|(_, (_, stake))| *stake)
2621                                    .map(|(digest, _)| digest.to_string())
2622                                    .unwrap_or_else(|| self.digest.to_string())
2623                            });
2624
2625                        overrides.insert(self.summary.sequence_number, correct_digest.clone());
2626
2627                        tracing::error!(
2628                            fatal = true,
2629                            "Fork recovery test: detected split-brain for sequence number: {}, using digest: {}",
2630                            self.summary.sequence_number,
2631                            correct_digest
2632                        );
2633                    }
2634                }
2635            });
2636
2637            debug_fatal!(
2638                "Split brain detected in checkpoint signature aggregation for checkpoint {:?}. Remaining stake: {:?}, Digests by stake: {:?}",
2639                self.summary.sequence_number,
2640                self.signatures_by_digest.uncommitted_stake(),
2641                digests_by_stake_messages
2642            );
2643            self.metrics.split_brain_checkpoint_forks.inc();
2644
2645            let all_unique_values = self.signatures_by_digest.get_all_unique_values();
2646            let local_summary = self.summary.clone();
2647            let state = self.state.clone();
2648            let tables = self.store.clone();
2649
2650            tokio::spawn(async move {
2651                diagnose_split_brain(all_unique_values, local_summary, state, tables).await;
2652            });
2653        }
2654    }
2655}
2656
2657/// Create data dump containing relevant data for diagnosing cause of the
2658/// split brain by querying one disagreeing validator for full checkpoint contents.
2659/// To minimize peer chatter, we only query one validator at random from each
2660/// disagreeing faction, as all honest validators that participated in this round may
2661/// inevitably run the same process.
2662async fn diagnose_split_brain(
2663    all_unique_values: BTreeMap<CheckpointDigest, (Vec<AuthorityName>, StakeUnit)>,
2664    local_summary: CheckpointSummary,
2665    state: Arc<AuthorityState>,
2666    tables: Arc<CheckpointStore>,
2667) {
2668    debug!(
2669        checkpoint_seq = local_summary.sequence_number,
2670        "Running split brain diagnostics..."
2671    );
2672    let time = SystemTime::now();
2673    // collect one random disagreeing validator per differing digest
2674    let digest_to_validator = all_unique_values
2675        .iter()
2676        .filter_map(|(digest, (validators, _))| {
2677            if *digest != local_summary.digest() {
2678                let random_validator = validators.choose(&mut get_rng()).unwrap();
2679                Some((*digest, *random_validator))
2680            } else {
2681                None
2682            }
2683        })
2684        .collect::<HashMap<_, _>>();
2685    if digest_to_validator.is_empty() {
2686        panic!(
2687            "Given split brain condition, there should be at \
2688                least one validator that disagrees with local signature"
2689        );
2690    }
2691
2692    let epoch_store = state.load_epoch_store_one_call_per_task();
2693    let committee = epoch_store
2694        .epoch_start_state()
2695        .get_sui_committee_with_network_metadata();
2696    let network_config = default_mysten_network_config();
2697    let network_clients =
2698        make_network_authority_clients_with_network_config(&committee, &network_config);
2699
2700    // Query all disagreeing validators
2701    let response_futures = digest_to_validator
2702        .values()
2703        .cloned()
2704        .map(|validator| {
2705            let client = network_clients
2706                .get(&validator)
2707                .expect("Failed to get network client");
2708            let request = CheckpointRequestV2 {
2709                sequence_number: Some(local_summary.sequence_number),
2710                request_content: true,
2711                certified: false,
2712            };
2713            client.handle_checkpoint_v2(request)
2714        })
2715        .collect::<Vec<_>>();
2716
2717    let digest_name_pair = digest_to_validator.iter();
2718    let response_data = futures::future::join_all(response_futures)
2719        .await
2720        .into_iter()
2721        .zip_debug_eq(digest_name_pair)
2722        .filter_map(|(response, (digest, name))| match response {
2723            Ok(response) => match response {
2724                CheckpointResponseV2 {
2725                    checkpoint: Some(CheckpointSummaryResponse::Pending(summary)),
2726                    contents: Some(contents),
2727                } => Some((*name, *digest, summary, contents)),
2728                CheckpointResponseV2 {
2729                    checkpoint: Some(CheckpointSummaryResponse::Certified(_)),
2730                    contents: _,
2731                } => {
2732                    panic!("Expected pending checkpoint, but got certified checkpoint");
2733                }
2734                CheckpointResponseV2 {
2735                    checkpoint: None,
2736                    contents: _,
2737                } => {
2738                    error!(
2739                        "Summary for checkpoint {:?} not found on validator {:?}",
2740                        local_summary.sequence_number, name
2741                    );
2742                    None
2743                }
2744                CheckpointResponseV2 {
2745                    checkpoint: _,
2746                    contents: None,
2747                } => {
2748                    error!(
2749                        "Contents for checkpoint {:?} not found on validator {:?}",
2750                        local_summary.sequence_number, name
2751                    );
2752                    None
2753                }
2754            },
2755            Err(e) => {
2756                error!(
2757                    "Failed to get checkpoint contents from validator for fork diagnostics: {:?}",
2758                    e
2759                );
2760                None
2761            }
2762        })
2763        .collect::<Vec<_>>();
2764
2765    let local_checkpoint_contents = tables
2766        .get_checkpoint_contents(&local_summary.content_digest)
2767        .unwrap_or_else(|_| {
2768            panic!(
2769                "Could not find checkpoint contents for digest {:?}",
2770                local_summary.digest()
2771            )
2772        })
2773        .unwrap_or_else(|| {
2774            panic!(
2775                "Could not find local full checkpoint contents for checkpoint {:?}, digest {:?}",
2776                local_summary.sequence_number,
2777                local_summary.digest()
2778            )
2779        });
2780    let local_contents_text = format!("{local_checkpoint_contents:?}");
2781
2782    let local_summary_text = format!("{local_summary:?}");
2783    let local_validator = state.name.concise();
2784    let diff_patches = response_data
2785        .iter()
2786        .map(|(name, other_digest, other_summary, contents)| {
2787            let other_contents_text = format!("{contents:?}");
2788            let other_summary_text = format!("{other_summary:?}");
2789            let (local_transactions, local_effects): (Vec<_>, Vec<_>) = local_checkpoint_contents
2790                .enumerate_transactions(&local_summary)
2791                .map(|(_, exec_digest)| (exec_digest.transaction, exec_digest.effects))
2792                .unzip();
2793            let (other_transactions, other_effects): (Vec<_>, Vec<_>) = contents
2794                .enumerate_transactions(other_summary)
2795                .map(|(_, exec_digest)| (exec_digest.transaction, exec_digest.effects))
2796                .unzip();
2797            let summary_patch = create_patch(&local_summary_text, &other_summary_text);
2798            let contents_patch = create_patch(&local_contents_text, &other_contents_text);
2799            let local_transactions_text = format!("{local_transactions:#?}");
2800            let other_transactions_text = format!("{other_transactions:#?}");
2801            let transactions_patch =
2802                create_patch(&local_transactions_text, &other_transactions_text);
2803            let local_effects_text = format!("{local_effects:#?}");
2804            let other_effects_text = format!("{other_effects:#?}");
2805            let effects_patch = create_patch(&local_effects_text, &other_effects_text);
2806            let seq_number = local_summary.sequence_number;
2807            let local_digest = local_summary.digest();
2808            let other_validator = name.concise();
2809            format!(
2810                "Checkpoint: {seq_number:?}\n\
2811                Local validator (original): {local_validator:?}, digest: {local_digest:?}\n\
2812                Other validator (modified): {other_validator:?}, digest: {other_digest:?}\n\n\
2813                Summary Diff: \n{summary_patch}\n\n\
2814                Contents Diff: \n{contents_patch}\n\n\
2815                Transactions Diff: \n{transactions_patch}\n\n\
2816                Effects Diff: \n{effects_patch}",
2817            )
2818        })
2819        .collect::<Vec<_>>()
2820        .join("\n\n\n");
2821
2822    let header = format!(
2823        "Checkpoint Fork Dump - Authority {local_validator:?}: \n\
2824        Datetime: {:?}",
2825        time
2826    );
2827    let fork_logs_text = format!("{header}\n\n{diff_patches}\n\n");
2828    let path = tempfile::tempdir()
2829        .expect("Failed to create tempdir")
2830        .keep()
2831        .join(Path::new("checkpoint_fork_dump.txt"));
2832    let mut file = File::create(path).unwrap();
2833    write!(file, "{}", fork_logs_text).unwrap();
2834    debug!("{}", fork_logs_text);
2835}
2836
2837pub trait CheckpointServiceNotify {
2838    fn notify_checkpoint_signature(&self, info: &CheckpointSignatureMessage) -> SuiResult;
2839
2840    fn notify_checkpoint(&self) -> SuiResult;
2841}
2842
2843#[allow(clippy::large_enum_variant)]
2844enum CheckpointServiceState {
2845    Unstarted(
2846        (
2847            CheckpointBuilder,
2848            CheckpointAggregator,
2849            CheckpointStateHasher,
2850        ),
2851    ),
2852    Started,
2853}
2854
2855impl CheckpointServiceState {
2856    fn take_unstarted(
2857        &mut self,
2858    ) -> (
2859        CheckpointBuilder,
2860        CheckpointAggregator,
2861        CheckpointStateHasher,
2862    ) {
2863        let mut state = CheckpointServiceState::Started;
2864        std::mem::swap(self, &mut state);
2865
2866        match state {
2867            CheckpointServiceState::Unstarted((builder, aggregator, hasher)) => {
2868                (builder, aggregator, hasher)
2869            }
2870            CheckpointServiceState::Started => panic!("CheckpointServiceState is already started"),
2871        }
2872    }
2873}
2874
2875pub struct CheckpointService {
2876    tables: Arc<CheckpointStore>,
2877    notify_builder: Arc<Notify>,
2878    signature_sender: mpsc::UnboundedSender<CheckpointSignatureMessage>,
2879    // A notification for the current highest built sequence number.
2880    highest_currently_built_seq_tx: watch::Sender<CheckpointSequenceNumber>,
2881    // The highest sequence number that had already been built at the time CheckpointService
2882    // was constructed
2883    highest_previously_built_seq: CheckpointSequenceNumber,
2884    metrics: Arc<CheckpointMetrics>,
2885    state: Mutex<CheckpointServiceState>,
2886}
2887
2888impl CheckpointService {
2889    /// Constructs a new CheckpointService in an un-started state.
2890    // The signature channel is unbounded because notify_checkpoint_signature is called from a
2891    // sync context (consensus_validator.rs implements a sync external trait) and cannot block.
2892    // The channel is consumed by a single async aggregator task that drains it continuously, so
2893    // unbounded growth is not a concern in practice.
2894    #[allow(clippy::disallowed_methods)]
2895    pub fn build(
2896        state: Arc<AuthorityState>,
2897        checkpoint_store: Arc<CheckpointStore>,
2898        epoch_store: Arc<AuthorityPerEpochStore>,
2899        effects_store: Arc<dyn TransactionCacheRead>,
2900        global_state_hasher: Weak<GlobalStateHasher>,
2901        checkpoint_output: Box<dyn CheckpointOutput>,
2902        certified_checkpoint_output: Box<dyn CertifiedCheckpointOutput>,
2903        metrics: Arc<CheckpointMetrics>,
2904    ) -> Arc<Self> {
2905        info!("Starting checkpoint service");
2906        Self::initialize_accumulator_account_metrics(&state, &epoch_store, &metrics);
2907        let notify_builder = Arc::new(Notify::new());
2908        let notify_aggregator = Arc::new(Notify::new());
2909
2910        // We may have built higher checkpoint numbers before restarting.
2911        let highest_previously_built_seq = checkpoint_store
2912            .get_latest_locally_computed_checkpoint()
2913            .expect("failed to get latest locally computed checkpoint")
2914            .map(|s| s.sequence_number)
2915            .unwrap_or(0);
2916
2917        let highest_currently_built_seq =
2918            CheckpointBuilder::load_last_built_checkpoint_summary(&epoch_store, &checkpoint_store)
2919                .expect("epoch should not have ended")
2920                .map(|(seq, _)| seq)
2921                .unwrap_or(0);
2922
2923        let (highest_currently_built_seq_tx, _) = watch::channel(highest_currently_built_seq);
2924
2925        let (signature_sender, signature_receiver) = mpsc::unbounded_channel();
2926
2927        let aggregator = CheckpointAggregator::new(
2928            checkpoint_store.clone(),
2929            epoch_store.clone(),
2930            notify_aggregator.clone(),
2931            signature_receiver,
2932            certified_checkpoint_output,
2933            state.clone(),
2934            metrics.clone(),
2935        );
2936
2937        let (send_to_hasher, receive_from_builder) = mpsc::channel(16);
2938
2939        let ckpt_state_hasher = CheckpointStateHasher::new(
2940            epoch_store.clone(),
2941            global_state_hasher.clone(),
2942            receive_from_builder,
2943        );
2944
2945        let builder = CheckpointBuilder::new(
2946            state.clone(),
2947            checkpoint_store.clone(),
2948            epoch_store.clone(),
2949            notify_builder.clone(),
2950            effects_store,
2951            global_state_hasher,
2952            send_to_hasher,
2953            checkpoint_output,
2954            notify_aggregator.clone(),
2955            highest_currently_built_seq_tx.clone(),
2956            metrics.clone(),
2957        );
2958
2959        Arc::new(Self {
2960            tables: checkpoint_store,
2961            notify_builder,
2962            signature_sender,
2963            highest_currently_built_seq_tx,
2964            highest_previously_built_seq,
2965            metrics,
2966            state: Mutex::new(CheckpointServiceState::Unstarted((
2967                builder,
2968                aggregator,
2969                ckpt_state_hasher,
2970            ))),
2971        })
2972    }
2973
2974    fn initialize_accumulator_account_metrics(
2975        state: &AuthorityState,
2976        epoch_store: &AuthorityPerEpochStore,
2977        metrics: &CheckpointMetrics,
2978    ) {
2979        if !epoch_store.protocol_config().enable_accumulators() {
2980            return;
2981        }
2982
2983        let object_store = state.get_object_store();
2984        match accumulator_metadata::get_accumulator_object_count(object_store.as_ref()) {
2985            Ok(Some(count)) => metrics.initialize_accumulator_accounts_live(count),
2986            Ok(None) => {}
2987            Err(e) => fatal!("failed to initialize accumulator account metrics: {e}"),
2988        }
2989    }
2990
2991    /// Starts the CheckpointService.
2992    ///
2993    /// This function blocks until the CheckpointBuilder re-builds all checkpoints that had
2994    /// been built before the most recent restart. You can think of this as a WAL replay
2995    /// operation. Upon startup, we may have a number of consensus commits and resulting
2996    /// checkpoints that were built but not committed to disk. We want to reprocess the
2997    /// commits and rebuild the checkpoints before starting normal operation.
2998    pub async fn spawn(
2999        &self,
3000        epoch_store: Arc<AuthorityPerEpochStore>,
3001        consensus_replay_waiter: Option<ReplayWaiter>,
3002    ) {
3003        let (builder, aggregator, state_hasher) = self.state.lock().take_unstarted();
3004
3005        // Clean up state hashes computed after the last built checkpoint
3006        // This prevents ECMH divergence after fork recovery restarts
3007
3008        // Note: there is a rare crash recovery edge case where we write the builder
3009        // summary, but crash before we can bump the highest executed checkpoint.
3010        // If we committed the builder summary, it was certified and unforked, so there
3011        // is no need to clear that state hash. If we do clear it, then checkpoint executor
3012        // will wait forever for checkpoint builder to produce the state hash, which will
3013        // never happen.
3014        let last_persisted_builder_seq = epoch_store
3015            .last_persisted_checkpoint_builder_summary()
3016            .expect("epoch should not have ended")
3017            .map(|s| s.summary.sequence_number);
3018
3019        let last_executed_seq = self
3020            .tables
3021            .get_highest_executed_checkpoint()
3022            .expect("Failed to get highest executed checkpoint")
3023            .map(|checkpoint| *checkpoint.sequence_number());
3024
3025        if let Some(last_committed_seq) = last_persisted_builder_seq.max(last_executed_seq) {
3026            if let Err(e) = builder
3027                .epoch_store
3028                .clear_state_hashes_after_checkpoint(last_committed_seq)
3029            {
3030                error!(
3031                    "Failed to clear state hashes after checkpoint {}: {:?}",
3032                    last_committed_seq, e
3033                );
3034            } else {
3035                info!(
3036                    "Cleared state hashes after checkpoint {} to ensure consistent ECMH computation",
3037                    last_committed_seq
3038                );
3039            }
3040        }
3041
3042        let (builder_finished_tx, builder_finished_rx) = tokio::sync::oneshot::channel();
3043
3044        let state_hasher_task = spawn_monitored_task!(state_hasher.run());
3045        let aggregator_task = spawn_monitored_task!(aggregator.run());
3046
3047        spawn_monitored_task!(async move {
3048            epoch_store
3049                .within_alive_epoch(async move {
3050                    builder.run(consensus_replay_waiter).await;
3051                    builder_finished_tx.send(()).ok();
3052                })
3053                .await
3054                .ok();
3055
3056            // state hasher will terminate as soon as it has finished processing all messages from builder
3057            state_hasher_task
3058                .await
3059                .expect("state hasher should exit normally");
3060
3061            // builder must shut down before aggregator and state_hasher, since it sends
3062            // messages to them
3063            aggregator_task.abort();
3064            aggregator_task.await.ok();
3065        });
3066
3067        // If this times out, the validator may still start up. The worst that can
3068        // happen is that we will crash later on instead of immediately. The eventual
3069        // crash would occur because we may be missing transactions that are below the
3070        // highest_synced_checkpoint watermark, which can cause a crash in
3071        // `CheckpointExecutor::extract_randomness_rounds`.
3072        if tokio::time::timeout(Duration::from_secs(120), async move {
3073            tokio::select! {
3074                _ = builder_finished_rx => { debug!("CheckpointBuilder finished"); }
3075                _ = self.wait_for_rebuilt_checkpoints() => (),
3076            }
3077        })
3078        .await
3079        .is_err()
3080        {
3081            debug_fatal!("Timed out waiting for checkpoints to be rebuilt");
3082        }
3083    }
3084}
3085
3086impl CheckpointService {
3087    /// Waits until all checkpoints had been built before the node restarted
3088    /// are rebuilt. This is required to preserve the invariant that all checkpoints
3089    /// (and their transactions) below the highest_synced_checkpoint watermark are
3090    /// available. Once the checkpoints are constructed, we can be sure that the
3091    /// transactions have also been executed.
3092    pub async fn wait_for_rebuilt_checkpoints(&self) {
3093        let highest_previously_built_seq = self.highest_previously_built_seq;
3094        let mut rx = self.highest_currently_built_seq_tx.subscribe();
3095        let mut highest_currently_built_seq = *rx.borrow_and_update();
3096        info!(
3097            "Waiting for checkpoints to be rebuilt, previously built seq: {highest_previously_built_seq}, currently built seq: {highest_currently_built_seq}"
3098        );
3099        loop {
3100            if highest_currently_built_seq >= highest_previously_built_seq {
3101                info!("Checkpoint rebuild complete");
3102                break;
3103            }
3104            rx.changed().await.unwrap();
3105            highest_currently_built_seq = *rx.borrow_and_update();
3106        }
3107    }
3108
3109    #[cfg(test)]
3110    fn write_and_notify_checkpoint_for_testing(
3111        &self,
3112        epoch_store: &AuthorityPerEpochStore,
3113        checkpoint: PendingCheckpoint,
3114    ) -> SuiResult {
3115        use crate::authority::authority_per_epoch_store::consensus_quarantine::ConsensusCommitOutput;
3116
3117        let mut output = ConsensusCommitOutput::new(0);
3118        epoch_store.write_pending_checkpoint(&mut output, &checkpoint);
3119        output.set_default_commit_stats_for_testing();
3120        epoch_store.push_consensus_output_for_tests(output);
3121        self.notify_checkpoint()?;
3122        Ok(())
3123    }
3124}
3125
3126impl CheckpointServiceNotify for CheckpointService {
3127    fn notify_checkpoint_signature(&self, info: &CheckpointSignatureMessage) -> SuiResult {
3128        let sequence = info.summary.sequence_number;
3129        let signer = info.summary.auth_sig().authority.concise();
3130
3131        if let Some(highest_verified_checkpoint) = self
3132            .tables
3133            .get_highest_verified_checkpoint()?
3134            .map(|x| *x.sequence_number())
3135            && sequence <= highest_verified_checkpoint
3136        {
3137            trace!(
3138                checkpoint_seq = sequence,
3139                "Ignore checkpoint signature from {} - already certified", signer,
3140            );
3141            self.metrics
3142                .last_ignored_checkpoint_signature_received
3143                .set(sequence as i64);
3144            return Ok(());
3145        }
3146        trace!(
3147            checkpoint_seq = sequence,
3148            "Received checkpoint signature, digest {} from {}",
3149            info.summary.digest(),
3150            signer,
3151        );
3152        self.metrics
3153            .last_received_checkpoint_signatures
3154            .with_label_values(&[&signer.to_string()])
3155            .set(sequence as i64);
3156        self.signature_sender.send(info.clone()).ok();
3157        Ok(())
3158    }
3159
3160    fn notify_checkpoint(&self) -> SuiResult {
3161        self.notify_builder.notify_one();
3162        Ok(())
3163    }
3164}
3165
3166// test helper
3167pub struct CheckpointServiceNoop {}
3168impl CheckpointServiceNotify for CheckpointServiceNoop {
3169    fn notify_checkpoint_signature(&self, _: &CheckpointSignatureMessage) -> SuiResult {
3170        Ok(())
3171    }
3172
3173    fn notify_checkpoint(&self) -> SuiResult {
3174        Ok(())
3175    }
3176}
3177
3178impl PendingCheckpoint {
3179    pub fn height(&self) -> CheckpointHeight {
3180        self.details.checkpoint_height
3181    }
3182
3183    pub(crate) fn num_roots(&self) -> usize {
3184        self.roots.iter().map(|r| r.tx_roots.len()).sum()
3185    }
3186}
3187
3188pin_project! {
3189    pub struct PollCounter<Fut> {
3190        #[pin]
3191        future: Fut,
3192        count: usize,
3193    }
3194}
3195
3196impl<Fut> PollCounter<Fut> {
3197    pub fn new(future: Fut) -> Self {
3198        Self { future, count: 0 }
3199    }
3200
3201    pub fn count(&self) -> usize {
3202        self.count
3203    }
3204}
3205
3206impl<Fut: Future> Future for PollCounter<Fut> {
3207    type Output = (usize, Fut::Output);
3208
3209    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
3210        let this = self.project();
3211        *this.count += 1;
3212        match this.future.poll(cx) {
3213            Poll::Ready(output) => Poll::Ready((*this.count, output)),
3214            Poll::Pending => Poll::Pending,
3215        }
3216    }
3217}
3218
3219fn poll_count<Fut>(future: Fut) -> PollCounter<Fut> {
3220    PollCounter::new(future)
3221}
3222
3223#[cfg(test)]
3224mod tests {
3225    use super::*;
3226    use crate::authority::test_authority_builder::TestAuthorityBuilder;
3227    use futures::FutureExt as _;
3228    use futures::future::BoxFuture;
3229    use std::collections::HashMap;
3230    use std::ops::Deref;
3231    use sui_macros::sim_test;
3232    use sui_protocol_config::{Chain, ProtocolConfig};
3233    use sui_types::accumulator_event::AccumulatorEvent;
3234    use sui_types::base_types::{SequenceNumber, TransactionEffectsDigest};
3235    use sui_types::crypto::Signature;
3236    use sui_types::effects::{TransactionEffects, TransactionEvents};
3237    use sui_types::messages_checkpoint::SignedCheckpointSummary;
3238    use sui_types::transaction::VerifiedTransaction;
3239    use tokio::sync::mpsc;
3240
3241    #[tokio::test]
3242    async fn test_clear_locally_computed_checkpoints_from_deletes_inclusive_range() {
3243        let store = CheckpointStore::new_for_tests();
3244        let protocol = sui_protocol_config::ProtocolConfig::get_for_max_version_UNSAFE();
3245        for seq in 70u64..=80u64 {
3246            let contents =
3247                sui_types::messages_checkpoint::CheckpointContents::new_with_digests_only_for_tests(
3248                    [sui_types::base_types::ExecutionDigests::new(
3249                        sui_types::digests::TransactionDigest::random(),
3250                        sui_types::digests::TransactionEffectsDigest::ZERO,
3251                    )],
3252                );
3253            let summary = sui_types::messages_checkpoint::CheckpointSummary::new(
3254                &protocol,
3255                0,
3256                seq,
3257                0,
3258                &contents,
3259                None,
3260                sui_types::gas::GasCostSummary::default(),
3261                None,
3262                0,
3263                Vec::new(),
3264                Vec::new(),
3265            );
3266            store
3267                .tables
3268                .locally_computed_checkpoints
3269                .insert(&seq, &summary)
3270                .unwrap();
3271        }
3272
3273        store
3274            .clear_locally_computed_checkpoints_from(76)
3275            .expect("clear should succeed");
3276
3277        // Explicit boundary checks: 75 must remain, 76 must be deleted
3278        assert!(
3279            store
3280                .tables
3281                .locally_computed_checkpoints
3282                .get(&75)
3283                .unwrap()
3284                .is_some()
3285        );
3286        assert!(
3287            store
3288                .tables
3289                .locally_computed_checkpoints
3290                .get(&76)
3291                .unwrap()
3292                .is_none()
3293        );
3294
3295        for seq in 70u64..76u64 {
3296            assert!(
3297                store
3298                    .tables
3299                    .locally_computed_checkpoints
3300                    .get(&seq)
3301                    .unwrap()
3302                    .is_some()
3303            );
3304        }
3305        for seq in 76u64..=80u64 {
3306            assert!(
3307                store
3308                    .tables
3309                    .locally_computed_checkpoints
3310                    .get(&seq)
3311                    .unwrap()
3312                    .is_none()
3313            );
3314        }
3315    }
3316
3317    #[tokio::test]
3318    async fn test_fork_detection_storage() {
3319        let store = CheckpointStore::new_for_tests();
3320        store.set_binary_version("v1");
3321        // checkpoint fork
3322        let seq_num = 42;
3323        let digest = CheckpointDigest::random();
3324        let certified_digest = CheckpointDigest::random();
3325
3326        assert!(store.get_checkpoint_fork_detected().unwrap().is_none());
3327
3328        store
3329            .record_checkpoint_fork_detected(seq_num, digest, Some(certified_digest))
3330            .unwrap();
3331
3332        let retrieved = store.get_checkpoint_fork_detected().unwrap();
3333        assert!(retrieved.is_some());
3334        let info = retrieved.unwrap();
3335        assert_eq!(info.checkpoint_seq, seq_num);
3336        assert_eq!(info.checkpoint_digest, digest);
3337        assert_eq!(info.certified_checkpoint_digest, Some(certified_digest));
3338        assert_eq!(info.binary_version, "v1");
3339
3340        store.clear_checkpoint_fork_detected().unwrap();
3341        assert!(store.get_checkpoint_fork_detected().unwrap().is_none());
3342
3343        // txn fork
3344        let tx_digest = TransactionDigest::random();
3345        let expected_effects = TransactionEffectsDigest::random();
3346        let actual_effects = TransactionEffectsDigest::random();
3347
3348        assert!(store.get_transaction_fork_detected().unwrap().is_none());
3349
3350        store
3351            .record_transaction_fork_detected(tx_digest, expected_effects, actual_effects, Some(7))
3352            .unwrap();
3353
3354        let retrieved = store.get_transaction_fork_detected().unwrap();
3355        assert!(retrieved.is_some());
3356        let info = retrieved.unwrap();
3357        assert_eq!(info.tx_digest, tx_digest);
3358        assert_eq!(info.expected_effects_digest, expected_effects);
3359        assert_eq!(info.actual_effects_digest, actual_effects);
3360        assert_eq!(info.certified_checkpoint_seq, Some(7));
3361        assert_eq!(info.binary_version, "v1");
3362
3363        store.clear_transaction_fork_detected().unwrap();
3364        assert!(store.get_transaction_fork_detected().unwrap().is_none());
3365    }
3366
3367    #[sim_test]
3368    pub async fn checkpoint_builder_test() {
3369        telemetry_subscribers::init_for_testing();
3370
3371        let mut protocol_config =
3372            ProtocolConfig::get_for_version(ProtocolVersion::max(), Chain::Unknown);
3373        protocol_config.disable_accumulators_for_testing();
3374        // This fixture supplies historical effects dependencies rather than consensus input order.
3375        protocol_config.set_disable_effects_tx_dependencies_for_testing(false);
3376        let state = TestAuthorityBuilder::new()
3377            .with_protocol_config(protocol_config)
3378            .build()
3379            .await;
3380
3381        let dummy_tx = VerifiedTransaction::new_authenticator_state_update(
3382            0,
3383            0,
3384            vec![],
3385            SequenceNumber::new(),
3386        );
3387
3388        for i in 0..20 {
3389            state
3390                .database_for_testing()
3391                .perpetual_tables
3392                .transactions
3393                .insert(&d(i), dummy_tx.serializable_ref())
3394                .unwrap();
3395        }
3396
3397        let mut store = HashMap::<TransactionDigest, TransactionEffects>::new();
3398        commit_cert_for_test(
3399            &mut store,
3400            state.clone(),
3401            d(1),
3402            vec![d(2), d(3)],
3403            GasCostSummary::new(11, 12, 11, 1),
3404        );
3405        commit_cert_for_test(
3406            &mut store,
3407            state.clone(),
3408            d(2),
3409            vec![d(3), d(4)],
3410            GasCostSummary::new(21, 22, 21, 1),
3411        );
3412        commit_cert_for_test(
3413            &mut store,
3414            state.clone(),
3415            d(3),
3416            vec![],
3417            GasCostSummary::new(31, 32, 31, 1),
3418        );
3419        commit_cert_for_test(
3420            &mut store,
3421            state.clone(),
3422            d(4),
3423            vec![],
3424            GasCostSummary::new(41, 42, 41, 1),
3425        );
3426        for i in [5, 6, 7, 10, 11, 12, 13] {
3427            commit_cert_for_test(
3428                &mut store,
3429                state.clone(),
3430                d(i),
3431                vec![],
3432                GasCostSummary::new(41, 42, 41, 1),
3433            );
3434        }
3435        for i in [15, 16, 17] {
3436            commit_cert_for_test(
3437                &mut store,
3438                state.clone(),
3439                d(i),
3440                vec![],
3441                GasCostSummary::new(51, 52, 51, 1),
3442            );
3443        }
3444        let all_digests: Vec<_> = store.keys().copied().collect();
3445        for digest in all_digests {
3446            let signature = Signature::Ed25519SuiSignature(Default::default()).into();
3447            state
3448                .epoch_store_for_testing()
3449                .test_insert_user_signature(digest, vec![(signature, None)]);
3450        }
3451
3452        let (output, mut result) = mpsc::channel::<(CheckpointContents, CheckpointSummary)>(10);
3453        let (certified_output, mut certified_result) =
3454            mpsc::channel::<CertifiedCheckpointSummary>(10);
3455        let store = Arc::new(store);
3456
3457        let ckpt_dir = tempfile::tempdir().unwrap();
3458        let checkpoint_store =
3459            CheckpointStore::new(ckpt_dir.path(), Arc::new(PrunerWatermarks::default()));
3460        let epoch_store = state.epoch_store_for_testing();
3461
3462        let global_state_hasher = Arc::new(GlobalStateHasher::new_for_tests(
3463            state.get_global_state_hash_store().clone(),
3464        ));
3465
3466        let checkpoint_service = CheckpointService::build(
3467            state.clone(),
3468            checkpoint_store,
3469            epoch_store.clone(),
3470            store,
3471            Arc::downgrade(&global_state_hasher),
3472            Box::new(output),
3473            Box::new(certified_output),
3474            CheckpointMetrics::new_for_tests(),
3475        );
3476        checkpoint_service.spawn(epoch_store.clone(), None).await;
3477
3478        checkpoint_service
3479            .write_and_notify_checkpoint_for_testing(&epoch_store, p(0, vec![4], 0))
3480            .unwrap();
3481        checkpoint_service
3482            .write_and_notify_checkpoint_for_testing(&epoch_store, p(1, vec![1, 3], 2000))
3483            .unwrap();
3484        checkpoint_service
3485            .write_and_notify_checkpoint_for_testing(&epoch_store, p(2, vec![10, 11, 12, 13], 3000))
3486            .unwrap();
3487        checkpoint_service
3488            .write_and_notify_checkpoint_for_testing(&epoch_store, p(3, vec![15, 16, 17], 4000))
3489            .unwrap();
3490        checkpoint_service
3491            .write_and_notify_checkpoint_for_testing(&epoch_store, p(4, vec![5], 4001))
3492            .unwrap();
3493        checkpoint_service
3494            .write_and_notify_checkpoint_for_testing(&epoch_store, p(5, vec![6], 5000))
3495            .unwrap();
3496
3497        let (c1c, c1s) = result.recv().await.unwrap();
3498        let (c2c, c2s) = result.recv().await.unwrap();
3499
3500        let c1t = c1c.iter().map(|d| d.transaction).collect::<Vec<_>>();
3501        let c2t = c2c.iter().map(|d| d.transaction).collect::<Vec<_>>();
3502        assert_eq!(c1t, vec![d(4)]);
3503        assert_eq!(c1s.previous_digest, None);
3504        assert_eq!(c1s.sequence_number, 0);
3505        assert_eq!(
3506            c1s.epoch_rolling_gas_cost_summary,
3507            GasCostSummary::new(41, 42, 41, 1)
3508        );
3509
3510        // Causal order places d(3) before d(1), which depends on it.
3511        assert_eq!(c2t, vec![d(3), d(1)]);
3512        assert_eq!(c2s.previous_digest, Some(c1s.digest()));
3513        assert_eq!(c2s.sequence_number, 1);
3514        assert_eq!(
3515            c2s.epoch_rolling_gas_cost_summary,
3516            GasCostSummary::new(83, 86, 83, 3)
3517        );
3518
3519        // Each pending checkpoint produces exactly one checkpoint; splitting is
3520        // done in the consensus handler.
3521        let (c3c, c3s) = result.recv().await.unwrap();
3522        let c3t = c3c.iter().map(|d| d.transaction).collect::<Vec<_>>();
3523        assert_eq!(c3s.sequence_number, 2);
3524        assert_eq!(c3s.previous_digest, Some(c2s.digest()));
3525        assert_eq!(c3t, vec![d(10), d(11), d(12), d(13)]);
3526
3527        let (c4c, c4s) = result.recv().await.unwrap();
3528        let c4t = c4c.iter().map(|d| d.transaction).collect::<Vec<_>>();
3529        assert_eq!(c4s.sequence_number, 3);
3530        assert_eq!(c4s.previous_digest, Some(c3s.digest()));
3531        assert_eq!(c4t, vec![d(15), d(16), d(17)]);
3532
3533        let (c5c, c5s) = result.recv().await.unwrap();
3534        let c5t = c5c.iter().map(|d| d.transaction).collect::<Vec<_>>();
3535        assert_eq!(c5s.sequence_number, 4);
3536        assert_eq!(c5s.previous_digest, Some(c4s.digest()));
3537        assert_eq!(c5t, vec![d(5)]);
3538
3539        let (c6c, c6s) = result.recv().await.unwrap();
3540        let c6t = c6c.iter().map(|d| d.transaction).collect::<Vec<_>>();
3541        assert_eq!(c6s.sequence_number, 5);
3542        assert_eq!(c6s.previous_digest, Some(c5s.digest()));
3543        assert_eq!(c6t, vec![d(6)]);
3544
3545        let c1ss = SignedCheckpointSummary::new(c1s.epoch, c1s, state.secret.deref(), state.name);
3546        let c2ss = SignedCheckpointSummary::new(c2s.epoch, c2s, state.secret.deref(), state.name);
3547
3548        checkpoint_service
3549            .notify_checkpoint_signature(&CheckpointSignatureMessage { summary: c2ss })
3550            .unwrap();
3551        checkpoint_service
3552            .notify_checkpoint_signature(&CheckpointSignatureMessage { summary: c1ss })
3553            .unwrap();
3554
3555        let c1sc = certified_result.recv().await.unwrap();
3556        let c2sc = certified_result.recv().await.unwrap();
3557        assert_eq!(c1sc.sequence_number, 0);
3558        assert_eq!(c2sc.sequence_number, 1);
3559    }
3560
3561    impl TransactionCacheRead for HashMap<TransactionDigest, TransactionEffects> {
3562        fn notify_read_executed_effects_may_fail(
3563            &self,
3564            _: &str,
3565            digests: &[TransactionDigest],
3566        ) -> BoxFuture<'_, SuiResult<Vec<TransactionEffects>>> {
3567            std::future::ready(Ok(digests
3568                .iter()
3569                .map(|d| self.get(d).expect("effects not found").clone())
3570                .collect()))
3571            .boxed()
3572        }
3573
3574        fn notify_read_executed_effects_digests(
3575            &self,
3576            _: &str,
3577            digests: &[TransactionDigest],
3578        ) -> BoxFuture<'_, Vec<TransactionEffectsDigest>> {
3579            std::future::ready(
3580                digests
3581                    .iter()
3582                    .map(|d| {
3583                        self.get(d)
3584                            .map(|fx| fx.digest())
3585                            .expect("effects not found")
3586                    })
3587                    .collect(),
3588            )
3589            .boxed()
3590        }
3591
3592        fn multi_get_executed_effects(
3593            &self,
3594            digests: &[TransactionDigest],
3595        ) -> Vec<Option<TransactionEffects>> {
3596            digests.iter().map(|d| self.get(d).cloned()).collect()
3597        }
3598
3599        // Unimplemented methods - its unfortunate to have this big blob of useless code, but it wasn't
3600        // worth it to keep EffectsNotifyRead around just for these tests, as it caused a ton of
3601        // complication in non-test code. (e.g. had to implement EFfectsNotifyRead for all
3602        // ExecutionCacheRead implementors).
3603
3604        fn multi_get_transaction_blocks(
3605            &self,
3606            _: &[TransactionDigest],
3607        ) -> Vec<Option<Arc<VerifiedTransaction>>> {
3608            unimplemented!()
3609        }
3610
3611        fn multi_get_executed_effects_digests(
3612            &self,
3613            _: &[TransactionDigest],
3614        ) -> Vec<Option<TransactionEffectsDigest>> {
3615            unimplemented!()
3616        }
3617
3618        fn multi_get_effects(
3619            &self,
3620            _: &[TransactionEffectsDigest],
3621        ) -> Vec<Option<TransactionEffects>> {
3622            unimplemented!()
3623        }
3624
3625        fn multi_get_events(&self, _: &[TransactionDigest]) -> Vec<Option<TransactionEvents>> {
3626            unimplemented!()
3627        }
3628
3629        fn take_accumulator_events(&self, _: &TransactionDigest) -> Option<Vec<AccumulatorEvent>> {
3630            unimplemented!()
3631        }
3632
3633        fn get_unchanged_loaded_runtime_objects(
3634            &self,
3635            _digest: &TransactionDigest,
3636        ) -> Option<Vec<sui_types::storage::ObjectKey>> {
3637            unimplemented!()
3638        }
3639
3640        fn transaction_executed_in_last_epoch(&self, _: &TransactionDigest, _: EpochId) -> bool {
3641            unimplemented!()
3642        }
3643    }
3644
3645    #[async_trait::async_trait]
3646    impl CheckpointOutput for mpsc::Sender<(CheckpointContents, CheckpointSummary)> {
3647        async fn checkpoint_created(
3648            &self,
3649            summary: &CheckpointSummary,
3650            contents: &CheckpointContents,
3651            _epoch_store: &Arc<AuthorityPerEpochStore>,
3652            _checkpoint_store: &Arc<CheckpointStore>,
3653        ) -> SuiResult {
3654            self.try_send((contents.clone(), summary.clone())).unwrap();
3655            Ok(())
3656        }
3657    }
3658
3659    #[async_trait::async_trait]
3660    impl CertifiedCheckpointOutput for mpsc::Sender<CertifiedCheckpointSummary> {
3661        async fn certified_checkpoint_created(
3662            &self,
3663            summary: &CertifiedCheckpointSummary,
3664        ) -> SuiResult {
3665            self.try_send(summary.clone()).unwrap();
3666            Ok(())
3667        }
3668    }
3669
3670    fn p(i: u64, t: Vec<u8>, timestamp_ms: u64) -> PendingCheckpoint {
3671        PendingCheckpoint {
3672            roots: vec![CheckpointRoots {
3673                tx_roots: t
3674                    .into_iter()
3675                    .map(|t| TransactionKey::Digest(d(t)))
3676                    .collect(),
3677                settlement_root: None,
3678                height: i,
3679            }],
3680            details: PendingCheckpointInfo {
3681                timestamp_ms,
3682                last_of_epoch: false,
3683                checkpoint_height: i,
3684                consensus_commit_ref: CommitRef::default(),
3685                rejected_transactions_digest: Digest::default(),
3686                checkpoint_seq: i,
3687            },
3688        }
3689    }
3690
3691    fn d(i: u8) -> TransactionDigest {
3692        let mut bytes: [u8; 32] = Default::default();
3693        bytes[0] = i;
3694        TransactionDigest::new(bytes)
3695    }
3696
3697    fn e(
3698        transaction_digest: TransactionDigest,
3699        dependencies: Vec<TransactionDigest>,
3700        gas_used: GasCostSummary,
3701    ) -> TransactionEffects {
3702        let mut effects = TransactionEffects::default();
3703        *effects.transaction_digest_mut_for_testing() = transaction_digest;
3704        *effects.dependencies_mut_for_testing() = dependencies;
3705        *effects.gas_cost_summary_mut_for_testing() = gas_used;
3706        effects
3707    }
3708
3709    fn commit_cert_for_test(
3710        store: &mut HashMap<TransactionDigest, TransactionEffects>,
3711        state: Arc<AuthorityState>,
3712        digest: TransactionDigest,
3713        dependencies: Vec<TransactionDigest>,
3714        gas_used: GasCostSummary,
3715    ) {
3716        let epoch_store = state.epoch_store_for_testing();
3717        let effects = e(digest, dependencies, gas_used);
3718        store.insert(digest, effects.clone());
3719        epoch_store.insert_executed_in_epoch(&digest);
3720    }
3721}