Skip to main content

iota_core/checkpoints/
mod.rs

1// Copyright (c) Mysten Labs, Inc.
2// Modifications Copyright (c) 2024 IOTA Stiftung
3// SPDX-License-Identifier: Apache-2.0
4
5mod causal_order;
6pub mod checkpoint_executor;
7mod checkpoint_output;
8mod epoch_info;
9mod full_checkpoint_contents_cache;
10mod metrics;
11
12use std::{
13    collections::{BTreeMap, BTreeSet, HashMap, HashSet},
14    fs::File,
15    future::Future,
16    io::Write,
17    path::Path,
18    pin::Pin,
19    sync::{Arc, Weak},
20    task::{Context, Poll},
21    time::{Duration, SystemTime},
22};
23
24use diffy::create_patch;
25use iota_common::{
26    debug_fatal, fatal,
27    random::get_rng,
28    sync::notify_read::{CHECKPOINT_BUILDER_NOTIFY_READ_TASK_NAME, NotifyRead},
29};
30use iota_metrics::{MonitoredFutureExt, monitored_future, monitored_scope};
31use iota_network::default_iota_network_config;
32use iota_sdk_types::{
33    CheckpointCommitment, CheckpointContents, CheckpointContentsDigest, CheckpointDigest,
34    CheckpointSummary, EndOfEpochData, GasCostSummary, TransactionDigest, TransactionEffects,
35    TransactionKind, UserSignature,
36};
37use iota_types::{
38    base_types::{AuthorityName, ConciseableName, EpochId, ExecutionData},
39    committee::StakeUnit,
40    crypto::AuthorityStrongQuorumSignInfo,
41    effects::{TransactionEffectsAPI, TransactionEffectsExt},
42    error::{IotaError, IotaResult},
43    event::SystemEpochInfoEvent,
44    iota_system_state::{
45        IotaSystemState, IotaSystemStateTrait,
46        epoch_start_iota_system_state::EpochStartSystemStateTrait,
47    },
48    messages_checkpoint::{
49        CertifiedCheckpointSummary, CheckpointContentsExt, CheckpointRequest, CheckpointResponse,
50        CheckpointSequenceNumber, CheckpointSignatureMessage, CheckpointSummaryExt,
51        CheckpointSummaryResponse, CheckpointTimestamp, FullCheckpointContents,
52        SignedCheckpointSummary, TrustedCheckpoint, VerifiedCheckpoint, VerifiedCheckpointContents,
53    },
54    messages_consensus::ConsensusTransactionKey,
55    storage::EpochInfoV2,
56    transaction::{TransactionAPI, TransactionEnvelope, TransactionKey},
57};
58use itertools::Itertools;
59use nonempty::NonEmpty;
60use once_cell::sync::Lazy;
61use parking_lot::Mutex;
62use pin_project_lite::pin_project;
63use rand::seq::IndexedRandom;
64use serde::{Deserialize, Serialize};
65use tokio::{
66    sync::{Notify, mpsc, watch},
67    task::JoinSet,
68    time::timeout,
69};
70use tracing::{debug, error, info, instrument, trace, warn};
71use typed_store::{
72    DBMapUtils, Map, TypedStoreError,
73    rocks::{DBMap, MetricConf},
74};
75
76pub use crate::checkpoints::{
77    checkpoint_output::{
78        LogCheckpointOutput, SendCheckpointToStateSync, SubmitCheckpointToConsensus,
79    },
80    full_checkpoint_contents_cache::{
81        FullCheckpointContentsCache, FullCheckpointContentsCacheMetrics,
82    },
83    metrics::CheckpointMetrics,
84};
85use crate::{
86    authority::{
87        AuthorityState,
88        authority_per_epoch_store::{AuthorityPerEpochStore, scorer::MAX_SCORE},
89    },
90    authority_client::{
91        make_network_authority_clients_with_network_config, validator_peer::ValidatorPeerAPI,
92    },
93    checkpoints::{
94        causal_order::CausalOrder,
95        checkpoint_output::{CertifiedCheckpointOutput, CheckpointOutput},
96    },
97    consensus_handler::SequencedConsensusTransactionKey,
98    consensus_manager::ReplayWaiter,
99    execution_cache::TransactionCacheRead,
100    global_state_hasher::GlobalStateHasher,
101    stake_aggregator::{InsertResult, MultiStakeAggregator},
102};
103
104pub type CheckpointHeight = u64;
105
106/// Digest of checkpoint contents with no transactions. Every empty checkpoint
107/// has this contents digest, so they all share one `checkpoint_content` row,
108/// which must never be deleted.
109pub(crate) static EMPTY_CHECKPOINT_CONTENTS_DIGEST: Lazy<CheckpointContentsDigest> =
110    Lazy::new(|| empty_checkpoint_contents().digest());
111
112fn empty_checkpoint_contents() -> CheckpointContents {
113    CheckpointContents::new_with_digests_and_signatures([], Vec::new())
114}
115
116pub struct EpochStats {
117    pub checkpoint_count: u64,
118    pub transaction_count: u64,
119    pub total_gas_reward: u64,
120}
121
122#[derive(Clone, Debug, Serialize, Deserialize)]
123pub struct PendingCheckpointInfo {
124    pub timestamp_ms: CheckpointTimestamp,
125    pub last_of_epoch: bool,
126    pub checkpoint_height: CheckpointHeight,
127}
128
129#[derive(Clone, Debug, Serialize, Deserialize)]
130pub enum PendingCheckpoint {
131    // This is an enum for future upgradability, though at the moment there is only one variant.
132    V1(PendingCheckpointContentsV1),
133}
134
135#[derive(Clone, Debug, Serialize, Deserialize)]
136pub struct PendingCheckpointContentsV1 {
137    pub roots: Vec<TransactionKey>,
138    pub details: PendingCheckpointInfo,
139}
140
141impl PendingCheckpoint {
142    pub fn as_v1(&self) -> &PendingCheckpointContentsV1 {
143        match self {
144            PendingCheckpoint::V1(contents) => contents,
145        }
146    }
147
148    pub fn into_v1(self) -> PendingCheckpointContentsV1 {
149        match self {
150            PendingCheckpoint::V1(contents) => contents,
151        }
152    }
153
154    pub fn roots(&self) -> &Vec<TransactionKey> {
155        &self.as_v1().roots
156    }
157
158    pub fn details(&self) -> &PendingCheckpointInfo {
159        &self.as_v1().details
160    }
161
162    pub fn height(&self) -> CheckpointHeight {
163        self.details().checkpoint_height
164    }
165}
166
167#[derive(Clone, Debug, Serialize, Deserialize)]
168pub struct BuilderCheckpointSummary {
169    pub summary: CheckpointSummary,
170    // Height at which this checkpoint summary was built. None for genesis checkpoint
171    pub checkpoint_height: Option<CheckpointHeight>,
172    pub position_in_commit: usize,
173}
174
175#[derive(DBMapUtils)]
176pub struct CheckpointStoreTables {
177    /// Maps checkpoint contents digest to checkpoint contents
178    pub(crate) checkpoint_content: DBMap<CheckpointContentsDigest, CheckpointContents>,
179
180    /// Deprecated: the contents-digest to sequence-number mapping moved to
181    /// the in-memory [`FullCheckpointContentsCache`]. Dropped on open; not
182    /// migrated (entries were a short-lived cache).
183    #[allow(dead_code)]
184    #[deprecated_db_map]
185    checkpoint_sequence_by_contents_digest: Option<DBMap<(), ()>>,
186
187    /// Deprecated: full checkpoint contents moved to the in-memory
188    /// [`FullCheckpointContentsCache`]. Dropped on open; not migrated
189    /// (entries were a short-lived cache, and readers can reconstruct full
190    /// contents from `checkpoint_content` plus the transaction and effects
191    /// stores).
192    #[allow(dead_code)]
193    #[deprecated_db_map]
194    full_checkpoint_content: Option<DBMap<(), ()>>,
195
196    /// Stores certified checkpoints
197    pub(crate) certified_checkpoints: DBMap<CheckpointSequenceNumber, TrustedCheckpoint>,
198    /// Map from checkpoint digest to certified checkpoint
199    pub(crate) checkpoint_by_digest: DBMap<CheckpointDigest, TrustedCheckpoint>,
200
201    /// Store locally computed checkpoint summaries so that we can detect forks
202    /// and log useful information. Can be pruned as soon as we verify that
203    /// we are in agreement with the latest certified checkpoint.
204    pub(crate) locally_computed_checkpoints: DBMap<CheckpointSequenceNumber, CheckpointSummary>,
205
206    /// A map from epoch ID to the sequence number of the last checkpoint in
207    /// that epoch.
208    ///
209    /// Written at certification time, before the boundary executes, so it
210    /// records the closing sequence number independently of the later-finalized
211    /// `epoch_info` row. Never pruned.
212    epoch_last_checkpoint_map: DBMap<EpochId, CheckpointSequenceNumber>,
213
214    /// Per-epoch verified data (start-of-epoch identity plus its close-of-epoch
215    /// proof bundle ([`EpochInfoV2`])) keyed by epoch ID. Populated at every
216    /// epoch boundary by the checkpoint executor and seeded from a formal
217    /// snapshot's `EPOCH_INFO` on restore.
218    ///
219    /// Intentionally not pruned: callers (the snapshot writer, the gRPC API,
220    /// etc.) need full `[0, snapshot_epoch]` coverage, so
221    /// this table grows unboundedly with epoch count (one row per epoch, ever)
222    /// by design. Do not add it to `prune_checkpoints`.
223    ///
224    /// Completeness is tracked by `epoch_info_watermark`.
225    epoch_info: DBMap<EpochId, EpochInfoV2>,
226
227    /// Highest epoch whose `epoch_info` row is finalized (close-of-epoch proof
228    /// present), as a contiguous prefix `[0, watermark]`. Singleton keyed by
229    /// `()`. Advanced atomically with the close-of-epoch upsert in
230    /// `index_epoch`, and recomputed over a seeded prefix by
231    /// `insert_epoch_info`; only ever raised. Never pruned.
232    epoch_info_watermark: DBMap<(), EpochId>,
233
234    /// Watermarks used to determine the highest verified, fully synced, and
235    /// fully executed checkpoints
236    pub(crate) watermarks: DBMap<CheckpointWatermark, (CheckpointSequenceNumber, CheckpointDigest)>,
237}
238
239impl CheckpointStoreTables {
240    pub fn new(path: &Path, metric_name: &'static str) -> Self {
241        Self::open_tables_read_write(path.to_path_buf(), MetricConf::new(metric_name), None, None)
242    }
243    pub fn open_readonly(path: &Path) -> CheckpointStoreTablesReadOnly {
244        Self::get_read_only_handle(
245            path.to_path_buf(),
246            None,
247            None,
248            MetricConf::new("checkpoint_readonly"),
249        )
250    }
251}
252
253pub struct CheckpointStore {
254    pub(crate) tables: CheckpointStoreTables,
255    full_checkpoint_contents_cache: FullCheckpointContentsCache,
256    synced_checkpoint_notify_read: NotifyRead<CheckpointSequenceNumber, VerifiedCheckpoint>,
257    executed_checkpoint_notify_read: NotifyRead<CheckpointSequenceNumber, VerifiedCheckpoint>,
258}
259
260impl CheckpointStore {
261    pub fn new(path: &Path) -> Arc<Self> {
262        Self::new_with_contents_cache(path, FullCheckpointContentsCache::default())
263    }
264
265    pub fn new_with_contents_cache(
266        path: &Path,
267        contents_cache: FullCheckpointContentsCache,
268    ) -> Arc<Self> {
269        let tables = CheckpointStoreTables::new(path, "checkpoint");
270        // Written on every open so that the shared empty contents row also
271        // exists on a database where it was deleted while empty checkpoints
272        // still used it.
273        tables
274            .checkpoint_content
275            .insert(
276                &EMPTY_CHECKPOINT_CONTENTS_DIGEST,
277                &empty_checkpoint_contents(),
278            )
279            .expect("inserting the empty checkpoint contents should succeed");
280        Arc::new(Self {
281            tables,
282            full_checkpoint_contents_cache: contents_cache,
283            synced_checkpoint_notify_read: NotifyRead::new(),
284            executed_checkpoint_notify_read: NotifyRead::new(),
285        })
286    }
287
288    pub fn new_for_tests() -> Arc<Self> {
289        let storage_dir = iota_common::tempdir().keep();
290        CheckpointStore::new(storage_dir.as_path())
291    }
292
293    pub fn open_readonly(path: &Path) -> CheckpointStoreTablesReadOnly {
294        CheckpointStoreTables::open_readonly(path)
295    }
296
297    #[instrument(level = "info", skip_all)]
298    pub fn insert_genesis_checkpoint(
299        &self,
300        checkpoint: VerifiedCheckpoint,
301        contents: CheckpointContents,
302        epoch_store: &AuthorityPerEpochStore,
303    ) {
304        assert_eq!(
305            checkpoint.epoch(),
306            0,
307            "can't call insert_genesis_checkpoint with a checkpoint not in epoch 0"
308        );
309        assert_eq!(
310            checkpoint.sequence_number(),
311            0,
312            "can't call insert_genesis_checkpoint with a checkpoint that doesn't have a sequence number of 0"
313        );
314
315        // Only insert the genesis checkpoint if the DB is empty and doesn't have it
316        // already
317        if self
318            .get_checkpoint_by_digest(checkpoint.digest())
319            .unwrap()
320            .is_none()
321        {
322            if epoch_store.epoch() == checkpoint.epoch {
323                epoch_store
324                    .put_genesis_checkpoint_in_builder(checkpoint.data(), &contents)
325                    .unwrap();
326            } else {
327                debug!(
328                    validator_epoch =% epoch_store.epoch(),
329                    genesis_epoch =% checkpoint.epoch(),
330                    "Not inserting checkpoint builder data for genesis checkpoint",
331                );
332            }
333            self.insert_checkpoint_contents(contents).unwrap();
334            self.insert_verified_checkpoint(&checkpoint).unwrap();
335            self.update_highest_synced_checkpoint(&checkpoint).unwrap();
336        }
337    }
338
339    pub fn get_checkpoint_by_digest(
340        &self,
341        digest: &CheckpointDigest,
342    ) -> Result<Option<VerifiedCheckpoint>, TypedStoreError> {
343        self.tables
344            .checkpoint_by_digest
345            .get(digest)
346            .map(|maybe_checkpoint| maybe_checkpoint.map(|c| c.into()))
347    }
348
349    pub fn get_checkpoint_by_sequence_number(
350        &self,
351        sequence_number: CheckpointSequenceNumber,
352    ) -> Result<Option<VerifiedCheckpoint>, TypedStoreError> {
353        self.tables
354            .certified_checkpoints
355            .get(&sequence_number)
356            .map(|maybe_checkpoint| maybe_checkpoint.map(|c| c.into()))
357    }
358
359    pub fn get_locally_computed_checkpoint(
360        &self,
361        sequence_number: CheckpointSequenceNumber,
362    ) -> Result<Option<CheckpointSummary>, TypedStoreError> {
363        self.tables
364            .locally_computed_checkpoints
365            .get(&sequence_number)
366    }
367
368    /// Get full checkpoint contents by contents digest from the in-memory
369    /// contents cache.
370    ///
371    /// Returns `None` once the entry has been evicted; callers reconstruct
372    /// full contents from `checkpoint_content` and the transaction and
373    /// effects stores instead.
374    pub fn get_full_checkpoint_contents_by_digest(
375        &self,
376        digest: &CheckpointContentsDigest,
377    ) -> Option<Arc<FullCheckpointContents>> {
378        self.full_checkpoint_contents_cache.get_by_digest(digest)
379    }
380
381    pub fn get_latest_certified_checkpoint(
382        &self,
383    ) -> Result<Option<VerifiedCheckpoint>, TypedStoreError> {
384        Ok(self
385            .tables
386            .certified_checkpoints
387            .safe_range_iter_reversed(..)
388            .next()
389            .transpose()?
390            .map(|(_, v)| v.into()))
391    }
392
393    pub fn get_latest_locally_computed_checkpoint(
394        &self,
395    ) -> Result<Option<CheckpointSummary>, TypedStoreError> {
396        Ok(self
397            .tables
398            .locally_computed_checkpoints
399            .safe_range_iter_reversed(..)
400            .next()
401            .transpose()?
402            .map(|(_, v)| v))
403    }
404
405    pub fn multi_get_checkpoint_by_sequence_number(
406        &self,
407        sequence_numbers: &[CheckpointSequenceNumber],
408    ) -> Result<Vec<Option<VerifiedCheckpoint>>, TypedStoreError> {
409        let checkpoints = self
410            .tables
411            .certified_checkpoints
412            .multi_get(sequence_numbers)?
413            .into_iter()
414            .map(|maybe_checkpoint| maybe_checkpoint.map(|c| c.into()))
415            .collect();
416
417        Ok(checkpoints)
418    }
419
420    pub fn multi_get_checkpoint_content(
421        &self,
422        contents_digest: &[CheckpointContentsDigest],
423    ) -> Result<Vec<Option<CheckpointContents>>, TypedStoreError> {
424        self.tables.checkpoint_content.multi_get(contents_digest)
425    }
426
427    /// The row `watermark` holds — a sequence number and a digest — and `None`
428    /// when it has never been written.
429    ///
430    /// Only the verified, synced and executed rows name a checkpoint. The
431    /// pruner writes a random digest for `HighestPruned`, since nothing reads
432    /// it, so that row's sequence number is the only part worth having.
433    fn get_watermark(
434        &self,
435        watermark: CheckpointWatermark,
436    ) -> Result<Option<(CheckpointSequenceNumber, CheckpointDigest)>, TypedStoreError> {
437        self.tables.watermarks.get(&watermark)
438    }
439
440    /// The checkpoint `watermark` names.
441    ///
442    /// Resolved by digest, so it costs a lookup the row itself does not, and
443    /// answers `None` for a checkpoint the store no longer holds. A caller
444    /// that only compares positions wants [`Self::get_watermark_seq_number`].
445    fn get_watermark_checkpoint(
446        &self,
447        watermark: CheckpointWatermark,
448    ) -> Result<Option<VerifiedCheckpoint>, TypedStoreError> {
449        let Some((_sequence_number, digest)) = self.get_watermark(watermark)? else {
450            return Ok(None);
451        };
452        self.get_checkpoint_by_digest(&digest)
453    }
454
455    /// The sequence number of the checkpoint `watermark` names, read from the
456    /// row rather than from the checkpoint.
457    fn get_watermark_seq_number(
458        &self,
459        watermark: CheckpointWatermark,
460    ) -> Result<Option<CheckpointSequenceNumber>, TypedStoreError> {
461        Ok(self
462            .get_watermark(watermark)?
463            .map(|(sequence_number, _digest)| sequence_number))
464    }
465
466    pub fn get_highest_verified_checkpoint(
467        &self,
468    ) -> Result<Option<VerifiedCheckpoint>, TypedStoreError> {
469        self.get_watermark_checkpoint(CheckpointWatermark::HighestVerified)
470    }
471
472    pub fn get_highest_verified_checkpoint_seq_number(
473        &self,
474    ) -> Result<Option<CheckpointSequenceNumber>, TypedStoreError> {
475        self.get_watermark_seq_number(CheckpointWatermark::HighestVerified)
476    }
477
478    pub fn get_highest_synced_checkpoint(
479        &self,
480    ) -> Result<Option<VerifiedCheckpoint>, TypedStoreError> {
481        self.get_watermark_checkpoint(CheckpointWatermark::HighestSynced)
482    }
483
484    pub fn get_highest_synced_checkpoint_seq_number(
485        &self,
486    ) -> Result<Option<CheckpointSequenceNumber>, TypedStoreError> {
487        self.get_watermark_seq_number(CheckpointWatermark::HighestSynced)
488    }
489
490    pub fn get_highest_executed_checkpoint(
491        &self,
492    ) -> Result<Option<VerifiedCheckpoint>, TypedStoreError> {
493        self.get_watermark_checkpoint(CheckpointWatermark::HighestExecuted)
494    }
495
496    pub fn get_highest_executed_checkpoint_seq_number(
497        &self,
498    ) -> Result<Option<CheckpointSequenceNumber>, TypedStoreError> {
499        self.get_watermark_seq_number(CheckpointWatermark::HighestExecuted)
500    }
501
502    pub fn get_highest_pruned_checkpoint_seq_number(
503        &self,
504    ) -> Result<Option<CheckpointSequenceNumber>, TypedStoreError> {
505        self.get_watermark_seq_number(CheckpointWatermark::HighestPruned)
506    }
507
508    pub fn get_checkpoint_contents(
509        &self,
510        digest: &CheckpointContentsDigest,
511    ) -> Result<Option<CheckpointContents>, TypedStoreError> {
512        self.tables.checkpoint_content.get(digest)
513    }
514
515    /// Get full checkpoint contents from the in-memory contents cache.
516    ///
517    /// Returns `None` once the entry has been evicted; callers fall back to
518    /// loading the individual transactions and effects from their stores.
519    pub fn get_full_checkpoint_contents_by_sequence_number(
520        &self,
521        seq: CheckpointSequenceNumber,
522    ) -> Option<Arc<FullCheckpointContents>> {
523        self.full_checkpoint_contents_cache.get_by_seq(seq)
524    }
525
526    fn prune_local_summaries(&self) -> IotaResult {
527        if let Some((last_local_summary, _)) = self
528            .tables
529            .locally_computed_checkpoints
530            .safe_range_iter_reversed(..)
531            .next()
532            .transpose()?
533        {
534            let mut batch = self.tables.locally_computed_checkpoints.batch();
535            batch.schedule_delete_range(
536                &self.tables.locally_computed_checkpoints,
537                &0,
538                &last_local_summary,
539            )?;
540            batch.write()?;
541            info!("Pruned local summaries up to {:?}", last_local_summary);
542        }
543        Ok(())
544    }
545
546    #[instrument(level = "trace", skip_all)]
547    fn check_for_checkpoint_fork(
548        &self,
549        local_checkpoint: &CheckpointSummary,
550        verified_checkpoint: &VerifiedCheckpoint,
551    ) {
552        if local_checkpoint != verified_checkpoint.data() {
553            let verified_contents = self
554                .get_checkpoint_contents(&verified_checkpoint.contents_digest)
555                .map(|opt_contents| {
556                    opt_contents
557                        .map(|contents| format!("{contents:?}"))
558                        .unwrap_or_else(|| {
559                            format!(
560                                "Verified checkpoint contents not found, digest: {:?}",
561                                verified_checkpoint.contents_digest,
562                            )
563                        })
564                })
565                .map_err(|e| {
566                    format!(
567                        "Failed to get verified checkpoint contents, digest: {:?} error: {:?}",
568                        verified_checkpoint.contents_digest, e
569                    )
570                })
571                .unwrap_or_else(|err_msg| err_msg);
572
573            let local_contents = self
574                .get_checkpoint_contents(&local_checkpoint.contents_digest)
575                .map(|opt_contents| {
576                    opt_contents
577                        .map(|contents| format!("{contents:?}"))
578                        .unwrap_or_else(|| {
579                            format!(
580                                "Local checkpoint contents not found, digest: {:?}",
581                                local_checkpoint.contents_digest
582                            )
583                        })
584                })
585                .map_err(|e| {
586                    format!(
587                        "Failed to get local checkpoint contents, digest: {:?} error: {:?}",
588                        local_checkpoint.contents_digest, e
589                    )
590                })
591                .unwrap_or_else(|err_msg| err_msg);
592
593            // checkpoint contents may be too large for panic message.
594            error!(
595                verified_checkpoint = ?verified_checkpoint.data(),
596                ?verified_contents,
597                ?local_checkpoint,
598                ?local_contents,
599                "Local checkpoint fork detected!",
600            );
601            fatal!(
602                "Local checkpoint fork detected for sequence number: {}",
603                local_checkpoint.sequence_number()
604            );
605        }
606    }
607
608    // Called by consensus (ConsensusAggregator).
609    // Different from `insert_verified_checkpoint`, it does not touch
610    // the highest_verified_checkpoint watermark such that state sync
611    // will have a chance to process this checkpoint and perform some
612    // state-sync only things.
613    pub fn insert_certified_checkpoint(
614        &self,
615        checkpoint: &VerifiedCheckpoint,
616    ) -> Result<(), TypedStoreError> {
617        self.multi_insert_certified_checkpoints(std::slice::from_ref(checkpoint))
618    }
619
620    /// Inserts a batch of certified checkpoints in a single write, like
621    /// [`Self::insert_certified_checkpoint`] for each of them.
622    pub fn multi_insert_certified_checkpoints(
623        &self,
624        checkpoints: &[VerifiedCheckpoint],
625    ) -> Result<(), TypedStoreError> {
626        let Some(last) = checkpoints.last() else {
627            return Ok(());
628        };
629        debug!(
630            checkpoint_seq = last.sequence_number(),
631            count = checkpoints.len(),
632            "Inserting certified checkpoints",
633        );
634        let mut batch = self.tables.certified_checkpoints.batch();
635        batch
636            .insert_batch(
637                &self.tables.certified_checkpoints,
638                checkpoints
639                    .iter()
640                    .map(|c| (c.sequence_number(), c.serializable_ref())),
641            )?
642            .insert_batch(
643                &self.tables.checkpoint_by_digest,
644                checkpoints
645                    .iter()
646                    .map(|c| (c.digest(), c.serializable_ref())),
647            )?
648            .insert_batch(
649                &self.tables.epoch_last_checkpoint_map,
650                checkpoints
651                    .iter()
652                    .filter(|c| c.next_epoch_committee().is_some())
653                    .map(|c| (c.epoch(), c.sequence_number())),
654            )?;
655        batch.write()?;
656
657        for checkpoint in checkpoints {
658            if let Some(local_checkpoint) = self
659                .tables
660                .locally_computed_checkpoints
661                .get(&checkpoint.sequence_number())?
662            {
663                self.check_for_checkpoint_fork(&local_checkpoint, checkpoint);
664            }
665        }
666
667        Ok(())
668    }
669
670    // Called by state sync, apart from inserting the checkpoint and updating
671    // related tables, it also bumps the highest_verified_checkpoint watermark.
672    #[instrument(level = "trace", skip_all)]
673    pub fn insert_verified_checkpoint(
674        &self,
675        checkpoint: &VerifiedCheckpoint,
676    ) -> Result<(), TypedStoreError> {
677        self.insert_certified_checkpoint(checkpoint)?;
678        self.update_highest_verified_checkpoint(checkpoint)
679    }
680
681    pub fn update_highest_verified_checkpoint(
682        &self,
683        checkpoint: &VerifiedCheckpoint,
684    ) -> Result<(), TypedStoreError> {
685        if Some(checkpoint.sequence_number()) > self.get_highest_verified_checkpoint_seq_number()? {
686            debug!(
687                checkpoint_seq = checkpoint.sequence_number(),
688                "Updating highest verified checkpoint",
689            );
690            self.tables.watermarks.insert(
691                &CheckpointWatermark::HighestVerified,
692                &(checkpoint.sequence_number(), *checkpoint.digest()),
693            )?;
694        }
695
696        Ok(())
697    }
698
699    pub fn update_highest_synced_checkpoint(
700        &self,
701        checkpoint: &VerifiedCheckpoint,
702    ) -> Result<(), TypedStoreError> {
703        self.multi_update_highest_synced_checkpoint(std::slice::from_ref(checkpoint))
704    }
705
706    /// Marks a consecutive run of checkpoints as synced: writes the watermark
707    /// once, for the last checkpoint, and notifies the waiters that write
708    /// releases. Does nothing when another writer has already carried the
709    /// watermark past the run.
710    pub fn multi_update_highest_synced_checkpoint(
711        &self,
712        checkpoints: &[VerifiedCheckpoint],
713    ) -> Result<(), TypedStoreError> {
714        let Some(last) = checkpoints.last() else {
715            return Ok(());
716        };
717
718        // Only ever forwards, like the other two: another writer advances this
719        // as well, and moving it back offers contents that may not be there.
720        let previous = self.get_highest_synced_checkpoint_seq_number()?;
721        if previous.is_none_or(|previous_seq_number| previous_seq_number < last.sequence_number()) {
722            debug!(
723                checkpoint_seq = last.sequence_number(),
724                "Updating highest synced checkpoint",
725            );
726            self.tables.watermarks.insert(
727                &CheckpointWatermark::HighestSynced,
728                &(last.sequence_number(), *last.digest()),
729            )?;
730
731            // A reader parks only below the watermark, so only what this
732            // write has just passed can have anyone waiting on it.
733            for checkpoint in checkpoints.iter().filter(|c| {
734                previous.is_none_or(|previous_seq_number| previous_seq_number < c.sequence_number())
735            }) {
736                self.synced_checkpoint_notify_read
737                    .notify(&checkpoint.sequence_number(), checkpoint);
738            }
739        }
740        Ok(())
741    }
742
743    async fn notify_read_checkpoint_watermark<F>(
744        &self,
745        notify_read: &NotifyRead<CheckpointSequenceNumber, VerifiedCheckpoint>,
746        seq: CheckpointSequenceNumber,
747        get_watermark: F,
748    ) -> VerifiedCheckpoint
749    where
750        F: Fn() -> Option<CheckpointSequenceNumber>,
751    {
752        type ReadResult = Result<Vec<Option<VerifiedCheckpoint>>, TypedStoreError>;
753
754        notify_read
755            .read("notify_read_checkpoint_watermark", &[seq], |seqs| {
756                let seq = seqs[0];
757                let Some(highest) = get_watermark() else {
758                    return Ok(vec![None]) as ReadResult;
759                };
760                if highest < seq {
761                    return Ok(vec![None]) as ReadResult;
762                }
763                let checkpoint = self
764                    .get_checkpoint_by_sequence_number(seq)
765                    .expect("db error")
766                    .expect("checkpoint not found");
767                Ok(vec![Some(checkpoint)]) as ReadResult
768            })
769            .await
770            .unwrap()
771            .into_iter()
772            .next()
773            .unwrap()
774    }
775
776    pub async fn notify_read_synced_checkpoint(
777        &self,
778        seq: CheckpointSequenceNumber,
779    ) -> VerifiedCheckpoint {
780        self.notify_read_checkpoint_watermark(&self.synced_checkpoint_notify_read, seq, || {
781            self.get_highest_synced_checkpoint_seq_number()
782                .expect("db error")
783        })
784        .await
785    }
786
787    pub async fn notify_read_executed_checkpoint(
788        &self,
789        seq: CheckpointSequenceNumber,
790    ) -> VerifiedCheckpoint {
791        self.notify_read_checkpoint_watermark(&self.executed_checkpoint_notify_read, seq, || {
792            self.get_highest_executed_checkpoint_seq_number()
793                .expect("db error")
794        })
795        .await
796    }
797
798    pub fn update_highest_executed_checkpoint(
799        &self,
800        checkpoint: &VerifiedCheckpoint,
801    ) -> Result<(), TypedStoreError> {
802        if let Some(seq_number) = self.get_highest_executed_checkpoint_seq_number()? {
803            if seq_number >= checkpoint.sequence_number() {
804                return Ok(());
805            }
806            assert_eq!(
807                seq_number + 1,
808                checkpoint.sequence_number(),
809                "Cannot update highest executed checkpoint to {} when current highest executed checkpoint is {}",
810                checkpoint.sequence_number(),
811                seq_number
812            );
813        }
814        let seq = checkpoint.sequence_number();
815        debug!(checkpoint_seq = seq, "Updating highest executed checkpoint",);
816        self.tables.watermarks.insert(
817            &CheckpointWatermark::HighestExecuted,
818            &(seq, *checkpoint.digest()),
819        )?;
820        self.executed_checkpoint_notify_read
821            .notify(&seq, checkpoint);
822        Ok(())
823    }
824
825    pub fn update_highest_pruned_checkpoint(
826        &self,
827        checkpoint: &VerifiedCheckpoint,
828    ) -> Result<(), TypedStoreError> {
829        self.tables.watermarks.insert(
830            &CheckpointWatermark::HighestPruned,
831            &(checkpoint.sequence_number(), *checkpoint.digest()),
832        )
833    }
834
835    /// Sets the verified watermark to `checkpoint` whether or not that moves
836    /// it forwards.
837    ///
838    /// Only for tooling that deliberately rewinds it. The node writes this row
839    /// through [`Self::update_highest_verified_checkpoint`], which declines to
840    /// move it backwards, under a lock its callers share.
841    pub fn set_highest_verified_checkpoint_subtle(
842        &self,
843        checkpoint: &VerifiedCheckpoint,
844    ) -> Result<(), TypedStoreError> {
845        self.tables.watermarks.insert(
846            &CheckpointWatermark::HighestVerified,
847            &(checkpoint.sequence_number(), *checkpoint.digest()),
848        )
849    }
850
851    /// Sets the synced watermark to `checkpoint` whether or not that moves it
852    /// forwards.
853    ///
854    /// Only for tooling that deliberately rewinds it. The node writes this row
855    /// through [`Self::multi_update_highest_synced_checkpoint`], which declines
856    /// to move it backwards, under a lock its callers share.
857    pub fn set_highest_synced_checkpoint_subtle(
858        &self,
859        checkpoint: &VerifiedCheckpoint,
860    ) -> Result<(), TypedStoreError> {
861        self.tables.watermarks.insert(
862            &CheckpointWatermark::HighestSynced,
863            &(checkpoint.sequence_number(), *checkpoint.digest()),
864        )
865    }
866
867    /// Sets highest executed checkpoint to any value.
868    ///
869    /// WARNING: This method is very subtle and can corrupt the database if used
870    /// incorrectly. It should only be used in one-off cases or tests after
871    /// fully understanding the risk.
872    pub fn set_highest_executed_checkpoint_subtle(
873        &self,
874        checkpoint: &VerifiedCheckpoint,
875    ) -> Result<(), TypedStoreError> {
876        self.tables.watermarks.insert(
877            &CheckpointWatermark::HighestExecuted,
878            &(checkpoint.sequence_number(), *checkpoint.digest()),
879        )
880    }
881
882    pub fn insert_checkpoint_contents(
883        &self,
884        contents: CheckpointContents,
885    ) -> Result<(), TypedStoreError> {
886        debug!(
887            checkpoint_seq = ?contents.digest(),
888            "Inserting checkpoint contents",
889        );
890        self.tables
891            .checkpoint_content
892            .insert(&contents.digest(), &contents)
893    }
894
895    /// Persists the checkpoint contents in digest form and caches the full
896    /// contents in memory, where they serve the checkpoint executor's bulk
897    /// transaction loads and contents requests from state-sync peers.
898    ///
899    /// INVARIANT: See [`Self::cache_full_checkpoint_contents`].
900    pub fn insert_verified_checkpoint_contents(
901        &self,
902        checkpoint: &VerifiedCheckpoint,
903        full_contents: VerifiedCheckpointContents,
904    ) -> Result<(), TypedStoreError> {
905        self.multi_insert_verified_checkpoint_contents(vec![(checkpoint.clone(), full_contents)])
906    }
907
908    /// Batch variant of [`Self::insert_verified_checkpoint_contents`]:
909    /// persists all contents rows in a single write, then caches the full
910    /// contents.
911    ///
912    /// INVARIANT: See [`Self::cache_full_checkpoint_contents`].
913    pub fn multi_insert_verified_checkpoint_contents(
914        &self,
915        checkpoints: Vec<(VerifiedCheckpoint, VerifiedCheckpointContents)>,
916    ) -> Result<(), TypedStoreError> {
917        let checkpoints: Vec<_> = checkpoints
918            .into_iter()
919            .map(|(checkpoint, full_contents)| (checkpoint, full_contents.into_inner()))
920            .collect();
921
922        let mut batch = self.tables.checkpoint_content.batch();
923        for (checkpoint, full_contents) in &checkpoints {
924            let contents = full_contents.checkpoint_contents();
925            assert_eq!(checkpoint.contents_digest, contents.digest());
926            batch.insert_batch(
927                &self.tables.checkpoint_content,
928                [(contents.digest(), contents)],
929            )?;
930        }
931        batch.write()?;
932
933        for (checkpoint, full_contents) in checkpoints {
934            self.cache_full_checkpoint_contents(
935                checkpoint.sequence_number(),
936                checkpoint.contents_digest,
937                full_contents,
938            );
939        }
940        Ok(())
941    }
942
943    /// Whether [`Self::cache_full_checkpoint_contents`] would retain contents
944    /// for this sequence number, so callers can skip assembling contents that
945    /// the cache would reject (disabled cache, or an entry the lowest-seq
946    /// eviction would remove immediately).
947    pub fn should_cache_full_checkpoint_contents(&self, seq: CheckpointSequenceNumber) -> bool {
948        self.full_checkpoint_contents_cache.should_cache(seq)
949    }
950
951    /// Caches full checkpoint contents in memory without writing anything to
952    /// disk, so state-sync peers can be served without reconstructing the
953    /// contents. `contents_digest` must be the digest of `full_contents`.
954    ///
955    /// INVARIANT: The caller must have durably written the matching
956    /// `checkpoint_content` row (and the contained transactions and
957    /// effects) first for two reasons:
958    /// 1. once the cache evicts the entry (or after a restart), readers
959    ///    reconstruct the full contents from those stores.
960    /// 2. state-sync treats available contents as proof of that row and skips
961    ///    its own durable write, and the checkpoint executor panics on a
962    ///    missing row.
963    ///
964    /// Best-effort: a serialization failure is logged and the insert skipped;
965    /// readers fall back to reconstructing the contents from the durable
966    /// stores.
967    pub fn cache_full_checkpoint_contents(
968        &self,
969        sequence_number: CheckpointSequenceNumber,
970        contents_digest: CheckpointContentsDigest,
971        full_contents: FullCheckpointContents,
972    ) {
973        let size = match bcs::serialized_size(&full_contents) {
974            Ok(size) => size,
975            Err(e) => {
976                warn!(
977                    sequence_number,
978                    "failed to serialize full checkpoint contents for caching: {e}"
979                );
980                return;
981            }
982        };
983        self.full_checkpoint_contents_cache.insert(
984            sequence_number,
985            contents_digest,
986            Arc::new(full_contents),
987            size,
988        );
989    }
990
991    pub fn get_epoch_last_checkpoint(
992        &self,
993        epoch_id: EpochId,
994    ) -> IotaResult<Option<VerifiedCheckpoint>> {
995        let seq = self.get_epoch_last_checkpoint_seq_number(epoch_id)?;
996        let checkpoint = match seq {
997            Some(seq) => self.get_checkpoint_by_sequence_number(seq)?,
998            None => None,
999        };
1000        Ok(checkpoint)
1001    }
1002
1003    /// Sequence number of `epoch_id`'s last checkpoint. Unlike
1004    /// [`Self::get_epoch_last_checkpoint`], this does not require the summary
1005    /// itself to still be present: the underlying map is never pruned.
1006    pub fn get_epoch_last_checkpoint_seq_number(
1007        &self,
1008        epoch_id: EpochId,
1009    ) -> Result<Option<CheckpointSequenceNumber>, TypedStoreError> {
1010        self.tables.epoch_last_checkpoint_map.get(&epoch_id)
1011    }
1012
1013    pub fn insert_epoch_last_checkpoint(
1014        &self,
1015        epoch_id: EpochId,
1016        checkpoint: &VerifiedCheckpoint,
1017    ) -> IotaResult {
1018        self.tables
1019            .epoch_last_checkpoint_map
1020            .insert(&epoch_id, &checkpoint.sequence_number())?;
1021        Ok(())
1022    }
1023
1024    pub fn get_epoch_state_commitments(
1025        &self,
1026        epoch: EpochId,
1027    ) -> IotaResult<Option<Vec<CheckpointCommitment>>> {
1028        let commitments = self.get_epoch_last_checkpoint(epoch)?.map(|checkpoint| {
1029            checkpoint
1030                .end_of_epoch_data
1031                .as_ref()
1032                .expect("Last checkpoint of epoch expected to have EndOfEpochData")
1033                .epoch_commitments
1034                .clone()
1035        });
1036        Ok(commitments)
1037    }
1038
1039    /// Given the epoch ID, and the last checkpoint of the epoch, derive a few
1040    /// statistics of the epoch.
1041    pub fn get_epoch_stats(
1042        &self,
1043        epoch: EpochId,
1044        last_checkpoint: &CheckpointSummary,
1045    ) -> Option<EpochStats> {
1046        let (first_checkpoint, prev_epoch_network_transactions) = if epoch == 0 {
1047            (0, 0)
1048        } else if let Ok(Some(checkpoint)) = self.get_epoch_last_checkpoint(epoch - 1) {
1049            (
1050                checkpoint.sequence_number + 1,
1051                checkpoint.network_total_transactions,
1052            )
1053        } else {
1054            return None;
1055        };
1056        Some(EpochStats {
1057            checkpoint_count: last_checkpoint.sequence_number - first_checkpoint + 1,
1058            transaction_count: last_checkpoint.network_total_transactions
1059                - prev_epoch_network_transactions,
1060            total_gas_reward: last_checkpoint
1061                .epoch_rolling_gas_cost_summary
1062                .computation_cost,
1063        })
1064    }
1065}
1066
1067#[derive(Copy, Clone, Debug, Serialize, Deserialize)]
1068pub enum CheckpointWatermark {
1069    HighestVerified,
1070    HighestSynced,
1071    HighestExecuted,
1072    HighestPruned,
1073}
1074
1075struct CheckpointStateHasher {
1076    epoch_store: Arc<AuthorityPerEpochStore>,
1077    hasher: Weak<GlobalStateHasher>,
1078    receive_from_builder: mpsc::Receiver<(CheckpointSequenceNumber, Vec<TransactionEffects>)>,
1079}
1080
1081impl CheckpointStateHasher {
1082    fn new(
1083        epoch_store: Arc<AuthorityPerEpochStore>,
1084        hasher: Weak<GlobalStateHasher>,
1085        receive_from_builder: mpsc::Receiver<(CheckpointSequenceNumber, Vec<TransactionEffects>)>,
1086    ) -> Self {
1087        Self {
1088            epoch_store,
1089            hasher,
1090            receive_from_builder,
1091        }
1092    }
1093
1094    async fn run(self) {
1095        let Self {
1096            epoch_store,
1097            hasher,
1098            mut receive_from_builder,
1099        } = self;
1100        while let Some((seq, effects)) = receive_from_builder.recv().await {
1101            let Some(hasher) = hasher.upgrade() else {
1102                info!("GlobalStateHash was dropped, stopping checkpoint accumulation");
1103                break;
1104            };
1105            hasher
1106                .accumulate_checkpoint(&effects, seq, &epoch_store)
1107                .expect("epoch ended while accumulating checkpoint");
1108        }
1109    }
1110}
1111
1112/// A checkpoint produced by [`CheckpointBuilder::create_checkpoints`],
1113/// written and cached by [`CheckpointBuilder::write_checkpoints`].
1114struct BuiltCheckpoint {
1115    summary: CheckpointSummary,
1116    contents: CheckpointContents,
1117    /// Full contents for the in-memory cache, assembled only when the cache
1118    /// would retain them.
1119    full_contents: Option<FullCheckpointContents>,
1120}
1121
1122#[derive(Debug)]
1123pub enum CheckpointBuilderError {
1124    ChangeEpochTxAlreadyExecuted,
1125    SystemPackagesMissing,
1126    Retry(anyhow::Error),
1127}
1128
1129impl<IotaError: std::error::Error + Send + Sync + 'static> From<IotaError>
1130    for CheckpointBuilderError
1131{
1132    fn from(e: IotaError) -> Self {
1133        Self::Retry(e.into())
1134    }
1135}
1136
1137pub type CheckpointBuilderResult<T = ()> = Result<T, CheckpointBuilderError>;
1138
1139pub struct CheckpointBuilder {
1140    state: Arc<AuthorityState>,
1141    store: Arc<CheckpointStore>,
1142    epoch_store: Arc<AuthorityPerEpochStore>,
1143    notify: Arc<Notify>,
1144    notify_aggregator: Arc<Notify>,
1145    last_built: watch::Sender<CheckpointSequenceNumber>,
1146    effects_store: Arc<dyn TransactionCacheRead>,
1147    global_state_hasher: Weak<GlobalStateHasher>,
1148    send_to_hasher: mpsc::Sender<(CheckpointSequenceNumber, Vec<TransactionEffects>)>,
1149    output: Box<dyn CheckpointOutput>,
1150    metrics: Arc<CheckpointMetrics>,
1151    max_transactions_per_checkpoint: usize,
1152    max_checkpoint_size_bytes: usize,
1153}
1154
1155pub struct CheckpointAggregator {
1156    store: Arc<CheckpointStore>,
1157    epoch_store: Arc<AuthorityPerEpochStore>,
1158    notify: Arc<Notify>,
1159    current: Option<CheckpointSignatureAggregator>,
1160    output: Box<dyn CertifiedCheckpointOutput>,
1161    state: Arc<AuthorityState>,
1162    metrics: Arc<CheckpointMetrics>,
1163}
1164
1165// This holds information to aggregate signatures for one checkpoint
1166pub struct CheckpointSignatureAggregator {
1167    next_index: u64,
1168    summary: CheckpointSummary,
1169    digest: CheckpointDigest,
1170    /// Aggregates voting stake for each signed checkpoint proposal by authority
1171    signatures_by_digest: MultiStakeAggregator<CheckpointDigest, CheckpointSummary, true>,
1172    store: Arc<CheckpointStore>,
1173    state: Arc<AuthorityState>,
1174    metrics: Arc<CheckpointMetrics>,
1175}
1176
1177impl CheckpointBuilder {
1178    fn new(
1179        state: Arc<AuthorityState>,
1180        store: Arc<CheckpointStore>,
1181        epoch_store: Arc<AuthorityPerEpochStore>,
1182        notify: Arc<Notify>,
1183        effects_store: Arc<dyn TransactionCacheRead>,
1184        // for synchronous accumulation of end-of-epoch checkpoint
1185        global_state_hasher: Weak<GlobalStateHasher>,
1186        // for asynchronous/concurrent accumulation of regular checkpoints
1187        send_to_hasher: mpsc::Sender<(CheckpointSequenceNumber, Vec<TransactionEffects>)>,
1188        output: Box<dyn CheckpointOutput>,
1189        notify_aggregator: Arc<Notify>,
1190        last_built: watch::Sender<CheckpointSequenceNumber>,
1191        metrics: Arc<CheckpointMetrics>,
1192        max_transactions_per_checkpoint: usize,
1193        max_checkpoint_size_bytes: usize,
1194    ) -> Self {
1195        Self {
1196            state,
1197            store,
1198            epoch_store,
1199            notify,
1200            effects_store,
1201            global_state_hasher,
1202            send_to_hasher,
1203            output,
1204            notify_aggregator,
1205            last_built,
1206            metrics,
1207            max_transactions_per_checkpoint,
1208            max_checkpoint_size_bytes,
1209        }
1210    }
1211
1212    /// This function first waits for ConsensusCommitHandler to finish
1213    /// reprocessing commits that have been processed before the last
1214    /// restart, if consensus_replay_waiter is supplied. Then it starts
1215    /// building checkpoints in a loop.
1216    ///
1217    /// It is optional to pass in consensus_replay_waiter, to make it easier to
1218    /// attribute if slow recovery of previously built checkpoints is due to
1219    /// consensus replay or checkpoint building.
1220    async fn run(mut self, consensus_replay_waiter: Option<ReplayWaiter>) {
1221        if let Some(replay_waiter) = consensus_replay_waiter {
1222            info!("Waiting for consensus commits to replay ...");
1223            replay_waiter.wait_for_replay().await;
1224            info!("Consensus commits finished replaying");
1225        }
1226        info!("Starting CheckpointBuilder");
1227        loop {
1228            match self.maybe_build_checkpoints().await {
1229                Ok(()) => {}
1230                err @ Err(
1231                    CheckpointBuilderError::ChangeEpochTxAlreadyExecuted
1232                    | CheckpointBuilderError::SystemPackagesMissing,
1233                ) => {
1234                    info!("CheckpointBuilder stopping: {:?}", err);
1235                    return;
1236                }
1237                Err(CheckpointBuilderError::Retry(inner)) => {
1238                    let msg = format!("{inner:?}");
1239                    debug_fatal!("Error while making checkpoint, will retry in 1s: {}", msg);
1240                    tokio::time::sleep(Duration::from_secs(1)).await;
1241                    self.metrics.checkpoint_errors.inc();
1242                    continue;
1243                }
1244            }
1245
1246            self.notify.notified().await;
1247        }
1248    }
1249
1250    async fn maybe_build_checkpoints(&mut self) -> CheckpointBuilderResult {
1251        let _scope = monitored_scope("BuildCheckpoints");
1252
1253        // Collect info about the most recently built checkpoint.
1254        let summary = self
1255            .epoch_store
1256            .last_built_checkpoint_builder_summary()
1257            .expect("epoch should not have ended");
1258        let mut last_height = summary.as_ref().and_then(|s| s.checkpoint_height);
1259        let mut last_timestamp = summary.as_ref().map(|s| s.summary.timestamp_ms);
1260        let mut last_seq = summary.map(|s| s.summary.sequence_number);
1261
1262        let min_checkpoint_interval_ms = self
1263            .epoch_store
1264            .protocol_config()
1265            .min_checkpoint_interval_ms_as_option()
1266            .unwrap_or_default();
1267        // When set, the interval may additionally be amortized over this many
1268        // recent checkpoints of the current epoch.
1269        let checkpoint_rate_window_size = self
1270            .epoch_store
1271            .protocol_config()
1272            .checkpoint_rate_window_size_as_option();
1273        let mut grouped_pending_checkpoints = Vec::new();
1274        let mut checkpoints_iter = self
1275            .epoch_store
1276            .get_pending_checkpoints(last_height)
1277            .expect("unexpected epoch store error")
1278            .into_iter()
1279            .peekable();
1280        while let Some((height, pending)) = checkpoints_iter.next() {
1281            let current_timestamp = pending.details().timestamp_ms;
1282            // Strict interval against the immediately preceding checkpoint.
1283            let adjacent_interval_elapsed = match last_timestamp {
1284                Some(last_timestamp) => {
1285                    current_timestamp >= last_timestamp + min_checkpoint_interval_ms
1286                }
1287                None => true,
1288            };
1289            // Windowed arm: also allow building once the checkpoint `window`
1290            // back is at least `window * interval` older. This recycles the
1291            // slack the strict arm loses to discrete commit timestamps,
1292            // holding the sustained rate at the ceiling, while the strict arm
1293            // keeps every quiet gap — and thus checkpoint sizes — within one
1294            // interval. Only checkpoints built in the current epoch are
1295            // consulted: every validator builds all of them, so the look-back
1296            // resolves identically everywhere and the gate stays
1297            // deterministic. Until the window fills (epoch start, chain
1298            // genesis) the arm is inert.
1299            let interval_elapsed = adjacent_interval_elapsed
1300                || checkpoint_rate_window_size.is_some_and(|window| {
1301                    last_seq
1302                        .and_then(|seq| (seq + 1).checked_sub(window))
1303                        .and_then(|window_start_seq| {
1304                            self.epoch_store
1305                                .get_built_checkpoint_summary(window_start_seq)
1306                                .expect("epoch store should not error reading a built checkpoint summary")
1307                        })
1308                        .is_some_and(|window_start| {
1309                            current_timestamp
1310                                >= window_start.timestamp_ms + window * min_checkpoint_interval_ms
1311                        })
1312                });
1313            // Group PendingCheckpoints until:
1314            // - the minimum interval has elapsed ...
1315            let can_build = interval_elapsed
1316                // - or, next PendingCheckpoint is last-of-epoch (since the last-of-epoch checkpoint
1317                //   should be written separately) ...
1318                || checkpoints_iter
1319                .peek()
1320                .is_some_and(|(_, next_pending)| next_pending.details().last_of_epoch)
1321                // - or, we have reached end of epoch.
1322                || pending.details().last_of_epoch;
1323            grouped_pending_checkpoints.push(pending);
1324            if !can_build {
1325                debug!(
1326                    checkpoint_commit_height = height,
1327                    ?last_timestamp,
1328                    ?current_timestamp,
1329                    "waiting for more PendingCheckpoints: minimum interval not yet elapsed"
1330                );
1331                continue;
1332            }
1333
1334            // Min interval has elapsed, we can now coalesce and build a checkpoint.
1335            last_height = Some(height);
1336            last_timestamp = Some(current_timestamp);
1337            let commits_in_checkpoint = grouped_pending_checkpoints.len();
1338            debug!(
1339                checkpoint_commit_height_from = grouped_pending_checkpoints
1340                    .first()
1341                    .unwrap()
1342                    .details()
1343                    .checkpoint_height,
1344                checkpoint_commit_height_to = last_height,
1345                "Making checkpoint with commit height range"
1346            );
1347
1348            let seq = self
1349                .make_checkpoint(std::mem::take(&mut grouped_pending_checkpoints))
1350                .await?;
1351
1352            // Count only on success; a failed build retries the same
1353            // group and would otherwise double-count it.
1354            self.metrics
1355                .commits_per_checkpoint
1356                .observe(commits_in_checkpoint as f64);
1357            // Advance the window anchor to the highest checkpoint just
1358            // built (a single call may emit several when chunked).
1359            last_seq = Some(seq);
1360            self.last_built.send_if_modified(|cur| {
1361                // when rebuilding checkpoints at startup, seq can be for an old checkpoint
1362                if seq > *cur {
1363                    *cur = seq;
1364                    true
1365                } else {
1366                    false
1367                }
1368            });
1369
1370            // Ensure that the task can be cancelled at end of epoch, even if no other await
1371            // yields execution.
1372            tokio::task::yield_now().await;
1373        }
1374        debug!(
1375            "Waiting for more checkpoints from consensus after processing {last_height:?}; {} pending checkpoints left unprocessed until next interval",
1376            grouped_pending_checkpoints.len(),
1377        );
1378
1379        Ok(())
1380    }
1381
1382    #[instrument(level = "debug", skip_all, fields(last_height = pendings.last().unwrap().details().checkpoint_height
1383    ))]
1384    async fn make_checkpoint(
1385        &self,
1386        pendings: Vec<PendingCheckpoint>,
1387    ) -> CheckpointBuilderResult<CheckpointSequenceNumber> {
1388        let _scope = monitored_scope("CheckpointBuilder::make_checkpoint");
1389        let last_details = pendings.last().unwrap().details().clone();
1390
1391        // Stores the transactions that should be included in the checkpoint.
1392        // Transactions will be recorded in the checkpoint in this order.
1393        let highest_executed_sequence = self
1394            .store
1395            .get_highest_executed_checkpoint_seq_number()
1396            .expect("db error")
1397            .unwrap_or(0);
1398
1399        let (poll_count, result) = poll_count(self.resolve_checkpoint_transactions(pendings)).await;
1400        let (sorted_tx_effects_included_in_checkpoint, all_roots) = result?;
1401
1402        let new_checkpoints = self
1403            .create_checkpoints(
1404                sorted_tx_effects_included_in_checkpoint,
1405                &last_details,
1406                &all_roots,
1407            )
1408            .await?;
1409        let highest_sequence = new_checkpoints.last().summary.sequence_number();
1410        if highest_sequence <= highest_executed_sequence && poll_count > 1 {
1411            debug_fatal!(
1412                "resolve_checkpoint_transactions should be instantaneous when executed checkpoint is ahead of checkpoint builder"
1413            );
1414        }
1415
1416        self.write_checkpoints(last_details.checkpoint_height, new_checkpoints)
1417            .await?;
1418        Ok(highest_sequence)
1419    }
1420
1421    // Given the root transactions of the pending checkpoints, resolve the
1422    // transactions that should be included in the checkpoint, and return them in
1423    // the order they should be included in the checkpoint, together with all the
1424    // resolved roots.
1425    #[instrument(level = "debug", skip_all)]
1426    async fn resolve_checkpoint_transactions(
1427        &self,
1428        pending_checkpoints: Vec<PendingCheckpoint>,
1429    ) -> IotaResult<(Vec<TransactionEffects>, HashSet<TransactionDigest>)> {
1430        let _scope = monitored_scope("CheckpointBuilder::resolve_checkpoint_transactions");
1431
1432        // Keeps track of the effects that are already included in the current
1433        // checkpoint. This is used when there are multiple pending checkpoints
1434        // to create a single checkpoint because in such scenarios, dependencies
1435        // of a transaction may in earlier created checkpoints, or in earlier
1436        // pending checkpoints.
1437        let mut effects_in_current_checkpoint = BTreeSet::new();
1438
1439        let mut tx_effects = Vec::new();
1440        let mut tx_roots = HashSet::new();
1441
1442        for pending_checkpoint in pending_checkpoints.into_iter() {
1443            let pending = pending_checkpoint.into_v1();
1444            debug!(
1445                checkpoint_commit_height = pending.details.checkpoint_height,
1446                roots = ?pending.roots,
1447                "Resolving checkpoint transactions for pending checkpoint.",
1448            );
1449
1450            let roots = &pending.roots;
1451
1452            self.metrics
1453                .checkpoint_roots_count
1454                .inc_by(roots.len() as u64);
1455
1456            let root_digests = self
1457                .epoch_store
1458                .notify_read_tx_key_to_digest(roots)
1459                .in_monitored_scope("CheckpointNotifyDigests")
1460                .await?;
1461            let root_effects = self
1462                .effects_store
1463                .try_notify_read_executed_effects(
1464                    CHECKPOINT_BUILDER_NOTIFY_READ_TASK_NAME,
1465                    &root_digests,
1466                )
1467                .in_monitored_scope("CheckpointNotifyRead")
1468                .await?;
1469
1470            let consensus_commit_prologue = {
1471                // If the roots contains consensus commit prologue transaction, we want to
1472                // extract it, and put it to the front of the checkpoint.
1473
1474                let consensus_commit_prologue =
1475                    self.extract_consensus_commit_prologue(&root_digests, &root_effects)?;
1476
1477                // Get the un-included dependencies of the consensus commit prologue. We should
1478                // expect no other dependencies that haven't been included in
1479                // any previous checkpoints.
1480                if let Some((ccp_digest, ccp_effects)) = &consensus_commit_prologue {
1481                    let unsorted_ccp = self.complete_checkpoint_effects(
1482                        vec![ccp_effects.clone()],
1483                        &mut effects_in_current_checkpoint,
1484                    )?;
1485
1486                    // No other dependencies of this consensus commit prologue that haven't been
1487                    // included in any previous checkpoint.
1488                    if unsorted_ccp.len() != 1 {
1489                        fatal!(
1490                            "Expected 1 consensus commit prologue, got {:?}",
1491                            unsorted_ccp
1492                                .iter()
1493                                .map(|e| e.transaction_digest())
1494                                .collect::<Vec<_>>()
1495                        );
1496                    }
1497                    assert_eq!(unsorted_ccp.len(), 1);
1498                    assert_eq!(unsorted_ccp[0].transaction_digest(), ccp_digest);
1499                }
1500                consensus_commit_prologue
1501            };
1502
1503            let unsorted =
1504                self.complete_checkpoint_effects(root_effects, &mut effects_in_current_checkpoint)?;
1505
1506            let _scope = monitored_scope("CheckpointBuilder::causal_sort");
1507            let mut sorted: Vec<TransactionEffects> = Vec::with_capacity(unsorted.len() + 1);
1508            if let Some((_ccp_digest, ccp_effects)) = consensus_commit_prologue {
1509                #[cfg(debug_assertions)]
1510                {
1511                    // When consensus_commit_prologue is extracted, it should not be included in the
1512                    // `unsorted`.
1513                    for tx in unsorted.iter() {
1514                        assert!(tx.transaction_digest() != &_ccp_digest);
1515                    }
1516                }
1517                sorted.push(ccp_effects);
1518            }
1519            sorted.extend(CausalOrder::causal_sort(unsorted));
1520
1521            #[cfg(msim)]
1522            {
1523                // Check consensus commit prologue invariants in sim test.
1524                self.expensive_consensus_commit_prologue_invariants_check(&root_digests, &sorted);
1525            }
1526
1527            tx_effects.extend(sorted);
1528            tx_roots.extend(root_digests);
1529        }
1530
1531        Ok((tx_effects, tx_roots))
1532    }
1533
1534    // This function is used to extract the consensus commit prologue digest and
1535    // effects from the root transactions.
1536    // The consensus commit prologue is expected to be the first transaction in the
1537    // roots.
1538    fn extract_consensus_commit_prologue(
1539        &self,
1540        root_digests: &[TransactionDigest],
1541        root_effects: &[TransactionEffects],
1542    ) -> IotaResult<Option<(TransactionDigest, TransactionEffects)>> {
1543        let _scope = monitored_scope("CheckpointBuilder::extract_consensus_commit_prologue");
1544        if root_digests.is_empty() {
1545            return Ok(None);
1546        }
1547
1548        // Reads the first transaction in the roots, and checks whether it is a
1549        // consensus commit prologue transaction.
1550        // The consensus commit prologue transaction should be the first transaction
1551        // in the roots written by the consensus handler.
1552        let first_tx = self
1553            .state
1554            .get_transaction_cache_reader()
1555            .try_get_transaction_block(&root_digests[0])?
1556            .expect("Transaction block must exist");
1557
1558        Ok(match first_tx.transaction().kind() {
1559            TransactionKind::ConsensusCommitPrologueV1(_) => {
1560                assert_eq!(first_tx.digest(), root_effects[0].transaction_digest());
1561                Some((*first_tx.digest(), root_effects[0].clone()))
1562            }
1563            _ => None,
1564        })
1565    }
1566
1567    /// Writes the new checkpoints to the DB storage and processes them.
1568    #[instrument(level = "debug", skip_all)]
1569    async fn write_checkpoints(
1570        &self,
1571        height: CheckpointHeight,
1572        mut new_checkpoints: NonEmpty<BuiltCheckpoint>,
1573    ) -> IotaResult {
1574        let _scope = monitored_scope("CheckpointBuilder::write_checkpoints");
1575        let mut batch = self.store.tables.checkpoint_content.batch();
1576        let mut all_tx_digests =
1577            Vec::with_capacity(new_checkpoints.iter().map(|c| c.contents.len()).sum());
1578
1579        // Write the new checkpoints to the DB storage.
1580        for BuiltCheckpoint {
1581            summary, contents, ..
1582        } in &new_checkpoints
1583        {
1584            debug!(
1585                checkpoint_commit_height = height,
1586                checkpoint_seq = summary.sequence_number,
1587                contents_digest = ?contents.digest(),
1588                "writing checkpoint",
1589            );
1590
1591            if let Some(previously_computed_summary) = self
1592                .store
1593                .tables
1594                .locally_computed_checkpoints
1595                .get(&summary.sequence_number)?
1596            {
1597                if previously_computed_summary != *summary {
1598                    // Panic so that we don't send out an equivocating checkpoint sig.
1599                    fatal!(
1600                        "Checkpoint {} was previously built with a different result: {previously_computed_summary:?} vs {summary:?}",
1601                        summary.sequence_number,
1602                    );
1603                }
1604            }
1605
1606            all_tx_digests.extend(contents.iter().map(|digests| digests.transaction));
1607
1608            self.metrics
1609                .transactions_included_in_checkpoint
1610                .inc_by(contents.len() as u64);
1611            let sequence_number = summary.sequence_number;
1612            self.metrics
1613                .last_constructed_checkpoint
1614                .set(sequence_number as i64);
1615
1616            batch.insert_batch(
1617                &self.store.tables.checkpoint_content,
1618                [(contents.digest(), contents)],
1619            )?;
1620
1621            batch.insert_batch(
1622                &self.store.tables.locally_computed_checkpoints,
1623                [(sequence_number, summary)],
1624            )?;
1625        }
1626
1627        batch.write()?;
1628
1629        // Cache the full contents only now that the checkpoint_content rows
1630        // are durable
1631        for checkpoint in new_checkpoints.iter_mut() {
1632            if let Some(full_contents) = checkpoint.full_contents.take() {
1633                self.store.cache_full_checkpoint_contents(
1634                    checkpoint.summary.sequence_number,
1635                    checkpoint.summary.contents_digest,
1636                    full_contents,
1637                );
1638            }
1639        }
1640
1641        // Send all checkpoint sigs to consensus. The messages including
1642        // MisbehaviorReports are also sent in this step.
1643        for BuiltCheckpoint {
1644            summary, contents, ..
1645        } in &new_checkpoints
1646        {
1647            self.output
1648                .checkpoint_created(summary, contents, &self.epoch_store, &self.store)
1649                .await?;
1650        }
1651
1652        for BuiltCheckpoint {
1653            summary: local_checkpoint,
1654            ..
1655        } in &new_checkpoints
1656        {
1657            if let Some(certified_checkpoint) = self
1658                .store
1659                .tables
1660                .certified_checkpoints
1661                .get(&local_checkpoint.sequence_number())?
1662            {
1663                self.store
1664                    .check_for_checkpoint_fork(local_checkpoint, &certified_checkpoint.into());
1665            }
1666        }
1667
1668        self.notify_aggregator.notify_one();
1669        self.epoch_store.process_constructed_checkpoint(
1670            height,
1671            new_checkpoints.map(|c| (c.summary, c.contents)),
1672        );
1673        Ok(())
1674    }
1675
1676    #[expect(clippy::type_complexity)]
1677    fn split_checkpoint_chunks(
1678        &self,
1679        transactions_effects_and_sizes: Vec<(TransactionEnvelope, TransactionEffects, usize)>,
1680        signatures: Vec<Vec<UserSignature>>,
1681    ) -> CheckpointBuilderResult<
1682        Vec<Vec<(TransactionEnvelope, TransactionEffects, Vec<UserSignature>)>>,
1683    > {
1684        let _guard = monitored_scope("CheckpointBuilder::split_checkpoint_chunks");
1685        let mut chunks = Vec::new();
1686        let mut chunk = Vec::new();
1687        let mut chunk_size: usize = 0;
1688        for ((transaction, effects, transaction_size), signatures) in
1689            transactions_effects_and_sizes.into_iter().zip(signatures)
1690        {
1691            // Roll over to a new chunk after either max count or max size is reached.
1692            // The size calculation here is intended to estimate the size of the
1693            // FullCheckpointContents struct. If this code is modified, that struct
1694            // should also be updated accordingly.
1695            let size = transaction_size
1696                + bcs::serialized_size(&effects)?
1697                + bcs::serialized_size(&signatures)?;
1698            if chunk.len() == self.max_transactions_per_checkpoint
1699                || (chunk_size + size) > self.max_checkpoint_size_bytes
1700            {
1701                if chunk.is_empty() {
1702                    // Always allow at least one tx in a checkpoint.
1703                    warn!(
1704                        "Size of single transaction ({size}) exceeds max checkpoint size ({}); allowing excessively large checkpoint to go through.",
1705                        self.max_checkpoint_size_bytes
1706                    );
1707                } else {
1708                    chunks.push(chunk);
1709                    chunk = Vec::new();
1710                    chunk_size = 0;
1711                }
1712            }
1713
1714            chunk.push((transaction, effects, signatures));
1715            chunk_size += size;
1716        }
1717
1718        if !chunk.is_empty() || chunks.is_empty() {
1719            // We intentionally create an empty checkpoint if there is no content provided
1720            // to make a 'heartbeat' checkpoint.
1721            // Important: if some conditions are added here later, we need to make sure we
1722            // always have at least one chunk if last_pending_of_epoch is set
1723            chunks.push(chunk);
1724            // Note: empty checkpoints are ok - they shouldn't happen at all on
1725            // a network with even modest load. Even if they do
1726            // happen, it is still useful as it allows fullnodes to
1727            // distinguish between "no transactions have happened" and "i am not
1728            // receiving new checkpoints".
1729        }
1730        Ok(chunks)
1731    }
1732
1733    /// Creates checkpoints using the provided transaction effects and pending
1734    /// checkpoint information.
1735    fn load_last_built_checkpoint_summary(
1736        epoch_store: &AuthorityPerEpochStore,
1737        store: &CheckpointStore,
1738    ) -> IotaResult<Option<(CheckpointSequenceNumber, CheckpointSummary)>> {
1739        let mut last_checkpoint = epoch_store.last_built_checkpoint_summary()?;
1740        if last_checkpoint.is_none() {
1741            let epoch = epoch_store.epoch();
1742            if epoch > 0 {
1743                let previous_epoch = epoch - 1;
1744                let last_verified = store.get_epoch_last_checkpoint(previous_epoch)?;
1745                last_checkpoint = last_verified.map(VerifiedCheckpoint::into_summary_and_sequence);
1746                if let Some((ref seq, _)) = last_checkpoint {
1747                    debug!(
1748                        "No checkpoints in builder DB, taking checkpoint from previous epoch with sequence {seq}"
1749                    );
1750                } else {
1751                    // This is some serious bug with when CheckpointBuilder started so surfacing it
1752                    // via panic
1753                    panic!("Can not find last checkpoint for previous epoch {previous_epoch}");
1754                }
1755            }
1756        }
1757        Ok(last_checkpoint)
1758    }
1759
1760    #[instrument(level = "debug", skip_all)]
1761    async fn create_checkpoints(
1762        &self,
1763        all_effects: Vec<TransactionEffects>,
1764        details: &PendingCheckpointInfo,
1765        all_roots: &HashSet<TransactionDigest>,
1766    ) -> CheckpointBuilderResult<NonEmpty<BuiltCheckpoint>> {
1767        let _scope = monitored_scope("CheckpointBuilder::create_checkpoints");
1768
1769        let total = all_effects.len();
1770        let mut last_checkpoint =
1771            Self::load_last_built_checkpoint_summary(&self.epoch_store, &self.store)?;
1772        let last_checkpoint_seq = last_checkpoint.as_ref().map(|(seq, _)| *seq);
1773        debug!(
1774            checkpoint_commit_height = details.checkpoint_height,
1775            next_checkpoint_seq = last_checkpoint_seq.unwrap_or_default() + 1,
1776            checkpoint_timestamp = details.timestamp_ms,
1777            "Creating checkpoint(s) for {} transactions",
1778            all_effects.len(),
1779        );
1780
1781        let all_digests: Vec<_> = all_effects
1782            .iter()
1783            .map(|effect| *effect.transaction_digest())
1784            .collect();
1785        let transactions_and_sizes = self
1786            .state
1787            .get_transaction_cache_reader()
1788            .try_get_transactions_and_serialized_sizes(&all_digests)?;
1789        let mut all_effects_and_transaction_sizes = Vec::with_capacity(all_effects.len());
1790        let mut transactions = Vec::with_capacity(all_effects.len());
1791        let mut transaction_keys = Vec::with_capacity(all_effects.len());
1792        let mut randomness_rounds = BTreeMap::new();
1793        {
1794            let _guard = monitored_scope("CheckpointBuilder::wait_for_transactions_sequenced");
1795            debug!(
1796                ?last_checkpoint_seq,
1797                "Waiting for {:?} certificates to appear in consensus",
1798                all_effects.len()
1799            );
1800
1801            for (effects, transaction_and_size) in all_effects
1802                .into_iter()
1803                .zip(transactions_and_sizes.into_iter())
1804            {
1805                let (transaction, size) = transaction_and_size
1806                    .unwrap_or_else(|| panic!("Could not find executed transaction {effects:?}"));
1807                match transaction.inner().transaction().kind() {
1808                    #[allow(deprecated)]
1809                    TransactionKind::ConsensusCommitPrologueV1(_)
1810                    | TransactionKind::AuthenticatorStateUpdateV1Deprecated => {
1811                        // ConsensusCommitPrologue is guaranteed to be
1812                        // processed before we reach here.
1813                        //
1814                        // Deprecated: Authenticator state (JWK) is deprecated
1815                        // and was never enabled.
1816                        // These transaction kinds are retained
1817                        // only for BCS enum variant compatibility.
1818                    }
1819                    TransactionKind::RandomnessStateUpdate(rsu) => {
1820                        randomness_rounds
1821                            .insert(*effects.transaction_digest(), rsu.randomness_round);
1822                    }
1823                    _ => {
1824                        // Only transactions that are not roots should be included in the call to
1825                        // `consensus_messages_processed_notify`. roots come directly from the
1826                        // consensus commit and so are known to be processed
1827                        // already.
1828                        let digest = *effects.transaction_digest();
1829                        if !all_roots.contains(&digest) {
1830                            // In P-COOL flow, user transactions are submitted as
1831                            // UserTransaction, not as Certificate.
1832                            let key = if self.epoch_store.protocol_config().enable_pcool_flow() {
1833                                ConsensusTransactionKey::UserTransaction(digest)
1834                            } else {
1835                                ConsensusTransactionKey::Certificate(digest)
1836                            };
1837                            transaction_keys.push(SequencedConsensusTransactionKey::External(key));
1838                        }
1839                    }
1840                }
1841                transactions.push(transaction);
1842                all_effects_and_transaction_sizes.push((effects, size));
1843            }
1844
1845            self.epoch_store
1846                .consensus_messages_processed_notify(transaction_keys)
1847                .await?;
1848        }
1849
1850        let signatures = self
1851            .epoch_store
1852            .user_signatures_for_checkpoint(&transactions, &all_digests)?;
1853        debug!(
1854            ?last_checkpoint_seq,
1855            "Received {} checkpoint user signatures from consensus",
1856            signatures.len()
1857        );
1858
1859        let transactions_effects_and_sizes = transactions
1860            .into_iter()
1861            .zip(all_effects_and_transaction_sizes)
1862            .map(|(transaction, (effects, size))| (transaction.into_inner(), effects, size))
1863            .collect();
1864        let chunks = self.split_checkpoint_chunks(transactions_effects_and_sizes, signatures)?;
1865        let chunks_count = chunks.len();
1866
1867        let mut checkpoints = Vec::with_capacity(chunks_count);
1868        debug!(
1869            ?last_checkpoint_seq,
1870            "Creating {} checkpoints with {} transactions", chunks_count, total,
1871        );
1872
1873        let epoch = self.epoch_store.epoch();
1874        for (index, chunk) in chunks.into_iter().enumerate() {
1875            let first_checkpoint_of_epoch = index == 0
1876                && last_checkpoint
1877                    .as_ref()
1878                    .map(|(_, c)| c.epoch != epoch)
1879                    .unwrap_or(true);
1880            if first_checkpoint_of_epoch {
1881                self.epoch_store
1882                    .record_epoch_first_checkpoint_creation_time_metric();
1883            }
1884            let last_checkpoint_of_epoch = details.last_of_epoch && index == chunks_count - 1;
1885
1886            let sequence_number = last_checkpoint
1887                .as_ref()
1888                .map(|(_, c)| c.sequence_number + 1)
1889                .unwrap_or_default();
1890            let mut timestamp_ms = details.timestamp_ms;
1891            if let Some((_, last_checkpoint)) = &last_checkpoint {
1892                if last_checkpoint.timestamp_ms > timestamp_ms {
1893                    // The first consensus commit of an epoch can have zero timestamp.
1894                    debug!(
1895                        "Decrease of checkpoint timestamp, possibly due to epoch change. Sequence: {}, previous: {}, current: {}",
1896                        sequence_number, last_checkpoint.timestamp_ms, timestamp_ms,
1897                    );
1898                    timestamp_ms = last_checkpoint.timestamp_ms;
1899                }
1900            }
1901
1902            let (chunk_transactions, mut effects, mut signatures): (
1903                Vec<TransactionEnvelope>,
1904                Vec<TransactionEffects>,
1905                Vec<Vec<UserSignature>>,
1906            ) = chunk.into_iter().multiunzip();
1907            let epoch_rolling_gas_cost_summary =
1908                self.get_epoch_total_gas_cost(last_checkpoint.as_ref().map(|(_, c)| c), &effects);
1909
1910            let end_of_epoch_data = if last_checkpoint_of_epoch {
1911                let scores: Vec<u64> = if self
1912                    .epoch_store
1913                    .protocol_config()
1914                    .pass_calculated_validator_scores_to_advance_epoch()
1915                {
1916                    self.epoch_store.scoreboard.current_scores()
1917                } else {
1918                    // Give everyone in the committee the max score
1919                    vec![MAX_SCORE; self.epoch_store.committee().num_members()]
1920                };
1921
1922                let (system_state_obj, system_epoch_info_event) = self
1923                    .augment_epoch_last_checkpoint(
1924                        &epoch_rolling_gas_cost_summary,
1925                        timestamp_ms,
1926                        &mut effects,
1927                        &mut signatures,
1928                        sequence_number,
1929                        scores,
1930                    )
1931                    .await?;
1932
1933                // The system epoch info event can be `None` in case if the `advance_epoch`
1934                // Move function call failed and was executed in the safe mode.
1935                // In this case, the tokens supply should be unchanged.
1936                //
1937                // SAFETY: The number of minted and burnt tokens easily fit into an i64 and due
1938                // to those small numbers, no overflows will occur during conversion or
1939                // subtraction.
1940                let epoch_supply_change =
1941                    system_epoch_info_event.map_or(0, |event| event.supply_change());
1942
1943                let committee = system_state_obj
1944                    .get_current_epoch_committee()
1945                    .committee()
1946                    .clone();
1947
1948                // This must happen after the call to augment_epoch_last_checkpoint,
1949                // otherwise we will not capture the change_epoch tx.
1950                let root_state_digest = {
1951                    let state_acc = self
1952                        .global_state_hasher
1953                        .upgrade()
1954                        .expect("No checkpoints should be getting built after local configuration");
1955                    let acc = state_acc.accumulate_checkpoint(
1956                        &effects,
1957                        sequence_number,
1958                        &self.epoch_store,
1959                    )?;
1960
1961                    state_acc
1962                        .wait_for_previous_running_root(&self.epoch_store, sequence_number)
1963                        .await?;
1964
1965                    state_acc.accumulate_running_root(
1966                        &self.epoch_store,
1967                        sequence_number,
1968                        Some(acc),
1969                    )?;
1970                    state_acc
1971                        .digest_epoch(self.epoch_store.clone(), sequence_number)
1972                        .await?
1973                };
1974                self.metrics.highest_accumulated_epoch.set(epoch as i64);
1975                info!("Epoch {epoch} root state hash digest: {root_state_digest:?}");
1976
1977                let epoch_commitments = vec![CheckpointCommitment::EcmhLiveObjectSet {
1978                    digest: root_state_digest.digest,
1979                }];
1980
1981                Some(EndOfEpochData {
1982                    next_epoch_committee: committee.committee_members(),
1983                    next_epoch_protocol_version: system_state_obj.protocol_version(),
1984                    epoch_commitments,
1985                    epoch_supply_change,
1986                })
1987            } else {
1988                self.send_to_hasher
1989                    .send((sequence_number, effects.clone()))
1990                    .await?;
1991                None
1992            };
1993            let contents = CheckpointContents::new_with_digests_and_signatures(
1994                effects.iter().map(TransactionEffects::execution_digests),
1995                signatures,
1996            );
1997
1998            let num_txns = contents.len() as u64;
1999
2000            let network_total_transactions = last_checkpoint
2001                .as_ref()
2002                .map(|(_, c)| c.network_total_transactions + num_txns)
2003                .unwrap_or(num_txns);
2004
2005            let previous_digest = last_checkpoint.as_ref().map(|(_, c)| c.digest());
2006
2007            let matching_randomness_rounds: Vec<_> = effects
2008                .iter()
2009                .filter_map(|e| randomness_rounds.get(e.transaction_digest()))
2010                .copied()
2011                .collect();
2012
2013            let summary = CheckpointSummary::new_with_protocol_config(
2014                self.epoch_store.protocol_config(),
2015                epoch,
2016                sequence_number,
2017                network_total_transactions,
2018                &contents,
2019                previous_digest,
2020                epoch_rolling_gas_cost_summary,
2021                end_of_epoch_data,
2022                timestamp_ms,
2023                matching_randomness_rounds,
2024            );
2025            summary.report_checkpoint_age(&self.metrics.last_created_checkpoint_age);
2026            if last_checkpoint_of_epoch {
2027                info!(
2028                    checkpoint_seq = sequence_number,
2029                    "creating last checkpoint of epoch {}", epoch
2030                );
2031                if let Some(stats) = self.store.get_epoch_stats(epoch, &summary) {
2032                    self.epoch_store
2033                        .report_epoch_metrics_at_last_checkpoint(stats);
2034                }
2035            }
2036
2037            // Assemble the full contents for faster checkpoint propagation to
2038            // peers. End-of-epoch checkpoints carry an appended change-epoch
2039            // transaction not tracked here; the executor caches those via its
2040            // synced path.
2041            let full_contents = (!last_checkpoint_of_epoch
2042                && self
2043                    .store
2044                    .should_cache_full_checkpoint_contents(sequence_number))
2045            .then(|| {
2046                let execution_data = chunk_transactions
2047                    .into_iter()
2048                    .zip(effects.iter().cloned())
2049                    .map(|(transaction, effects)| ExecutionData::new(transaction, effects));
2050                FullCheckpointContents::from_contents_and_execution_data(
2051                    contents.clone(),
2052                    execution_data,
2053                )
2054                // `contents` was built from these same transactions and
2055                // effects a few lines above, so the two cannot differ in
2056                // length.
2057                .expect("checkpoint builder pins one signature set per transaction it includes")
2058            });
2059
2060            last_checkpoint = Some((sequence_number, summary.clone()));
2061            checkpoints.push(BuiltCheckpoint {
2062                summary,
2063                contents,
2064                full_contents,
2065            });
2066        }
2067
2068        Ok(NonEmpty::from_vec(checkpoints).expect("at least one checkpoint"))
2069    }
2070
2071    fn get_epoch_total_gas_cost(
2072        &self,
2073        last_checkpoint: Option<&CheckpointSummary>,
2074        cur_checkpoint_effects: &[TransactionEffects],
2075    ) -> GasCostSummary {
2076        let (previous_epoch, previous_gas_costs) = last_checkpoint
2077            .map(|c| (c.epoch, c.epoch_rolling_gas_cost_summary.clone()))
2078            .unwrap_or_default();
2079        let current_gas_costs =
2080            checked::new_gas_cost_summary_from_txn_effects(cur_checkpoint_effects.iter());
2081        if previous_epoch == self.epoch_store.epoch() {
2082            // sum only when we are within the same epoch
2083            GasCostSummary::new(
2084                previous_gas_costs.computation_cost + current_gas_costs.computation_cost,
2085                previous_gas_costs.computation_cost_burned
2086                    + current_gas_costs.computation_cost_burned,
2087                previous_gas_costs.storage_cost + current_gas_costs.storage_cost,
2088                previous_gas_costs.storage_rebate + current_gas_costs.storage_rebate,
2089                previous_gas_costs.non_refundable_storage_fee
2090                    + current_gas_costs.non_refundable_storage_fee,
2091            )
2092        } else {
2093            current_gas_costs
2094        }
2095    }
2096
2097    /// Augments the last checkpoint of the epoch by creating and executing an
2098    /// advance epoch transaction.
2099    #[instrument(level = "error", skip_all)]
2100    async fn augment_epoch_last_checkpoint(
2101        &self,
2102        epoch_total_gas_cost: &GasCostSummary,
2103        epoch_start_timestamp_ms: CheckpointTimestamp,
2104        checkpoint_effects: &mut Vec<TransactionEffects>,
2105        signatures: &mut Vec<Vec<UserSignature>>,
2106        checkpoint: CheckpointSequenceNumber,
2107        scores: Vec<u64>,
2108    ) -> CheckpointBuilderResult<(IotaSystemState, Option<SystemEpochInfoEvent>)> {
2109        let (system_state, system_epoch_info_event, effects) = self
2110            .state
2111            .create_and_execute_advance_epoch_tx(
2112                &self.epoch_store,
2113                epoch_total_gas_cost,
2114                checkpoint,
2115                epoch_start_timestamp_ms,
2116                scores,
2117            )
2118            .await?;
2119        checkpoint_effects.push(effects);
2120        signatures.push(vec![]);
2121        Ok((system_state, system_epoch_info_event))
2122    }
2123
2124    /// For the given roots return complete list of effects to include in
2125    /// checkpoint This list includes the roots and all their dependencies,
2126    /// which are not part of checkpoint already. Note that this function
2127    /// may be called multiple times to construct the checkpoint.
2128    /// `existing_tx_digests_in_checkpoint` is used to track the transactions
2129    /// that are already included in the checkpoint. Txs in `roots` that
2130    /// need to be included in the checkpoint will be added to
2131    /// `existing_tx_digests_in_checkpoint` after the call of this function.
2132    #[instrument(level = "debug", skip_all)]
2133    fn complete_checkpoint_effects(
2134        &self,
2135        mut roots: Vec<TransactionEffects>,
2136        existing_tx_digests_in_checkpoint: &mut BTreeSet<TransactionDigest>,
2137    ) -> IotaResult<Vec<TransactionEffects>> {
2138        let _scope = monitored_scope("CheckpointBuilder::complete_checkpoint_effects");
2139        let mut results = vec![];
2140        let mut seen = HashSet::new();
2141        loop {
2142            let mut pending = HashSet::new();
2143
2144            let transactions_included = self
2145                .epoch_store
2146                .builder_included_transactions_in_checkpoint(
2147                    roots.iter().map(|e| e.transaction_digest()),
2148                )?;
2149
2150            for (effect, tx_included) in roots.into_iter().zip(transactions_included.into_iter()) {
2151                let digest = effect.transaction_digest();
2152                // Unnecessary to read effects of a dependency if the effect is already
2153                // processed.
2154                seen.insert(*digest);
2155
2156                // Skip roots that are already included in the checkpoint.
2157                if existing_tx_digests_in_checkpoint.contains(effect.transaction_digest()) {
2158                    continue;
2159                }
2160
2161                // Skip roots already included in checkpoints or roots from previous epochs
2162                if tx_included || effect.epoch() < self.epoch_store.epoch() {
2163                    continue;
2164                }
2165
2166                let existing_effects = self
2167                    .epoch_store
2168                    .transactions_executed_in_cur_epoch(effect.dependencies())?;
2169
2170                for (dependency, effects_signature_exists) in
2171                    effect.dependencies().iter().zip(existing_effects.iter())
2172                {
2173                    // Skip here if dependency not executed in the current epoch.
2174                    // Note that the existence of an effects signature in the
2175                    // epoch store for the given digest indicates that the transaction
2176                    // was locally executed in the current epoch
2177                    if !effects_signature_exists {
2178                        continue;
2179                    }
2180                    if seen.insert(*dependency) {
2181                        pending.insert(*dependency);
2182                    }
2183                }
2184                results.push(effect);
2185            }
2186            if pending.is_empty() {
2187                break;
2188            }
2189            let pending = pending.into_iter().collect::<Vec<_>>();
2190            let effects = self
2191                .effects_store
2192                .try_multi_get_executed_effects(&pending)?;
2193            let effects = effects
2194                .into_iter()
2195                .zip(pending)
2196                .map(|(opt, digest)| match opt {
2197                    Some(x) => x,
2198                    None => panic!(
2199                        "Can not find effect for transaction {digest}, however transaction that depend on it was already executed"
2200                    ),
2201                })
2202                .collect::<Vec<_>>();
2203            roots = effects;
2204        }
2205
2206        existing_tx_digests_in_checkpoint.extend(results.iter().map(|e| e.transaction_digest()));
2207        Ok(results)
2208    }
2209
2210    // This function is used to check the invariants of the consensus commit
2211    // prologue transactions in the checkpoint in simtest.
2212    #[cfg(msim)]
2213    fn expensive_consensus_commit_prologue_invariants_check(
2214        &self,
2215        root_digests: &[TransactionDigest],
2216        sorted: &[TransactionEffects],
2217    ) {
2218        // Gets all the consensus commit prologue transactions from the roots.
2219        let root_txs = self
2220            .state
2221            .get_transaction_cache_reader()
2222            .multi_get_transaction_blocks(root_digests);
2223        let ccps = root_txs
2224            .iter()
2225            .filter_map(|tx| {
2226                tx.as_ref().filter(|tx| {
2227                    matches!(
2228                        tx.transaction().kind(),
2229                        TransactionKind::ConsensusCommitPrologueV1(_)
2230                    )
2231                })
2232            })
2233            .collect::<Vec<_>>();
2234
2235        // There should be at most one consensus commit prologue transaction in the
2236        // roots.
2237        assert!(ccps.len() <= 1);
2238
2239        // Get all the transactions in the checkpoint.
2240        let txs = self
2241            .state
2242            .get_transaction_cache_reader()
2243            .multi_get_transaction_blocks(
2244                &sorted
2245                    .iter()
2246                    .map(|tx| *tx.transaction_digest())
2247                    .collect::<Vec<_>>(),
2248            );
2249
2250        if ccps.is_empty() {
2251            // If there is no consensus commit prologue transaction in the roots, then there
2252            // should be no consensus commit prologue transaction in the
2253            // checkpoint.
2254            for tx in txs.iter().flatten() {
2255                assert!(!matches!(
2256                    tx.transaction().kind(),
2257                    TransactionKind::ConsensusCommitPrologueV1(_)
2258                ));
2259            }
2260        } else {
2261            // If there is one consensus commit prologue, it must be the first one in the
2262            // checkpoint.
2263            assert!(matches!(
2264                txs[0].as_ref().unwrap().transaction().kind(),
2265                TransactionKind::ConsensusCommitPrologueV1(_)
2266            ));
2267
2268            assert_eq!(ccps[0].digest(), txs[0].as_ref().unwrap().digest());
2269
2270            for tx in txs.iter().skip(1).flatten() {
2271                assert!(!matches!(
2272                    tx.transaction().kind(),
2273                    TransactionKind::ConsensusCommitPrologueV1(_)
2274                ));
2275            }
2276        }
2277    }
2278}
2279
2280impl CheckpointAggregator {
2281    fn new(
2282        tables: Arc<CheckpointStore>,
2283        epoch_store: Arc<AuthorityPerEpochStore>,
2284        notify: Arc<Notify>,
2285        output: Box<dyn CertifiedCheckpointOutput>,
2286        state: Arc<AuthorityState>,
2287        metrics: Arc<CheckpointMetrics>,
2288    ) -> Self {
2289        let current = None;
2290        Self {
2291            store: tables,
2292            epoch_store,
2293            notify,
2294            current,
2295            output,
2296            state,
2297            metrics,
2298        }
2299    }
2300
2301    /// Runs the `CheckpointAggregator` in an asynchronous loop, managing the
2302    /// aggregation of checkpoints.
2303    /// The function ensures continuous aggregation of checkpoints, handling
2304    /// errors and retries gracefully, and allowing for proper shutdown on
2305    /// receiving an exit signal.
2306    async fn run(mut self) {
2307        info!("Starting CheckpointAggregator");
2308        loop {
2309            if let Err(e) = self.run_and_notify().await {
2310                error!(
2311                    "Error while aggregating checkpoint, will retry in 1s: {:?}",
2312                    e
2313                );
2314                self.metrics.checkpoint_errors.inc();
2315                tokio::time::sleep(Duration::from_secs(1)).await;
2316                continue;
2317            }
2318
2319            let _ = timeout(Duration::from_secs(1), self.notify.notified()).await;
2320        }
2321    }
2322
2323    async fn run_and_notify(&mut self) -> IotaResult {
2324        let summaries = self.run_inner()?;
2325        for summary in summaries {
2326            self.output.certified_checkpoint_created(&summary).await?;
2327        }
2328        Ok(())
2329    }
2330
2331    fn run_inner(&mut self) -> IotaResult<Vec<CertifiedCheckpointSummary>> {
2332        let _scope = monitored_scope("CheckpointAggregator");
2333        let mut result = vec![];
2334        'outer: loop {
2335            let next_to_certify = self.next_checkpoint_to_certify()?;
2336            let current = if let Some(current) = &mut self.current {
2337                // It's possible that the checkpoint was already certified by
2338                // the rest of the network and we've already received the
2339                // certified checkpoint via StateSync. In this case, we reset
2340                // the current signature aggregator to the next checkpoint to
2341                // be certified
2342                if current.summary.sequence_number < next_to_certify {
2343                    self.current = None;
2344                    continue;
2345                }
2346                current
2347            } else {
2348                let Some(summary) = self
2349                    .epoch_store
2350                    .get_built_checkpoint_summary(next_to_certify)?
2351                else {
2352                    return Ok(result);
2353                };
2354                self.current = Some(CheckpointSignatureAggregator {
2355                    next_index: 0,
2356                    digest: summary.digest(),
2357                    summary,
2358                    signatures_by_digest: MultiStakeAggregator::new(
2359                        self.epoch_store.committee().clone(),
2360                    ),
2361                    store: self.store.clone(),
2362                    state: self.state.clone(),
2363                    metrics: self.metrics.clone(),
2364                });
2365                self.current.as_mut().unwrap()
2366            };
2367
2368            let epoch_tables = self
2369                .epoch_store
2370                .tables()
2371                .expect("should not run past end of epoch");
2372            let iter = epoch_tables
2373                .pending_checkpoint_signatures
2374                .safe_iter_with_bounds(
2375                    Some((current.summary.sequence_number, current.next_index)),
2376                    None,
2377                );
2378            for item in iter {
2379                let ((seq, index), data) = item?;
2380                if seq != current.summary.sequence_number {
2381                    trace!(
2382                        checkpoint_seq =? current.summary.sequence_number,
2383                        "Not enough checkpoint signatures",
2384                    );
2385                    // No more signatures (yet) for this checkpoint
2386                    return Ok(result);
2387                }
2388                trace!(
2389                    checkpoint_seq = current.summary.sequence_number,
2390                    "Processing signature for checkpoint (digest: {:?}) from {:?}",
2391                    current.summary.digest(),
2392                    data.summary.auth_sig().authority.concise()
2393                );
2394                self.metrics
2395                    .checkpoint_participation
2396                    .with_label_values(&[&format!(
2397                        "{:?}",
2398                        data.summary.auth_sig().authority.concise()
2399                    )])
2400                    .inc();
2401                if let Ok(auth_signature) = current.try_aggregate(data) {
2402                    debug!(
2403                        checkpoint_seq = current.summary.sequence_number,
2404                        "Successfully aggregated signatures for checkpoint (digest: {:?})",
2405                        current.summary.digest(),
2406                    );
2407                    let summary = VerifiedCheckpoint::new_unchecked(
2408                        CertifiedCheckpointSummary::new_from_data_and_sig(
2409                            current.summary.clone(),
2410                            auth_signature,
2411                        ),
2412                    );
2413
2414                    self.store.insert_certified_checkpoint(&summary)?;
2415                    self.metrics
2416                        .last_certified_checkpoint
2417                        .set(current.summary.sequence_number as i64);
2418                    current
2419                        .summary
2420                        .report_checkpoint_age(&self.metrics.last_certified_checkpoint_age);
2421                    result.push(summary.into_inner());
2422                    self.current = None;
2423                    continue 'outer;
2424                } else {
2425                    current.next_index = index + 1;
2426                }
2427            }
2428            break;
2429        }
2430        Ok(result)
2431    }
2432
2433    fn next_checkpoint_to_certify(&self) -> IotaResult<CheckpointSequenceNumber> {
2434        Ok(self
2435            .store
2436            .tables
2437            .certified_checkpoints
2438            .safe_range_iter_reversed(..)
2439            .next()
2440            .transpose()?
2441            .map(|(seq, _)| seq + 1)
2442            .unwrap_or_default())
2443    }
2444}
2445
2446impl CheckpointSignatureAggregator {
2447    #[expect(clippy::result_unit_err)]
2448    pub fn try_aggregate(
2449        &mut self,
2450        data: CheckpointSignatureMessage,
2451    ) -> Result<AuthorityStrongQuorumSignInfo, ()> {
2452        let their_digest = *data.summary.digest();
2453        let (_, signature) = data.summary.into_data_and_sig();
2454        let author = signature.authority;
2455        let envelope =
2456            SignedCheckpointSummary::new_from_data_and_sig(self.summary.clone(), signature);
2457        match self.signatures_by_digest.insert(their_digest, envelope) {
2458            // ignore repeated signatures
2459            InsertResult::Failed {
2460                error:
2461                    IotaError::StakeAggregatorRepeatedSigner {
2462                        conflicting_sig: false,
2463                        ..
2464                    },
2465            } => Err(()),
2466            InsertResult::Failed { error } => {
2467                warn!(
2468                    checkpoint_seq = self.summary.sequence_number,
2469                    "Failed to aggregate new signature from validator {:?}: {:?}",
2470                    author.concise(),
2471                    error
2472                );
2473                self.check_for_split_brain();
2474                Err(())
2475            }
2476            InsertResult::QuorumReached(cert) => {
2477                // It is not guaranteed that signature.authority == consensus_cert.author, but
2478                // we do verify the signature so we know that the author signed
2479                // the message at some point.
2480                if their_digest != self.digest {
2481                    self.metrics.remote_checkpoint_forks.inc();
2482                    warn!(
2483                        checkpoint_seq = self.summary.sequence_number,
2484                        "Validator {:?} has mismatching checkpoint digest {}, we have digest {}",
2485                        author.concise(),
2486                        their_digest,
2487                        self.digest
2488                    );
2489                    return Err(());
2490                }
2491                Ok(cert)
2492            }
2493            InsertResult::NotEnoughVotes {
2494                bad_votes: _,
2495                bad_authorities: _,
2496            } => {
2497                self.check_for_split_brain();
2498                Err(())
2499            }
2500        }
2501    }
2502
2503    /// Check if there is a split brain condition in checkpoint signature
2504    /// aggregation, defined as any state wherein it is no longer possible
2505    /// to achieve quorum on a checkpoint proposal, irrespective of the
2506    /// outcome of any outstanding votes.
2507    fn check_for_split_brain(&self) {
2508        debug!(
2509            checkpoint_seq = self.summary.sequence_number,
2510            "Checking for split brain condition"
2511        );
2512        if self.signatures_by_digest.quorum_unreachable() {
2513            // TODO: at this point we should immediately halt processing
2514            // of new transaction certificates to avoid building on top of
2515            // forked output
2516            // self.halt_all_execution();
2517
2518            let digests_by_stake_messages = self
2519                .signatures_by_digest
2520                .get_all_unique_values()
2521                .into_iter()
2522                .sorted_by_key(|(_, (_, stake))| -(*stake as i64))
2523                .map(|(digest, (_authorities, total_stake))| {
2524                    format!("{digest} (total stake: {total_stake})")
2525                })
2526                .collect::<Vec<String>>();
2527            debug_fatal!(
2528                "Split brain detected in checkpoint signature aggregation for checkpoint {:?}. Remaining stake: {:?}, Digests by stake: {:?}",
2529                self.summary.sequence_number,
2530                self.signatures_by_digest.uncommitted_stake(),
2531                digests_by_stake_messages,
2532            );
2533            self.metrics.split_brain_checkpoint_forks.inc();
2534
2535            let all_unique_values = self.signatures_by_digest.get_all_unique_values();
2536            let local_summary = self.summary.clone();
2537            let state = self.state.clone();
2538            let tables = self.store.clone();
2539
2540            tokio::spawn(async move {
2541                diagnose_split_brain(all_unique_values, local_summary, state, tables).await;
2542            });
2543        }
2544    }
2545}
2546
2547/// Create data dump containing relevant data for diagnosing cause of the
2548/// split brain by querying one disagreeing validator for full checkpoint
2549/// contents. To minimize peer chatter, we only query one validator at random
2550/// from each disagreeing faction, as all honest validators that participated in
2551/// this round may inevitably run the same process.
2552async fn diagnose_split_brain(
2553    all_unique_values: BTreeMap<CheckpointDigest, (Vec<AuthorityName>, StakeUnit)>,
2554    local_summary: CheckpointSummary,
2555    state: Arc<AuthorityState>,
2556    tables: Arc<CheckpointStore>,
2557) {
2558    debug!(
2559        checkpoint_seq = local_summary.sequence_number,
2560        "Running split brain diagnostics..."
2561    );
2562    let time = SystemTime::now();
2563    // collect one random disagreeing validator per differing digest
2564    let digest_to_validator = all_unique_values
2565        .iter()
2566        .filter_map(|(digest, (validators, _))| {
2567            if *digest != local_summary.digest() {
2568                let random_validator = validators.choose(&mut get_rng()).unwrap();
2569                Some((*digest, *random_validator))
2570            } else {
2571                None
2572            }
2573        })
2574        .collect::<HashMap<_, _>>();
2575    if digest_to_validator.is_empty() {
2576        panic!(
2577            "Given split brain condition, there should be at \
2578                least one validator that disagrees with local signature"
2579        );
2580    }
2581
2582    let epoch_store = state.load_epoch_store_one_call_per_task();
2583    let committee = epoch_store
2584        .epoch_start_state()
2585        .get_iota_committee_with_network_metadata();
2586    let network_config = default_iota_network_config();
2587    let network_clients =
2588        make_network_authority_clients_with_network_config(&committee, &network_config);
2589
2590    // Query all disagreeing validators
2591    let response_futures = digest_to_validator
2592        .values()
2593        .cloned()
2594        .map(|validator| {
2595            let client = network_clients
2596                .get(&validator)
2597                .expect("Failed to get network client");
2598            let request = CheckpointRequest {
2599                sequence_number: Some(local_summary.sequence_number),
2600                request_content: true,
2601                certified: false,
2602            };
2603            client.get_checkpoint_v2(request)
2604        })
2605        .collect::<Vec<_>>();
2606
2607    let digest_name_pair = digest_to_validator.iter();
2608    let response_data = futures::future::join_all(response_futures)
2609        .await
2610        .into_iter()
2611        .zip(digest_name_pair)
2612        .filter_map(|(response, (digest, name))| match response {
2613            Ok(response) => match response {
2614                CheckpointResponse {
2615                    checkpoint: Some(CheckpointSummaryResponse::Pending(summary)),
2616                    contents: Some(contents),
2617                } => Some((*name, *digest, summary, contents)),
2618                CheckpointResponse {
2619                    checkpoint: Some(CheckpointSummaryResponse::Certified(_)),
2620                    contents: _,
2621                } => {
2622                    panic!("Expected pending checkpoint, but got certified checkpoint");
2623                }
2624                CheckpointResponse {
2625                    checkpoint: None,
2626                    contents: _,
2627                } => {
2628                    error!(
2629                        "Summary for checkpoint {:?} not found on validator {:?}",
2630                        local_summary.sequence_number, name
2631                    );
2632                    None
2633                }
2634                CheckpointResponse {
2635                    checkpoint: _,
2636                    contents: None,
2637                } => {
2638                    error!(
2639                        "Contents for checkpoint {:?} not found on validator {:?}",
2640                        local_summary.sequence_number, name
2641                    );
2642                    None
2643                }
2644            },
2645            Err(e) => {
2646                error!(
2647                    "Failed to get checkpoint contents from validator for fork diagnostics: {:?}",
2648                    e
2649                );
2650                None
2651            }
2652        })
2653        .collect::<Vec<_>>();
2654
2655    let local_checkpoint_contents = tables
2656        .get_checkpoint_contents(&local_summary.contents_digest)
2657        .unwrap_or_else(|_| {
2658            panic!(
2659                "Could not find checkpoint contents for digest {:?}",
2660                local_summary.digest()
2661            )
2662        })
2663        .unwrap_or_else(|| {
2664            panic!(
2665                "Could not find local full checkpoint contents for checkpoint {:?}, digest {:?}",
2666                local_summary.sequence_number,
2667                local_summary.digest()
2668            )
2669        });
2670    let local_contents_text = format!("{local_checkpoint_contents:?}");
2671
2672    let local_summary_text = format!("{local_summary:?}");
2673    let local_validator = state.name.concise();
2674    let diff_patches = response_data
2675        .iter()
2676        .map(|(name, other_digest, other_summary, contents)| {
2677            let other_contents_text = format!("{contents:?}");
2678            let other_summary_text = format!("{other_summary:?}");
2679            let (local_transactions, local_effects): (Vec<_>, Vec<_>) = local_checkpoint_contents
2680                .enumerate_transactions(&local_summary)
2681                .map(|(_, exec_digest)| (exec_digest.transaction, exec_digest.effects))
2682                .unzip();
2683            let (other_transactions, other_effects): (Vec<_>, Vec<_>) = contents
2684                .enumerate_transactions(other_summary)
2685                .map(|(_, exec_digest)| (exec_digest.transaction, exec_digest.effects))
2686                .unzip();
2687            let summary_patch = create_patch(&local_summary_text, &other_summary_text);
2688            let contents_patch = create_patch(&local_contents_text, &other_contents_text);
2689            let local_transactions_text = format!("{local_transactions:#?}");
2690            let other_transactions_text = format!("{other_transactions:#?}");
2691            let transactions_patch =
2692                create_patch(&local_transactions_text, &other_transactions_text);
2693            let local_effects_text = format!("{local_effects:#?}");
2694            let other_effects_text = format!("{other_effects:#?}");
2695            let effects_patch = create_patch(&local_effects_text, &other_effects_text);
2696            let seq_number = local_summary.sequence_number;
2697            let local_digest = local_summary.digest();
2698            let other_validator = name.concise();
2699            format!(
2700                "Checkpoint: {seq_number:?}\n\
2701                Local validator (original): {local_validator:?}, digest: {local_digest}\n\
2702                Other validator (modified): {other_validator:?}, digest: {other_digest}\n\n\
2703                Summary Diff: \n{summary_patch}\n\n\
2704                Contents Diff: \n{contents_patch}\n\n\
2705                Transactions Diff: \n{transactions_patch}\n\n\
2706                Effects Diff: \n{effects_patch}",
2707            )
2708        })
2709        .collect::<Vec<_>>()
2710        .join("\n\n\n");
2711
2712    let header = format!(
2713        "Checkpoint Fork Dump - Authority {local_validator:?}: \n\
2714        Datetime: {time:?}"
2715    );
2716    let fork_logs_text = format!("{header}\n\n{diff_patches}\n\n");
2717    let checkpoint_fork_dir = iota_common::tempdir().keep();
2718    let checkpoint_fork_file_path = checkpoint_fork_dir.join(Path::new("checkpoint_fork_dump.txt"));
2719    let mut file = File::create(checkpoint_fork_file_path).unwrap();
2720    write!(file, "{fork_logs_text}").unwrap();
2721    debug!("{}", fork_logs_text);
2722}
2723
2724pub trait CheckpointServiceNotify {
2725    fn notify_checkpoint_signature(
2726        &self,
2727        epoch_store: &AuthorityPerEpochStore,
2728        info: &CheckpointSignatureMessage,
2729    ) -> IotaResult;
2730
2731    fn notify_checkpoint(&self) -> IotaResult;
2732}
2733
2734enum CheckpointServiceState {
2735    Unstarted(
2736        Box<(
2737            CheckpointBuilder,
2738            CheckpointAggregator,
2739            CheckpointStateHasher,
2740        )>,
2741    ),
2742    Started,
2743}
2744
2745impl CheckpointServiceState {
2746    fn take_unstarted(
2747        &mut self,
2748    ) -> (
2749        CheckpointBuilder,
2750        CheckpointAggregator,
2751        CheckpointStateHasher,
2752    ) {
2753        let mut state = CheckpointServiceState::Started;
2754        std::mem::swap(self, &mut state);
2755
2756        match state {
2757            CheckpointServiceState::Unstarted(tup) => (tup.0, tup.1, tup.2),
2758            CheckpointServiceState::Started => panic!("CheckpointServiceState is already started"),
2759        }
2760    }
2761}
2762
2763pub struct CheckpointService {
2764    tables: Arc<CheckpointStore>,
2765    notify_builder: Arc<Notify>,
2766    notify_aggregator: Arc<Notify>,
2767    last_signature_index: Mutex<u64>,
2768    // A notification for the current highest built sequence number.
2769    highest_currently_built_seq_tx: watch::Sender<CheckpointSequenceNumber>,
2770    // The highest sequence number that had already been built at the time CheckpointService
2771    // was constructed
2772    highest_previously_built_seq: CheckpointSequenceNumber,
2773    metrics: Arc<CheckpointMetrics>,
2774    state: Mutex<CheckpointServiceState>,
2775}
2776
2777impl CheckpointService {
2778    /// Spawns the checkpoint service, initializing and starting the checkpoint
2779    /// builder and aggregator tasks.
2780    /// Constructs a new CheckpointService in an un-started state.
2781    pub fn build(
2782        state: Arc<AuthorityState>,
2783        checkpoint_store: Arc<CheckpointStore>,
2784        epoch_store: Arc<AuthorityPerEpochStore>,
2785        effects_store: Arc<dyn TransactionCacheRead>,
2786        global_state_hasher: Weak<GlobalStateHasher>,
2787        checkpoint_output: Box<dyn CheckpointOutput>,
2788        certified_checkpoint_output: Box<dyn CertifiedCheckpointOutput>,
2789        metrics: Arc<CheckpointMetrics>,
2790        max_transactions_per_checkpoint: usize,
2791        max_checkpoint_size_bytes: usize,
2792    ) -> Arc<Self> {
2793        info!(
2794            "Starting checkpoint service with {max_transactions_per_checkpoint} max_transactions_per_checkpoint and {max_checkpoint_size_bytes} max_checkpoint_size_bytes"
2795        );
2796        let notify_builder = Arc::new(Notify::new());
2797        let notify_aggregator = Arc::new(Notify::new());
2798
2799        // We may have built higher checkpoint numbers before restarting.
2800        let highest_previously_built_seq = checkpoint_store
2801            .get_latest_locally_computed_checkpoint()
2802            .expect("failed to get latest locally computed checkpoint")
2803            .map(|s| s.sequence_number)
2804            .unwrap_or(0);
2805
2806        let highest_currently_built_seq =
2807            CheckpointBuilder::load_last_built_checkpoint_summary(&epoch_store, &checkpoint_store)
2808                .expect("epoch should not have ended")
2809                .map(|(seq, _)| seq)
2810                .unwrap_or(0);
2811
2812        let (highest_currently_built_seq_tx, _) = watch::channel(highest_currently_built_seq);
2813
2814        let aggregator = CheckpointAggregator::new(
2815            checkpoint_store.clone(),
2816            epoch_store.clone(),
2817            notify_aggregator.clone(),
2818            certified_checkpoint_output,
2819            state.clone(),
2820            metrics.clone(),
2821        );
2822
2823        let (send_to_hasher, receive_from_builder) = mpsc::channel(16);
2824
2825        let ckpt_state_hasher = CheckpointStateHasher::new(
2826            epoch_store.clone(),
2827            global_state_hasher.clone(),
2828            receive_from_builder,
2829        );
2830
2831        let builder = CheckpointBuilder::new(
2832            state,
2833            checkpoint_store.clone(),
2834            epoch_store.clone(),
2835            notify_builder.clone(),
2836            effects_store,
2837            global_state_hasher,
2838            send_to_hasher,
2839            checkpoint_output,
2840            notify_aggregator.clone(),
2841            highest_currently_built_seq_tx.clone(),
2842            metrics.clone(),
2843            max_transactions_per_checkpoint,
2844            max_checkpoint_size_bytes,
2845        );
2846
2847        let last_signature_index = epoch_store
2848            .get_last_checkpoint_signature_index()
2849            .expect("should not cross end of epoch");
2850        let last_signature_index = Mutex::new(last_signature_index);
2851
2852        Arc::new(Self {
2853            tables: checkpoint_store,
2854            notify_builder,
2855            notify_aggregator,
2856            last_signature_index,
2857            highest_currently_built_seq_tx,
2858            highest_previously_built_seq,
2859            metrics,
2860            state: Mutex::new(CheckpointServiceState::Unstarted(Box::new((
2861                builder,
2862                aggregator,
2863                ckpt_state_hasher,
2864            )))),
2865        })
2866    }
2867
2868    /// Starts the CheckpointService.
2869    ///
2870    /// This function blocks until the CheckpointBuilder re-builds all
2871    /// checkpoints that had been built before the most recent restart. You
2872    /// can think of this as a WAL replay operation. Upon startup, we may
2873    /// have a number of consensus commits and resulting checkpoints that
2874    /// were built but not committed to disk. We want to reprocess the
2875    /// commits and rebuild the checkpoints before starting normal operation.
2876    pub async fn spawn(&self, consensus_replay_waiter: Option<ReplayWaiter>) -> JoinSet<()> {
2877        let mut tasks = JoinSet::new();
2878
2879        let (builder, aggregator, state_hasher) = self.state.lock().take_unstarted();
2880        let (builder_finished_tx, builder_finished_rx) = tokio::sync::oneshot::channel();
2881        tasks.spawn(monitored_future!(async move {
2882            builder.run(consensus_replay_waiter).await;
2883            builder_finished_tx.send(()).ok();
2884        }));
2885        tasks.spawn(monitored_future!(aggregator.run()));
2886        tasks.spawn(monitored_future!(state_hasher.run()));
2887
2888        // If this times out, the validator may still start up. The worst that can
2889        // happen is that we will crash later on instead of immediately. The eventual
2890        // crash would occur because we may be missing transactions that are below the
2891        // highest_synced_checkpoint watermark, which can cause a crash in
2892        // `CheckpointExecutor::extract_randomness_rounds`.
2893        if tokio::time::timeout(Duration::from_secs(120), async move {
2894            tokio::select! {
2895                _ = builder_finished_rx => { debug!("CheckpointBuilder finished"); }
2896                _ = self.wait_for_rebuilt_checkpoints() => (),
2897            }
2898        })
2899        .await
2900        .is_err()
2901        {
2902            debug_fatal!("Timed out waiting for checkpoints to be rebuilt");
2903        }
2904
2905        tasks
2906    }
2907}
2908
2909impl CheckpointService {
2910    /// Waits until all checkpoints had been built before the node restarted
2911    /// are rebuilt. This is required to preserve the invariant that all
2912    /// checkpoints (and their transactions) below the
2913    /// highest_synced_checkpoint watermark are available. Once the
2914    /// checkpoints are constructed, we can be sure that the transactions
2915    /// have also been executed.
2916    pub async fn wait_for_rebuilt_checkpoints(&self) {
2917        let highest_previously_built_seq = self.highest_previously_built_seq;
2918        let mut rx = self.highest_currently_built_seq_tx.subscribe();
2919        let mut highest_currently_built_seq = *rx.borrow_and_update();
2920        info!(
2921            "Waiting for checkpoints to be rebuilt, previously built seq: \
2922            {highest_previously_built_seq}, currently built seq: {highest_currently_built_seq}"
2923        );
2924        loop {
2925            if highest_currently_built_seq >= highest_previously_built_seq {
2926                info!("Checkpoint rebuild complete");
2927                break;
2928            }
2929            rx.changed().await.unwrap();
2930            highest_currently_built_seq = *rx.borrow_and_update();
2931        }
2932    }
2933
2934    #[cfg(test)]
2935    fn write_and_notify_checkpoint_for_testing(
2936        &self,
2937        epoch_store: &AuthorityPerEpochStore,
2938        checkpoint: PendingCheckpoint,
2939    ) -> IotaResult {
2940        use crate::authority::authority_per_epoch_store::consensus_quarantine::ConsensusCommitOutput;
2941
2942        let mut output = ConsensusCommitOutput::new(0);
2943        epoch_store.write_pending_checkpoint(&mut output, &checkpoint)?;
2944        output.set_default_commit_stats_for_testing();
2945        epoch_store.push_consensus_output_for_tests(output);
2946        self.notify_checkpoint()?;
2947        Ok(())
2948    }
2949}
2950
2951impl CheckpointServiceNotify for CheckpointService {
2952    fn notify_checkpoint_signature(
2953        &self,
2954        epoch_store: &AuthorityPerEpochStore,
2955        info: &CheckpointSignatureMessage,
2956    ) -> IotaResult {
2957        let sequence = info.summary.sequence_number;
2958        let signer = info.summary.auth_sig().authority.concise();
2959
2960        if let Some(highest_verified_checkpoint) =
2961            self.tables.get_highest_verified_checkpoint_seq_number()?
2962        {
2963            if sequence <= highest_verified_checkpoint {
2964                trace!(
2965                    checkpoint_seq = sequence,
2966                    "Ignore checkpoint signature from {} - already certified", signer,
2967                );
2968                self.metrics
2969                    .last_ignored_checkpoint_signature_received
2970                    .set(sequence as i64);
2971                return Ok(());
2972            }
2973        }
2974        trace!(
2975            checkpoint_seq = sequence,
2976            "Received checkpoint signature, digest {} from {}",
2977            info.summary.digest(),
2978            signer,
2979        );
2980        self.metrics
2981            .last_received_checkpoint_signatures
2982            .with_label_values(&[&signer.to_string()])
2983            .set(sequence as i64);
2984        // While it can be tempting to make last_signature_index into AtomicU64, this
2985        // won't work We need to make sure we write to `pending_signatures` and
2986        // trigger `notify_aggregator` without race conditions
2987        let mut index = self.last_signature_index.lock();
2988        *index += 1;
2989        epoch_store.insert_checkpoint_signature(sequence, *index, info)?;
2990        self.notify_aggregator.notify_one();
2991        Ok(())
2992    }
2993
2994    fn notify_checkpoint(&self) -> IotaResult {
2995        self.notify_builder.notify_one();
2996        Ok(())
2997    }
2998}
2999
3000#[iota_macros::with_checked_arithmetic]
3001mod checked {
3002    use iota_sdk_types::{GasCostSummary, TransactionEffects};
3003    use iota_types::effects::TransactionEffectsAPI;
3004    use itertools::MultiUnzip;
3005
3006    #[expect(clippy::type_complexity)]
3007    pub fn new_gas_cost_summary_from_txn_effects<'a>(
3008        transactions: impl Iterator<Item = &'a TransactionEffects>,
3009    ) -> GasCostSummary {
3010        let (
3011            storage_costs,
3012            computation_costs,
3013            computation_costs_burned,
3014            storage_rebates,
3015            non_refundable_storage_fee,
3016        ): (Vec<u64>, Vec<u64>, Vec<u64>, Vec<u64>, Vec<u64>) = transactions
3017            .map(|e| {
3018                (
3019                    e.gas_cost_summary().storage_cost,
3020                    e.gas_cost_summary().computation_cost,
3021                    e.gas_cost_summary().computation_cost_burned,
3022                    e.gas_cost_summary().storage_rebate,
3023                    e.gas_cost_summary().non_refundable_storage_fee,
3024                )
3025            })
3026            .multiunzip();
3027
3028        GasCostSummary::new(
3029            computation_costs.iter().sum(),
3030            computation_costs_burned.iter().sum(),
3031            storage_costs.iter().sum(),
3032            storage_rebates.iter().sum(),
3033            non_refundable_storage_fee.iter().sum(),
3034        )
3035    }
3036}
3037// test helper
3038pub struct CheckpointServiceNoop {}
3039impl CheckpointServiceNotify for CheckpointServiceNoop {
3040    fn notify_checkpoint_signature(
3041        &self,
3042        _: &AuthorityPerEpochStore,
3043        _: &CheckpointSignatureMessage,
3044    ) -> IotaResult {
3045        Ok(())
3046    }
3047
3048    fn notify_checkpoint(&self) -> IotaResult {
3049        Ok(())
3050    }
3051}
3052
3053pin_project! {
3054    pub struct PollCounter<Fut> {
3055        #[pin]
3056        future: Fut,
3057        count: usize,
3058    }
3059}
3060
3061impl<Fut> PollCounter<Fut> {
3062    pub fn new(future: Fut) -> Self {
3063        Self { future, count: 0 }
3064    }
3065
3066    pub fn count(&self) -> usize {
3067        self.count
3068    }
3069}
3070
3071impl<Fut: Future> Future for PollCounter<Fut> {
3072    type Output = (usize, Fut::Output);
3073
3074    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
3075        let this = self.project();
3076        *this.count += 1;
3077        match this.future.poll(cx) {
3078            Poll::Ready(output) => Poll::Ready((*this.count, output)),
3079            Poll::Pending => Poll::Pending,
3080        }
3081    }
3082}
3083
3084fn poll_count<Fut>(future: Fut) -> PollCounter<Fut> {
3085    PollCounter::new(future)
3086}
3087
3088/// A verified checkpoint over the given contents at the given sequence
3089/// number, with a placeholder signature; usable wherever verification is
3090/// not re-run and no committee is needed.
3091#[cfg(test)]
3092pub(crate) fn test_checkpoint_with_contents(
3093    sequence_number: CheckpointSequenceNumber,
3094    full_contents: &FullCheckpointContents,
3095) -> VerifiedCheckpoint {
3096    let contents = full_contents.checkpoint_contents();
3097    let summary = CheckpointSummary {
3098        epoch: 0,
3099        sequence_number,
3100        network_total_transactions: full_contents.size() as u64,
3101        contents_digest: contents.digest(),
3102        previous_digest: None,
3103        epoch_rolling_gas_cost_summary: GasCostSummary::default(),
3104        end_of_epoch_data: None,
3105        timestamp_ms: 0,
3106        version_specific_data: Vec::new(),
3107        checkpoint_commitments: Vec::new(),
3108    };
3109    let sig = AuthorityStrongQuorumSignInfo {
3110        epoch: 0,
3111        signature: Default::default(),
3112        signers_map: Default::default(),
3113    };
3114    VerifiedCheckpoint::new_unchecked(
3115        iota_types::message_envelope::Envelope::new_from_data_and_sig(summary, sig),
3116    )
3117}
3118
3119#[cfg(test)]
3120mod tests {
3121    use std::{
3122        collections::{BTreeMap, HashMap},
3123        ops::Deref,
3124    };
3125
3126    use futures::{FutureExt as _, future::BoxFuture};
3127    use iota_macros::sim_test;
3128    use iota_protocol_config::{Chain, ProtocolConfig, ProtocolVersion};
3129    use iota_sdk_types::{
3130        GenesisObject, Identifier, MovePackage, ObjectData, ObjectId, Owner, TransactionEffects,
3131        TransactionEffectsDigest, TransactionEvents, Version,
3132    };
3133    use iota_types::{
3134        effects::{TransactionEffectsAPIForTesting, TransactionEffectsExtForTesting},
3135        messages_checkpoint::SignedCheckpointSummary,
3136        transaction::VerifiedTransaction,
3137    };
3138    use tokio::sync::mpsc;
3139
3140    use super::*;
3141    use crate::authority::test_authority_builder::TestAuthorityBuilder;
3142
3143    #[tokio::test]
3144    async fn insert_verified_checkpoint_contents_persists_digests_and_caches_full_contents() {
3145        let tempdir = iota_common::tempdir();
3146        let path = tempdir.path();
3147
3148        let full_contents = FullCheckpointContents::random_for_testing();
3149        let checkpoint = test_checkpoint_with_contents(0, &full_contents);
3150        let contents_digest = checkpoint.contents_digest;
3151
3152        {
3153            let store = CheckpointStore::new(path);
3154            store
3155                .insert_verified_checkpoint_contents(
3156                    &checkpoint,
3157                    VerifiedCheckpointContents::new_unchecked(full_contents.clone()),
3158                )
3159                .unwrap();
3160
3161            // Full contents are served from the in-memory cache.
3162            assert_eq!(
3163                store
3164                    .get_full_checkpoint_contents_by_sequence_number(0)
3165                    .unwrap()
3166                    .as_ref(),
3167                &full_contents
3168            );
3169            assert_eq!(
3170                store
3171                    .get_full_checkpoint_contents_by_digest(&contents_digest)
3172                    .unwrap()
3173                    .as_ref(),
3174                &full_contents
3175            );
3176            // The digest-form contents are durable.
3177            assert_eq!(
3178                store
3179                    .get_checkpoint_contents(&contents_digest)
3180                    .unwrap()
3181                    .map(|c| c.digest()),
3182                Some(contents_digest)
3183            );
3184        }
3185
3186        // Reopening drops the in-memory cache, but the digest-form contents
3187        // survive for readers to reconstruct full contents from.
3188        let store = CheckpointStore::new(path);
3189        assert!(
3190            store
3191                .get_full_checkpoint_contents_by_sequence_number(0)
3192                .is_none()
3193        );
3194        assert!(
3195            store
3196                .get_full_checkpoint_contents_by_digest(&contents_digest)
3197                .is_none()
3198        );
3199        assert!(
3200            store
3201                .get_checkpoint_contents(&contents_digest)
3202                .unwrap()
3203                .is_some()
3204        );
3205    }
3206
3207    #[tokio::test]
3208    async fn open_restores_empty_checkpoint_contents() {
3209        let tempdir = iota_common::tempdir();
3210        let path = tempdir.path();
3211
3212        {
3213            let store = CheckpointStore::new(path);
3214            assert_eq!(
3215                store
3216                    .get_checkpoint_contents(&EMPTY_CHECKPOINT_CONTENTS_DIGEST)
3217                    .unwrap()
3218                    .map(|c| c.digest()),
3219                Some(*EMPTY_CHECKPOINT_CONTENTS_DIGEST)
3220            );
3221            store
3222                .tables
3223                .checkpoint_content
3224                .remove(&EMPTY_CHECKPOINT_CONTENTS_DIGEST)
3225                .unwrap();
3226        }
3227
3228        let store = CheckpointStore::new(path);
3229        assert!(
3230            store
3231                .get_checkpoint_contents(&EMPTY_CHECKPOINT_CONTENTS_DIGEST)
3232                .unwrap()
3233                .is_some_and(|c| c.is_empty())
3234        );
3235    }
3236
3237    #[tokio::test]
3238    async fn cache_full_checkpoint_contents_serves_reads_without_disk_writes() {
3239        let store = CheckpointStore::new_for_tests();
3240        let full_contents = FullCheckpointContents::random_for_testing();
3241        let checkpoint = test_checkpoint_with_contents(0, &full_contents);
3242        let contents_digest = checkpoint.contents_digest;
3243
3244        store.cache_full_checkpoint_contents(
3245            checkpoint.sequence_number(),
3246            contents_digest,
3247            full_contents.clone(),
3248        );
3249
3250        assert_eq!(
3251            store
3252                .get_full_checkpoint_contents_by_sequence_number(0)
3253                .unwrap()
3254                .as_ref(),
3255            &full_contents
3256        );
3257        assert_eq!(
3258            store
3259                .get_full_checkpoint_contents_by_digest(&contents_digest)
3260                .unwrap()
3261                .as_ref(),
3262            &full_contents
3263        );
3264        // Cache-only: the digest-form contents row is untouched.
3265        assert!(
3266            store
3267                .get_checkpoint_contents(&contents_digest)
3268                .unwrap()
3269                .is_none()
3270        );
3271    }
3272
3273    // State-sync skips its own durable `checkpoint_content` write when the
3274    // full contents are already available, and the checkpoint executor panics
3275    // on a missing row. So the builder must not serve contents from the cache
3276    // before the row is durable.
3277    #[tokio::test]
3278    async fn builder_caches_full_contents_only_after_durable_contents_write() {
3279        let state = TestAuthorityBuilder::new().build().await;
3280
3281        let tx = VerifiedTransaction::new_genesis_transaction(vec![], vec![]);
3282        let digest = *tx.digest();
3283        state
3284            .database_for_testing()
3285            .perpetual_tables
3286            .transactions
3287            .insert(&digest, tx.serializable_ref())
3288            .unwrap();
3289
3290        let mut effects_map = HashMap::new();
3291        commit_cert_for_test(
3292            &mut effects_map,
3293            state.clone(),
3294            digest,
3295            vec![],
3296            GasCostSummary::new(1, 1, 1, 1, 1),
3297        );
3298        let effects = effects_map[&digest].clone();
3299
3300        let signature = iota_types::crypto::zero_ed25519_signature().into();
3301        state
3302            .epoch_store_for_testing()
3303            .test_insert_user_signature(digest, vec![signature]);
3304
3305        let (output, _result) = mpsc::channel::<(CheckpointContents, CheckpointSummary)>(10);
3306        let (certified_output, _certified_result) = mpsc::channel::<CertifiedCheckpointSummary>(10);
3307
3308        let tmp_dir = iota_common::tempdir();
3309        let checkpoint_store = CheckpointStore::new(tmp_dir.path());
3310        let epoch_store = state.epoch_store_for_testing();
3311
3312        let global_state_hasher = Arc::new(GlobalStateHasher::new_for_tests(
3313            state.get_global_state_hash_store().clone(),
3314        ));
3315
3316        let checkpoint_service = CheckpointService::build(
3317            state.clone(),
3318            checkpoint_store.clone(),
3319            epoch_store.clone(),
3320            Arc::new(effects_map),
3321            Arc::downgrade(&global_state_hasher),
3322            Box::new(output),
3323            Box::new(certified_output),
3324            CheckpointMetrics::new_for_tests(),
3325            3,
3326            100_000,
3327        );
3328        // Drive the builder manually instead of spawning it, to observe the
3329        // store between checkpoint creation and the durable batch write. The
3330        // state hasher must stay alive: `create_checkpoints` sends to it.
3331        let (builder, _aggregator, _hasher) = checkpoint_service.state.lock().take_unstarted();
3332
3333        checkpoint_service
3334            .write_and_notify_checkpoint_for_testing(&epoch_store, p(0, vec![digest], 0))
3335            .unwrap();
3336
3337        let details = PendingCheckpointInfo {
3338            timestamp_ms: 0,
3339            last_of_epoch: false,
3340            checkpoint_height: 0,
3341        };
3342        let new_checkpoints = builder
3343            .create_checkpoints(vec![effects], &details, &HashSet::from([digest]))
3344            .await
3345            .unwrap();
3346        let summary = new_checkpoints.first().summary.clone();
3347
3348        assert!(
3349            checkpoint_store
3350                .get_full_checkpoint_contents_by_sequence_number(summary.sequence_number)
3351                .is_none(),
3352            "full contents must not be served before the checkpoint_content row is durable"
3353        );
3354        assert!(
3355            checkpoint_store
3356                .get_checkpoint_contents(&summary.contents_digest)
3357                .unwrap()
3358                .is_none()
3359        );
3360
3361        builder
3362            .write_checkpoints(details.checkpoint_height, new_checkpoints)
3363            .await
3364            .unwrap();
3365
3366        assert!(
3367            checkpoint_store
3368                .get_checkpoint_contents(&summary.contents_digest)
3369                .unwrap()
3370                .is_some()
3371        );
3372        let cached = checkpoint_store
3373            .get_full_checkpoint_contents_by_sequence_number(summary.sequence_number)
3374            .expect("builder should cache full contents once the row is durable");
3375        assert_eq!(
3376            cached.checkpoint_contents().digest(),
3377            summary.contents_digest
3378        );
3379    }
3380
3381    #[sim_test]
3382    pub async fn checkpoint_builder_test() {
3383        telemetry_subscribers::init_for_testing();
3384
3385        let mut protocol_config =
3386            ProtocolConfig::get_for_version(ProtocolVersion::max(), Chain::Unknown);
3387        protocol_config.set_min_checkpoint_interval_ms_for_testing(100);
3388        // This test exercises the strict adjacent-checkpoint interval.
3389        protocol_config.disable_checkpoint_rate_window_size_for_testing();
3390        let state = TestAuthorityBuilder::new()
3391            .with_protocol_config(protocol_config)
3392            .build()
3393            .await;
3394
3395        // Build distinct genesis transactions and assign their digests to
3396        // indices so that, within any pending checkpoint, digest order matches
3397        // index order: `CausalOrder` orders non-dependent transactions by
3398        // digest, and the assertions below rely on that order. Transactions
3399        // 15..20 carry a large payload to exercise size-based checkpoint
3400        // splitting.
3401        let make_tx = |seed: u8, payload_size: usize| {
3402            let mut id = [0u8; 32];
3403            id[0] = seed;
3404            VerifiedTransaction::new_genesis_transaction(
3405                vec![GenesisObject::new(
3406                    ObjectData::Package(MovePackage::new(
3407                        ObjectId::new(id),
3408                        Version::default(),
3409                        BTreeMap::from([(Identifier::new_unchecked("m"), vec![0u8; payload_size])]),
3410                        // no modules so empty type_origin_table as no types are defined in
3411                        // this package
3412                        Vec::new(),
3413                        // no modules so empty linkage_table as no dependencies of this package
3414                        // exist
3415                        BTreeMap::new(),
3416                    )),
3417                    Owner::Immutable,
3418                )],
3419                vec![],
3420            )
3421        };
3422
3423        let mut small: Vec<_> = (0..15).map(|seed| make_tx(seed, 1)).collect();
3424        small.sort_by_key(|tx| *tx.digest());
3425
3426        let mut large: Vec<_> = (0..5).map(|seed| make_tx(100 + seed, 40000)).collect();
3427        large.sort_by_key(|tx| *tx.digest());
3428
3429        let txns: Vec<_> = small.into_iter().chain(large).collect();
3430        let digests: Vec<TransactionDigest> = txns.iter().map(|tx| *tx.digest()).collect();
3431
3432        // Digest for test index `i`; ascending within the small (0..15) and
3433        // large (15..20) pools, so index order implies digest order per pool.
3434        let d = |i: u8| digests[i as usize];
3435
3436        for (tx, digest) in txns.iter().zip(&digests) {
3437            state
3438                .database_for_testing()
3439                .perpetual_tables
3440                .transactions
3441                .insert(digest, tx.serializable_ref())
3442                .unwrap();
3443        }
3444
3445        let mut store = HashMap::<TransactionDigest, TransactionEffects>::new();
3446        commit_cert_for_test(
3447            &mut store,
3448            state.clone(),
3449            d(1),
3450            vec![d(2), d(3)],
3451            GasCostSummary::new(11, 11, 12, 11, 1),
3452        );
3453        commit_cert_for_test(
3454            &mut store,
3455            state.clone(),
3456            d(2),
3457            vec![d(3), d(4)],
3458            GasCostSummary::new(21, 21, 22, 21, 1),
3459        );
3460        commit_cert_for_test(
3461            &mut store,
3462            state.clone(),
3463            d(3),
3464            vec![],
3465            GasCostSummary::new(31, 31, 32, 31, 1),
3466        );
3467        commit_cert_for_test(
3468            &mut store,
3469            state.clone(),
3470            d(4),
3471            vec![],
3472            GasCostSummary::new(41, 41, 42, 41, 1),
3473        );
3474        for i in [5, 6, 7, 10, 11, 12, 13] {
3475            commit_cert_for_test(
3476                &mut store,
3477                state.clone(),
3478                d(i),
3479                vec![],
3480                GasCostSummary::new(41, 41, 42, 41, 1),
3481            );
3482        }
3483        for i in [15, 16, 17] {
3484            commit_cert_for_test(
3485                &mut store,
3486                state.clone(),
3487                d(i),
3488                vec![],
3489                GasCostSummary::new(51, 51, 52, 51, 1),
3490            );
3491        }
3492        let all_digests: Vec<_> = store.keys().copied().collect();
3493        for digest in all_digests {
3494            let signature = iota_types::crypto::zero_ed25519_signature().into();
3495            state
3496                .epoch_store_for_testing()
3497                .test_insert_user_signature(digest, vec![signature]);
3498        }
3499
3500        let (output, mut result) = mpsc::channel::<(CheckpointContents, CheckpointSummary)>(10);
3501        let (certified_output, mut certified_result) =
3502            mpsc::channel::<CertifiedCheckpointSummary>(10);
3503        let store = Arc::new(store);
3504
3505        let tmp_dir = iota_common::tempdir();
3506        let checkpoint_store = CheckpointStore::new(tmp_dir.path());
3507        let epoch_store = state.epoch_store_for_testing();
3508
3509        let global_state_hasher = Arc::new(GlobalStateHasher::new_for_tests(
3510            state.get_global_state_hash_store().clone(),
3511        ));
3512
3513        let checkpoint_service = CheckpointService::build(
3514            state.clone(),
3515            checkpoint_store.clone(),
3516            epoch_store.clone(),
3517            store,
3518            Arc::downgrade(&global_state_hasher),
3519            Box::new(output),
3520            Box::new(certified_output),
3521            CheckpointMetrics::new_for_tests(),
3522            3,
3523            100_000,
3524        );
3525        let _tasks = checkpoint_service.spawn(None).await;
3526
3527        checkpoint_service
3528            .write_and_notify_checkpoint_for_testing(&epoch_store, p(0, vec![d(4)], 0))
3529            .unwrap();
3530        checkpoint_service
3531            .write_and_notify_checkpoint_for_testing(&epoch_store, p(1, vec![d(1), d(3)], 2000))
3532            .unwrap();
3533        checkpoint_service
3534            .write_and_notify_checkpoint_for_testing(
3535                &epoch_store,
3536                p(2, vec![d(10), d(11), d(12), d(13)], 3000),
3537            )
3538            .unwrap();
3539        checkpoint_service
3540            .write_and_notify_checkpoint_for_testing(
3541                &epoch_store,
3542                p(3, vec![d(15), d(16), d(17)], 4000),
3543            )
3544            .unwrap();
3545        checkpoint_service
3546            .write_and_notify_checkpoint_for_testing(&epoch_store, p(4, vec![d(5)], 4001))
3547            .unwrap();
3548        checkpoint_service
3549            .write_and_notify_checkpoint_for_testing(&epoch_store, p(5, vec![d(6)], 5000))
3550            .unwrap();
3551
3552        let (c1c, c1s) = result.recv().await.unwrap();
3553        let (c2c, c2s) = result.recv().await.unwrap();
3554
3555        let c1t = c1c.iter().map(|d| d.transaction).collect::<Vec<_>>();
3556        let c2t = c2c.iter().map(|d| d.transaction).collect::<Vec<_>>();
3557        assert_eq!(c1t, vec![d(4)]);
3558        assert_eq!(c1s.previous_digest, None);
3559        assert_eq!(c1s.sequence_number, 0);
3560        assert_eq!(
3561            c1s.epoch_rolling_gas_cost_summary,
3562            GasCostSummary::new(41, 41, 42, 41, 1)
3563        );
3564
3565        assert_eq!(c2t, vec![d(3), d(2), d(1)]);
3566        assert_eq!(c2s.previous_digest, Some(c1s.digest()));
3567        assert_eq!(c2s.sequence_number, 1);
3568        assert_eq!(
3569            c2s.epoch_rolling_gas_cost_summary,
3570            GasCostSummary::new(104, 104, 108, 104, 4)
3571        );
3572
3573        // The builder caches the full contents of each locally built checkpoint.
3574        // The cached contents must match the checkpoint's transactions and hash
3575        // to its content digest, so state-sync peers can be served from memory.
3576        for (summary, digest_contents) in [(&c1s, &c1c), (&c2s, &c2c)] {
3577            let cached = checkpoint_store
3578                .get_full_checkpoint_contents_by_sequence_number(summary.sequence_number)
3579                .expect("builder should cache full contents of a locally built checkpoint");
3580            assert_eq!(&cached.checkpoint_contents(), digest_contents);
3581            assert_eq!(
3582                cached.checkpoint_contents().digest(),
3583                summary.contents_digest
3584            );
3585            // A cached entry must imply a durable checkpoint_content row:
3586            // state-sync skips its own durable write when contents are
3587            // already available.
3588            assert!(
3589                checkpoint_store
3590                    .get_checkpoint_contents(&summary.contents_digest)
3591                    .unwrap()
3592                    .is_some()
3593            );
3594        }
3595
3596        // Pending at index 2 had 4 transactions, and we configured 3 transactions max.
3597        // Verify that we split into 2 checkpoints.
3598        let (c3c, c3s) = result.recv().await.unwrap();
3599        let c3t = c3c.iter().map(|d| d.transaction).collect::<Vec<_>>();
3600        let (c4c, c4s) = result.recv().await.unwrap();
3601        let c4t = c4c.iter().map(|d| d.transaction).collect::<Vec<_>>();
3602        assert_eq!(c3s.sequence_number, 2);
3603        assert_eq!(c3s.previous_digest, Some(c2s.digest()));
3604        assert_eq!(c4s.sequence_number, 3);
3605        assert_eq!(c4s.previous_digest, Some(c3s.digest()));
3606        assert_eq!(c3t, vec![d(10), d(11), d(12)]);
3607        assert_eq!(c4t, vec![d(13)]);
3608
3609        // Pending at index 3 had 3 transactions of 40K size, and we configured 100K
3610        // max. Verify that we split into 2 checkpoints.
3611        let (c5c, c5s) = result.recv().await.unwrap();
3612        let c5t = c5c.iter().map(|d| d.transaction).collect::<Vec<_>>();
3613        let (c6c, c6s) = result.recv().await.unwrap();
3614        let c6t = c6c.iter().map(|d| d.transaction).collect::<Vec<_>>();
3615        assert_eq!(c5s.sequence_number, 4);
3616        assert_eq!(c5s.previous_digest, Some(c4s.digest()));
3617        assert_eq!(c6s.sequence_number, 5);
3618        assert_eq!(c6s.previous_digest, Some(c5s.digest()));
3619        assert_eq!(c5t, vec![d(15), d(16)]);
3620        assert_eq!(c6t, vec![d(17)]);
3621
3622        // Pending at index 4 was too soon after the prior one and should be coalesced
3623        // into the next one.
3624        let (c7c, c7s) = result.recv().await.unwrap();
3625        let c7t = c7c.iter().map(|d| d.transaction).collect::<Vec<_>>();
3626        assert_eq!(c7t, vec![d(5), d(6)]);
3627        assert_eq!(c7s.previous_digest, Some(c6s.digest()));
3628        assert_eq!(c7s.sequence_number, 6);
3629
3630        let c1ss = SignedCheckpointSummary::new(c1s.epoch, c1s, state.secret.deref(), state.name);
3631        let c2ss = SignedCheckpointSummary::new(c2s.epoch, c2s, state.secret.deref(), state.name);
3632
3633        checkpoint_service
3634            .notify_checkpoint_signature(
3635                &epoch_store,
3636                &CheckpointSignatureMessage { summary: c2ss },
3637            )
3638            .unwrap();
3639        checkpoint_service
3640            .notify_checkpoint_signature(
3641                &epoch_store,
3642                &CheckpointSignatureMessage { summary: c1ss },
3643            )
3644            .unwrap();
3645
3646        let c1sc = certified_result.recv().await.unwrap();
3647        let c2sc = certified_result.recv().await.unwrap();
3648        assert_eq!(c1sc.sequence_number, 0);
3649        assert_eq!(c2sc.sequence_number, 1);
3650    }
3651
3652    #[sim_test]
3653    pub async fn checkpoint_builder_windowed_interval_test() {
3654        telemetry_subscribers::init_for_testing();
3655
3656        // 100ms interval with a window of 3 checkpoints: a checkpoint is built
3657        // when either 100ms passed since the previous one or the checkpoint 3
3658        // back is >= 300ms older. The windowed arm permits sub-interval bursts
3659        // while banked budget lasts; the adjacent arm caps every quiet gap at
3660        // one interval regardless of window state.
3661        let mut protocol_config =
3662            ProtocolConfig::get_for_version(ProtocolVersion::max(), Chain::Unknown);
3663        protocol_config.set_min_checkpoint_interval_ms_for_testing(100);
3664        protocol_config.set_checkpoint_rate_window_size_for_testing(3);
3665        let state = TestAuthorityBuilder::new()
3666            .with_protocol_config(protocol_config)
3667            .build()
3668            .await;
3669
3670        // Distinct real transactions: only build timing is under test, but the
3671        // builder pairs each stored transaction with its effects by digest, so
3672        // effects must reference the digests of actually stored transactions.
3673        let make_tx = |seed: u8| {
3674            let mut id = [0u8; 32];
3675            id[0] = seed;
3676            VerifiedTransaction::new_genesis_transaction(
3677                vec![GenesisObject::new(
3678                    ObjectData::Package(MovePackage::new(
3679                        ObjectId::new(id),
3680                        Version::default(),
3681                        BTreeMap::from([(Identifier::new_unchecked("m"), vec![0u8; 1])]),
3682                        // no modules so empty type_origin_table as no types are defined in
3683                        // this package
3684                        Vec::new(),
3685                        // no modules so empty linkage_table as no dependencies of this package
3686                        // exist
3687                        BTreeMap::new(),
3688                    )),
3689                    Owner::Immutable,
3690                )],
3691                vec![],
3692            )
3693        };
3694        let txns: Vec<_> = (1..=8).map(make_tx).collect();
3695        let digests: Vec<TransactionDigest> = txns.iter().map(|tx| *tx.digest()).collect();
3696        // Digest for test index `i` (1-based).
3697        let d = |i: u8| digests[(i - 1) as usize];
3698
3699        for tx in &txns {
3700            state
3701                .database_for_testing()
3702                .perpetual_tables
3703                .transactions
3704                .insert(tx.digest(), tx.serializable_ref())
3705                .unwrap();
3706        }
3707
3708        let mut store = HashMap::<TransactionDigest, TransactionEffects>::new();
3709        for i in 1..=8 {
3710            commit_cert_for_test(
3711                &mut store,
3712                state.clone(),
3713                d(i),
3714                vec![],
3715                GasCostSummary::new(11, 11, 12, 11, 1),
3716            );
3717        }
3718        let all_digests: Vec<_> = store.keys().copied().collect();
3719        for digest in all_digests {
3720            let signature = iota_types::crypto::zero_ed25519_signature().into();
3721            state
3722                .epoch_store_for_testing()
3723                .test_insert_user_signature(digest, vec![signature]);
3724        }
3725
3726        let (output, mut result) = mpsc::channel::<(CheckpointContents, CheckpointSummary)>(10);
3727        let (certified_output, _certified_result) = mpsc::channel::<CertifiedCheckpointSummary>(10);
3728        let store = Arc::new(store);
3729
3730        let tmp_dir = iota_common::tempdir();
3731        let checkpoint_store = CheckpointStore::new(tmp_dir.path());
3732        let epoch_store = state.epoch_store_for_testing();
3733
3734        let global_state_hasher = Arc::new(GlobalStateHasher::new_for_tests(
3735            state.get_global_state_hash_store().clone(),
3736        ));
3737
3738        let checkpoint_service = CheckpointService::build(
3739            state.clone(),
3740            checkpoint_store,
3741            epoch_store.clone(),
3742            store,
3743            Arc::downgrade(&global_state_hasher),
3744            Box::new(output),
3745            Box::new(certified_output),
3746            CheckpointMetrics::new_for_tests(),
3747            3,
3748            100_000,
3749        );
3750        let _tasks = checkpoint_service.spawn(None).await;
3751
3752        // The windowed arm is inert until 3 checkpoints exist; the first three
3753        // build via the adjacent arm, which their spread-out timestamps
3754        // (0, 1000, 2000) satisfy while accruing burst budget.
3755        // window_ms = 3 * 100 = 300.
3756        for (height, root, timestamp_ms) in [(0, 1u8, 0), (1, 2, 1000), (2, 3, 2000)] {
3757            checkpoint_service
3758                .write_and_notify_checkpoint_for_testing(
3759                    &epoch_store,
3760                    p(height, vec![d(root)], timestamp_ms),
3761                )
3762                .unwrap();
3763        }
3764        // p3 @ 2010: adjacent red (< C2 + 100 = 2100) but windowed green
3765        // (>= C0(0) + 300) -> builds, only 10ms after C2. A strict adjacent
3766        // rule would instead coalesce it.
3767        checkpoint_service
3768            .write_and_notify_checkpoint_for_testing(&epoch_store, p(3, vec![d(4)], 2010))
3769            .unwrap();
3770        // p4 @ 2020: windowed green (>= C1(1000) + 300) -> builds, another
3771        // sub-interval burst.
3772        checkpoint_service
3773            .write_and_notify_checkpoint_for_testing(&epoch_store, p(4, vec![d(5)], 2020))
3774            .unwrap();
3775        // Budget now spent: the window covers the tight 2000-2020 cluster.
3776        // p5 @ 2030: adjacent red (< C4 + 100 = 2120), windowed red
3777        // (< C2(2000) + 300 = 2300) -> coalesced.
3778        checkpoint_service
3779            .write_and_notify_checkpoint_for_testing(&epoch_store, p(5, vec![d(6)], 2030))
3780            .unwrap();
3781        // p6 @ 2130: windowed still red (< 2300), but adjacent green
3782        // (>= 2120) -> builds, coalescing p5 and p6. A window-only gate would
3783        // keep coalescing until 2300; the adjacent arm caps the gap.
3784        checkpoint_service
3785            .write_and_notify_checkpoint_for_testing(&epoch_store, p(6, vec![d(7)], 2130))
3786            .unwrap();
3787        // p7 @ 2300: windowed red (< C3(2010) + 300 = 2310), adjacent green
3788        // (>= C5 + 100 = 2230) -> builds.
3789        checkpoint_service
3790            .write_and_notify_checkpoint_for_testing(&epoch_store, p(7, vec![d(8)], 2300))
3791            .unwrap();
3792
3793        let mut built = Vec::new();
3794        for _ in 0..7 {
3795            let (contents, summary) = result.recv().await.unwrap();
3796            built.push((summary.sequence_number, contents.iter().count()));
3797        }
3798
3799        // Seven checkpoints, contiguous sequence numbers.
3800        let sequence_numbers: Vec<_> = built.iter().map(|(seq, _)| *seq).collect();
3801        assert_eq!(sequence_numbers, vec![0, 1, 2, 3, 4, 5, 6]);
3802        // Bursts (seq 3, 4) stand alone despite sub-interval spacing; once the
3803        // budget is spent, pendings coalesce only until the adjacent interval
3804        // elapses (seq 5), never longer. A strict adjacent rule would build
3805        // [0,1,2] and then coalesce the whole 2010-2130 cluster into one
3806        // checkpoint; a window-only rule would coalesce p5-p7 into one
3807        // checkpoint at 2300.
3808        let sizes: Vec<_> = built.iter().map(|(_, size)| *size).collect();
3809        assert_eq!(sizes, vec![1, 1, 1, 1, 1, 2, 1]);
3810    }
3811
3812    impl TransactionCacheRead for HashMap<TransactionDigest, TransactionEffects> {
3813        fn try_notify_read_executed_effects(
3814            &self,
3815            _: &str,
3816            digests: &[TransactionDigest],
3817        ) -> BoxFuture<'_, IotaResult<Vec<TransactionEffects>>> {
3818            std::future::ready(Ok(digests
3819                .iter()
3820                .map(|d| self.get(d).expect("effects not found").clone())
3821                .collect()))
3822            .boxed()
3823        }
3824
3825        fn try_notify_read_executed_effects_digests(
3826            &self,
3827            _: &str,
3828            digests: &[TransactionDigest],
3829        ) -> BoxFuture<'_, IotaResult<Vec<TransactionEffectsDigest>>> {
3830            std::future::ready(Ok(digests
3831                .iter()
3832                .map(|d| {
3833                    self.get(d)
3834                        .map(|fx| fx.digest())
3835                        .expect("effects not found")
3836                })
3837                .collect()))
3838            .boxed()
3839        }
3840
3841        fn try_multi_get_executed_effects(
3842            &self,
3843            digests: &[TransactionDigest],
3844        ) -> IotaResult<Vec<Option<TransactionEffects>>> {
3845            Ok(digests.iter().map(|d| self.get(d).cloned()).collect())
3846        }
3847
3848        // Unimplemented methods - its unfortunate to have this big blob of useless
3849        // code, but it wasn't worth it to keep EffectsNotifyRead around just
3850        // for these tests, as it caused a ton of complication in non-test code.
3851        // (e.g. had to implement EFfectsNotifyRead for all ExecutionCacheRead
3852        // implementors).
3853
3854        fn try_multi_get_transaction_blocks(
3855            &self,
3856            _: &[TransactionDigest],
3857        ) -> IotaResult<Vec<Option<Arc<VerifiedTransaction>>>> {
3858            unimplemented!()
3859        }
3860
3861        fn try_multi_get_executed_effects_digests(
3862            &self,
3863            _: &[TransactionDigest],
3864        ) -> IotaResult<Vec<Option<TransactionEffectsDigest>>> {
3865            unimplemented!()
3866        }
3867
3868        fn try_multi_get_effects(
3869            &self,
3870            _: &[TransactionEffectsDigest],
3871        ) -> IotaResult<Vec<Option<TransactionEffects>>> {
3872            unimplemented!()
3873        }
3874
3875        fn try_multi_get_events(
3876            &self,
3877            _: &[TransactionDigest],
3878        ) -> IotaResult<Vec<Option<TransactionEvents>>> {
3879            unimplemented!()
3880        }
3881    }
3882
3883    #[async_trait::async_trait]
3884    impl CheckpointOutput for mpsc::Sender<(CheckpointContents, CheckpointSummary)> {
3885        async fn checkpoint_created(
3886            &self,
3887            summary: &CheckpointSummary,
3888            contents: &CheckpointContents,
3889            _epoch_store: &Arc<AuthorityPerEpochStore>,
3890            _checkpoint_store: &Arc<CheckpointStore>,
3891        ) -> IotaResult {
3892            self.try_send((contents.clone(), summary.clone())).unwrap();
3893            Ok(())
3894        }
3895    }
3896
3897    #[async_trait::async_trait]
3898    impl CertifiedCheckpointOutput for mpsc::Sender<CertifiedCheckpointSummary> {
3899        async fn certified_checkpoint_created(
3900            &self,
3901            summary: &CertifiedCheckpointSummary,
3902        ) -> IotaResult {
3903            self.try_send(summary.clone()).unwrap();
3904            Ok(())
3905        }
3906    }
3907
3908    fn p(i: u64, roots: Vec<TransactionDigest>, timestamp_ms: u64) -> PendingCheckpoint {
3909        PendingCheckpoint::V1(PendingCheckpointContentsV1 {
3910            roots: roots.into_iter().map(TransactionKey::Digest).collect(),
3911            details: PendingCheckpointInfo {
3912                timestamp_ms,
3913                last_of_epoch: false,
3914                checkpoint_height: i,
3915            },
3916        })
3917    }
3918
3919    fn e(
3920        transaction_digest: TransactionDigest,
3921        dependencies: Vec<TransactionDigest>,
3922        gas_cost_summary: GasCostSummary,
3923    ) -> TransactionEffects {
3924        let mut effects = TransactionEffects::new_empty_v1_for_testing(transaction_digest);
3925        *effects.dependencies_mut_for_testing() = dependencies;
3926        *effects.gas_cost_summary_mut_for_testing() = gas_cost_summary;
3927        effects
3928    }
3929
3930    fn commit_cert_for_test(
3931        store: &mut HashMap<TransactionDigest, TransactionEffects>,
3932        state: Arc<AuthorityState>,
3933        digest: TransactionDigest,
3934        dependencies: Vec<TransactionDigest>,
3935        gas_cost_summary: GasCostSummary,
3936    ) {
3937        let epoch_store = state.epoch_store_for_testing();
3938        let effects = e(digest, dependencies, gas_cost_summary);
3939        store.insert(digest, effects);
3940        epoch_store.insert_executed_in_epoch(&digest);
3941    }
3942}