1mod 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
106pub(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 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 pub checkpoint_height: Option<CheckpointHeight>,
172 pub position_in_commit: usize,
173}
174
175#[derive(DBMapUtils)]
176pub struct CheckpointStoreTables {
177 pub(crate) checkpoint_content: DBMap<CheckpointContentsDigest, CheckpointContents>,
179
180 #[allow(dead_code)]
184 #[deprecated_db_map]
185 checkpoint_sequence_by_contents_digest: Option<DBMap<(), ()>>,
186
187 #[allow(dead_code)]
193 #[deprecated_db_map]
194 full_checkpoint_content: Option<DBMap<(), ()>>,
195
196 pub(crate) certified_checkpoints: DBMap<CheckpointSequenceNumber, TrustedCheckpoint>,
198 pub(crate) checkpoint_by_digest: DBMap<CheckpointDigest, TrustedCheckpoint>,
200
201 pub(crate) locally_computed_checkpoints: DBMap<CheckpointSequenceNumber, CheckpointSummary>,
205
206 epoch_last_checkpoint_map: DBMap<EpochId, CheckpointSequenceNumber>,
213
214 epoch_info: DBMap<EpochId, EpochInfoV2>,
226
227 epoch_info_watermark: DBMap<(), EpochId>,
233
234 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 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 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 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 fn get_watermark(
434 &self,
435 watermark: CheckpointWatermark,
436 ) -> Result<Option<(CheckpointSequenceNumber, CheckpointDigest)>, TypedStoreError> {
437 self.tables.watermarks.get(&watermark)
438 }
439
440 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 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 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 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 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 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 #[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 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 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 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 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 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 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 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 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 pub fn should_cache_full_checkpoint_contents(&self, seq: CheckpointSequenceNumber) -> bool {
948 self.full_checkpoint_contents_cache.should_cache(seq)
949 }
950
951 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 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 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
1112struct BuiltCheckpoint {
1115 summary: CheckpointSummary,
1116 contents: CheckpointContents,
1117 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
1165pub struct CheckpointSignatureAggregator {
1167 next_index: u64,
1168 summary: CheckpointSummary,
1169 digest: CheckpointDigest,
1170 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 global_state_hasher: Weak<GlobalStateHasher>,
1186 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 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 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 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 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 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 let can_build = interval_elapsed
1316 || checkpoints_iter
1319 .peek()
1320 .is_some_and(|(_, next_pending)| next_pending.details().last_of_epoch)
1321 || 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 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 self.metrics
1355 .commits_per_checkpoint
1356 .observe(commits_in_checkpoint as f64);
1357 last_seq = Some(seq);
1360 self.last_built.send_if_modified(|cur| {
1361 if seq > *cur {
1363 *cur = seq;
1364 true
1365 } else {
1366 false
1367 }
1368 });
1369
1370 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 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 #[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 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 let consensus_commit_prologue =
1475 self.extract_consensus_commit_prologue(&root_digests, &root_effects)?;
1476
1477 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 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 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 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 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 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 #[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 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 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 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 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 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 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 chunks.push(chunk);
1724 }
1730 Ok(chunks)
1731 }
1732
1733 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 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 }
1819 TransactionKind::RandomnessStateUpdate(rsu) => {
1820 randomness_rounds
1821 .insert(*effects.transaction_digest(), rsu.randomness_round);
1822 }
1823 _ => {
1824 let digest = *effects.transaction_digest();
1829 if !all_roots.contains(&digest) {
1830 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 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 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 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 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 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 .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 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 #[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 #[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 seen.insert(*digest);
2155
2156 if existing_tx_digests_in_checkpoint.contains(effect.transaction_digest()) {
2158 continue;
2159 }
2160
2161 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 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 #[cfg(msim)]
2213 fn expensive_consensus_commit_prologue_invariants_check(
2214 &self,
2215 root_digests: &[TransactionDigest],
2216 sorted: &[TransactionEffects],
2217 ) {
2218 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 assert!(ccps.len() <= 1);
2238
2239 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 for tx in txs.iter().flatten() {
2255 assert!(!matches!(
2256 tx.transaction().kind(),
2257 TransactionKind::ConsensusCommitPrologueV1(_)
2258 ));
2259 }
2260 } else {
2261 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 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 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 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 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 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 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 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
2547async 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 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 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 highest_currently_built_seq_tx: watch::Sender<CheckpointSequenceNumber>,
2770 highest_previously_built_seq: CheckpointSequenceNumber,
2773 metrics: Arc<CheckpointMetrics>,
2774 state: Mutex<CheckpointServiceState>,
2775}
2776
2777impl CheckpointService {
2778 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 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 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 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 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 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}
3037pub 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#[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 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 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 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 assert!(
3266 store
3267 .get_checkpoint_contents(&contents_digest)
3268 .unwrap()
3269 .is_none()
3270 );
3271 }
3272
3273 #[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 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 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 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 Vec::new(),
3413 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 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 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 assert!(
3589 checkpoint_store
3590 .get_checkpoint_contents(&summary.contents_digest)
3591 .unwrap()
3592 .is_some()
3593 );
3594 }
3595
3596 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 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 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 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 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 Vec::new(),
3685 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 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 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 checkpoint_service
3768 .write_and_notify_checkpoint_for_testing(&epoch_store, p(3, vec![d(4)], 2010))
3769 .unwrap();
3770 checkpoint_service
3773 .write_and_notify_checkpoint_for_testing(&epoch_store, p(4, vec![d(5)], 2020))
3774 .unwrap();
3775 checkpoint_service
3779 .write_and_notify_checkpoint_for_testing(&epoch_store, p(5, vec![d(6)], 2030))
3780 .unwrap();
3781 checkpoint_service
3785 .write_and_notify_checkpoint_for_testing(&epoch_store, p(6, vec![d(7)], 2130))
3786 .unwrap();
3787 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 let sequence_numbers: Vec<_> = built.iter().map(|(seq, _)| *seq).collect();
3801 assert_eq!(sequence_numbers, vec![0, 1, 2, 3, 4, 5, 6]);
3802 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 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}