1use std::{future::Future, sync::Arc, time::Instant};
25
26use futures::StreamExt;
27use iota_common::{debug_fatal, fatal};
28use iota_config::node::{CheckpointExecutorConfig, RunWithRange};
29use iota_macros::fail_point;
30use iota_sdk_types::{
31 CheckpointContents, RandomnessRound, TransactionDigest, TransactionEffects,
32 TransactionEffectsDigest, TransactionKind,
33};
34use iota_types::{
35 base_types::ExecutionData,
36 effects::TransactionEffectsAPI,
37 executable_transaction::VerifiedExecutableTransaction,
38 full_checkpoint_content::CheckpointData,
39 global_state_hash::GlobalStateHash,
40 messages_checkpoint::{
41 CheckpointContentsExt, CheckpointSequenceNumber, CheckpointSummaryExt,
42 FullCheckpointContents, VerifiedCheckpoint,
43 },
44 transaction::{SenderSignedTransactionAPI, TransactionAPI, VerifiedTransaction},
45};
46use parking_lot::Mutex;
47use tap::{TapFallible, TapOptional};
48use tracing::{debug, info, instrument};
49
50use crate::{
51 authority::{
52 AuthorityState, ExecutionEnv, authority_per_epoch_store::AuthorityPerEpochStore,
53 backpressure::BackpressureManager,
54 shared_object_version_manager::gas_object_cancellation_version_from_effects,
55 },
56 checkpoint_progress_tracker::CheckpointProgressTracker,
57 checkpoints::CheckpointStore,
58 execution_cache::{ObjectCacheRead, TransactionCacheRead},
59 execution_scheduler::{ExecutionSchedulerAPI, ExecutionSchedulerWrapper},
60 global_state_hasher::GlobalStateHasher,
61};
62
63mod data_ingestion_handler;
64pub mod metrics;
65pub(crate) mod utils;
66
67#[cfg(test)]
68pub(crate) mod tests;
69
70use data_ingestion_handler::{load_checkpoint_data, store_checkpoint_locally};
71use metrics::CheckpointExecutorMetrics;
72use utils::*;
73
74type CheckpointDataSender = Box<dyn Fn(&CheckpointData) + Send + Sync>;
75
76#[derive(PartialEq, Eq, Debug)]
77pub enum StopReason {
78 EpochComplete,
79 RunWithRangeCondition,
80}
81
82pub(crate) struct CheckpointExecutionData {
83 pub checkpoint: VerifiedCheckpoint,
84 pub checkpoint_contents: CheckpointContents,
85 pub tx_digests: Vec<TransactionDigest>,
86 pub fx_digests: Vec<TransactionEffectsDigest>,
87}
88
89pub(crate) struct CheckpointTransactionData {
90 pub transactions: Vec<VerifiedExecutableTransaction>,
91 pub effects: Vec<TransactionEffects>,
92}
93
94pub(crate) struct CheckpointExecutionState {
95 pub data: CheckpointExecutionData,
96 state_hash: Option<GlobalStateHash>,
97 full_data: Option<CheckpointData>,
98}
99
100impl CheckpointExecutionState {
101 pub fn new(data: CheckpointExecutionData) -> Self {
102 Self {
103 data,
104 state_hash: None,
105 full_data: None,
106 }
107 }
108
109 pub fn new_with_global_state_hash(
110 data: CheckpointExecutionData,
111 hash: GlobalStateHash,
112 ) -> Self {
113 Self {
114 data,
115 state_hash: Some(hash),
116 full_data: None,
117 }
118 }
119}
120
121macro_rules! finish_stage {
122 ($handle:expr, $stage:ident) => {
123 $handle.finish_stage(PipelineStage::$stage).await;
124 };
125}
126
127pub struct CheckpointExecutor {
128 epoch_store: Arc<AuthorityPerEpochStore>,
129 state: Arc<AuthorityState>,
130 checkpoint_store: Arc<CheckpointStore>,
131 object_cache_reader: Arc<dyn ObjectCacheRead>,
132 transaction_cache_reader: Arc<dyn TransactionCacheRead>,
133 execution_scheduler: Arc<ExecutionSchedulerWrapper>,
134 global_state_hasher: Arc<GlobalStateHasher>,
135 backpressure_manager: Arc<BackpressureManager>,
136 config: CheckpointExecutorConfig,
137 metrics: Arc<CheckpointExecutorMetrics>,
138 tps_estimator: Mutex<TPSEstimator>,
139 checkpoint_progress_tracker: Option<Arc<CheckpointProgressTracker>>,
140 data_sender: Option<CheckpointDataSender>,
141}
142
143impl CheckpointExecutor {
144 pub fn new(
145 epoch_store: Arc<AuthorityPerEpochStore>,
146 checkpoint_store: Arc<CheckpointStore>,
147 state: Arc<AuthorityState>,
148 global_state_hasher: Arc<GlobalStateHasher>,
149 backpressure_manager: Arc<BackpressureManager>,
150 config: CheckpointExecutorConfig,
151 metrics: Arc<CheckpointExecutorMetrics>,
152 data_sender: Option<CheckpointDataSender>,
153 checkpoint_progress_tracker: Option<Arc<CheckpointProgressTracker>>,
154 ) -> Self {
155 Self {
156 epoch_store,
157 state: state.clone(),
158 checkpoint_store,
159 object_cache_reader: state.get_object_cache_reader().clone(),
160 transaction_cache_reader: state.get_transaction_cache_reader().clone(),
161 execution_scheduler: state.execution_scheduler().clone(),
162 global_state_hasher,
163 backpressure_manager,
164 config,
165 metrics,
166 tps_estimator: Mutex::new(TPSEstimator::default()),
167 checkpoint_progress_tracker,
168 data_sender,
169 }
170 }
171
172 pub fn new_for_tests(
173 epoch_store: Arc<AuthorityPerEpochStore>,
174 checkpoint_store: Arc<CheckpointStore>,
175 state: Arc<AuthorityState>,
176 global_state_hasher: Arc<GlobalStateHasher>,
177 ) -> Self {
178 Self::new(
179 epoch_store,
180 checkpoint_store.clone(),
181 state,
182 global_state_hasher,
183 BackpressureManager::new_from_checkpoint_store(&checkpoint_store),
184 Default::default(),
185 CheckpointExecutorMetrics::new_for_tests(),
186 None, None, )
189 }
190
191 fn get_next_to_schedule(&self) -> Option<CheckpointSequenceNumber> {
194 let highest_executed = self
198 .checkpoint_store
199 .get_highest_executed_checkpoint()
200 .unwrap();
201
202 if let Some(highest_executed) = &highest_executed {
203 if self.epoch_store.epoch() == highest_executed.epoch()
204 && highest_executed.is_last_checkpoint_of_epoch()
205 {
206 info!(seq = ?highest_executed.sequence_number, "final checkpoint of epoch has already been executed");
209 return None;
210 }
211 }
212
213 Some(
214 highest_executed
215 .as_ref()
216 .map(|c| c.sequence_number() + 1)
217 .unwrap_or_else(|| {
218 assert_eq!(self.epoch_store.epoch(), 0);
220 0
222 }),
223 )
224 }
225
226 #[instrument(level = "error", skip_all, fields(epoch = ?self.epoch_store.epoch()))]
230 pub async fn run_epoch(self, run_with_range: Option<RunWithRange>) -> StopReason {
231 let _metrics_scope = iota_metrics::monitored_scope("CheckpointExecutor::run_epoch");
232 info!(?run_with_range, "CheckpointExecutor::run_epoch");
233 debug!(
234 "Checkpoint executor running for epoch {:?}",
235 self.epoch_store.epoch(),
236 );
237
238 if run_with_range.is_some_and(|rwr| rwr.is_epoch_gt(self.epoch_store.epoch())) {
242 info!("RunWithRange condition satisfied at {:?}", run_with_range,);
243 return StopReason::RunWithRangeCondition;
244 };
245
246 self.metrics
247 .checkpoint_exec_epoch
248 .set(self.epoch_store.epoch() as i64);
249
250 let Some(next_to_schedule) = self.get_next_to_schedule() else {
251 return StopReason::EpochComplete;
252 };
253
254 let this = Arc::new(self);
255
256 let concurrency = std::env::var("IOTA_CHECKPOINT_EXECUTION_MAX_CONCURRENCY")
257 .ok()
258 .and_then(|s| s.parse().ok())
259 .unwrap_or(this.config.checkpoint_execution_max_concurrency);
260
261 let pipeline_stages = PipelineStages::new(
262 next_to_schedule,
263 this.metrics.clone(),
264 this.checkpoint_progress_tracker.clone(),
265 );
266
267 let final_checkpoint_executed = stream_synced_checkpoints(
268 this.checkpoint_store.clone(),
269 next_to_schedule,
270 run_with_range.and_then(|rwr| rwr.into_checkpoint_bound()),
271 )
272 .map(|checkpoint| {
276 let this = this.clone();
277 let pipeline_stages = pipeline_stages.clone();
278 async move {
279 let seq = checkpoint.sequence_number();
280 let pipeline_handle = async move { pipeline_stages.handle(seq).await };
281 tokio::spawn(this.execute_checkpoint(checkpoint, pipeline_handle))
282 .await
283 .unwrap()
284 }
285 })
286 .buffered(concurrency)
287 .fold(false, |state, is_final_checkpoint| async move {
289 assert!(!state, "Cannot execute checkpoint after epoch end");
290 is_final_checkpoint
291 })
292 .await;
293
294 if final_checkpoint_executed {
295 StopReason::EpochComplete
296 } else {
297 StopReason::RunWithRangeCondition
298 }
299 }
300}
301
302impl CheckpointExecutor {
303 #[instrument(level = "debug", skip_all, fields(seq = ?checkpoint.sequence_number()))]
306 async fn execute_checkpoint(
307 self: Arc<Self>,
308 checkpoint: VerifiedCheckpoint,
309 pipeline_handle: impl Future<Output = PipelineHandle> + Send,
310 ) -> bool {
311 debug!("executing checkpoint");
312 let sequence_number = checkpoint.sequence_number;
313 let is_last_checkpoint_of_epoch = checkpoint.is_last_checkpoint_of_epoch();
314
315 checkpoint.report_checkpoint_age(&self.metrics.checkpoint_contents_age);
316 self.backpressure_manager
317 .update_highest_certified_checkpoint(sequence_number);
318
319 let loaded = (self.state.is_fullnode(&self.epoch_store) || is_last_checkpoint_of_epoch)
326 .then(|| {
327 let _scope = iota_metrics::monitored_scope(
328 "CheckpointExecutor::load_checkpoint_transactions",
329 );
330 self.load_checkpoint_transactions(checkpoint.clone())
331 });
332
333 let mut pipeline_handle = pipeline_handle.await;
334
335 if is_last_checkpoint_of_epoch && sequence_number > 0 {
336 let _wait_for_previous_checkpoints_guard =
337 iota_metrics::monitored_scope("CheckpointExecutor::wait_for_previous_checkpoints");
338
339 info!(
340 "Reached end of epoch checkpoint, waiting for all previous checkpoints to be executed"
341 );
342 self.checkpoint_store
343 .notify_read_executed_checkpoint(sequence_number - 1)
344 .await;
345 }
346
347 let _parallel_step_guard =
348 iota_metrics::monitored_scope("CheckpointExecutor::parallel_step");
349
350 let exec_start = Instant::now();
353 let ckpt_state = match loaded {
354 Some((ckpt_state, tx_data)) => {
355 self.execute_transactions_from_synced_checkpoint(
356 ckpt_state,
357 tx_data,
358 &mut pipeline_handle,
359 )
360 .await
361 }
362 None => {
363 self.verify_locally_built_checkpoint(checkpoint, &mut pipeline_handle)
364 .await
365 }
366 };
367
368 let tps = self.tps_estimator.lock().update(
369 Instant::now(),
370 ckpt_state.data.checkpoint.network_total_transactions,
371 );
372 self.metrics.checkpoint_exec_sync_tps.set(tps as i64);
373
374 self.backpressure_manager
375 .update_highest_executed_checkpoint(ckpt_state.data.checkpoint.sequence_number());
376
377 let is_final_checkpoint = ckpt_state.data.checkpoint.is_last_checkpoint_of_epoch();
378
379 let seq = ckpt_state.data.checkpoint.sequence_number;
380
381 let batch = self.state.get_cache_commit().build_db_batch(
382 self.epoch_store.epoch(),
383 seq,
384 &ckpt_state.data.tx_digests,
385 );
386
387 finish_stage!(pipeline_handle, BuildDbBatch);
388
389 let mut ckpt_state = tokio::task::spawn_blocking({
390 let this = self.clone();
391 move || {
392 let cache_commit = this.state.get_cache_commit();
394 debug!(?seq, "committing checkpoint transactions to disk");
395 cache_commit.commit_transaction_outputs(
396 this.epoch_store.epoch(),
397 batch,
398 &ckpt_state.data.tx_digests,
399 );
400 ckpt_state
401 }
402 })
403 .await
404 .unwrap();
405
406 finish_stage!(pipeline_handle, CommitTransactionOutputs);
407
408 self.epoch_store
409 .handle_finalized_checkpoint(&ckpt_state.data.checkpoint, &ckpt_state.data.tx_digests)
410 .expect("cannot fail");
411
412 let randomness_rounds = self.extract_randomness_rounds(
413 &ckpt_state.data.checkpoint,
414 &ckpt_state.data.checkpoint_contents,
415 );
416
417 if let Some(randomness_reporter) = self.epoch_store.randomness_reporter() {
422 for round in randomness_rounds {
423 debug!(
424 ?round,
425 "notifying RandomnessReporter that randomness update was executed in checkpoint"
426 );
427 randomness_reporter
428 .notify_randomness_in_checkpoint(round)
429 .expect("epoch cannot have ended");
430 }
431 }
432
433 finish_stage!(pipeline_handle, FinalizeCheckpoint);
434
435 if let Some(checkpoint_data) = ckpt_state.full_data.take() {
436 self.commit_index_updates(checkpoint_data);
437 }
438
439 finish_stage!(pipeline_handle, UpdateRpcIndex);
440
441 self.global_state_hasher
442 .accumulate_running_root(&self.epoch_store, seq, ckpt_state.state_hash)
443 .expect("Failed to accumulate running root");
444
445 if is_final_checkpoint {
446 self.checkpoint_store
447 .insert_epoch_last_checkpoint(self.epoch_store.epoch(), &ckpt_state.data.checkpoint)
448 .expect("Failed to insert epoch last checkpoint");
449
450 self.global_state_hasher
451 .accumulate_epoch(self.epoch_store.clone(), seq)
452 .expect("Accumulating epoch cannot fail");
453
454 self.checkpoint_store
455 .prune_local_summaries()
456 .tap_err(|e| debug_fatal!("Failed to prune local summaries: {}", e))
457 .ok();
458 }
459
460 fail_point!("crash");
461
462 self.bump_highest_executed_checkpoint(&ckpt_state.data.checkpoint);
463
464 self.broadcast_checkpoint(&ckpt_state.data, ckpt_state.full_data.as_ref());
465
466 self.state
469 .pruner()
470 .nudge(ckpt_state.data.checkpoint.sequence_number());
471
472 finish_stage!(pipeline_handle, BumpHighestExecutedCheckpoint);
473
474 if let Some(tracker) = &self.checkpoint_progress_tracker {
475 tracker.add_execution_time(exec_start.elapsed());
476 }
477
478 ckpt_state.data.checkpoint.is_last_checkpoint_of_epoch()
482 }
483
484 #[instrument(level = "info", skip_all)]
487 async fn verify_locally_built_checkpoint(
488 &self,
489 checkpoint: VerifiedCheckpoint,
490 pipeline_handle: &mut PipelineHandle,
491 ) -> CheckpointExecutionState {
492 assert!(
493 !checkpoint.is_last_checkpoint_of_epoch(),
494 "only fullnode path has end-of-epoch logic"
495 );
496
497 let sequence_number = checkpoint.sequence_number;
498 let locally_built_checkpoint = self
499 .checkpoint_store
500 .get_locally_computed_checkpoint(sequence_number)
501 .expect("db error");
502
503 let Some(locally_built_checkpoint) = locally_built_checkpoint else {
504 let (ckpt_state, tx_data) = self.load_checkpoint_transactions(checkpoint);
506 return self
507 .execute_transactions_from_synced_checkpoint(ckpt_state, tx_data, pipeline_handle)
508 .await;
509 };
510
511 self.metrics.checkpoint_executor_validator_path.inc();
512
513 assert_checkpoint_not_forked(
515 &locally_built_checkpoint,
516 &checkpoint,
517 &self.checkpoint_store,
518 );
519
520 let state_hash = {
523 let _metrics_scope =
524 iota_metrics::monitored_scope("CheckpointExecutor::notify_read_state_hash");
525 self.epoch_store
526 .notify_read_checkpoint_state_hasher(&[sequence_number])
527 .await
528 .unwrap()
529 .pop()
530 .unwrap()
531 };
532
533 let checkpoint_contents = self
537 .checkpoint_store
538 .get_checkpoint_contents(&checkpoint.contents_digest)
539 .expect("db error")
540 .expect("checkpoint contents not found");
541
542 let (tx_digests, fx_digests): (Vec<_>, Vec<_>) = checkpoint_contents
543 .iter()
544 .map(|digests| (digests.transaction, digests.effects))
545 .unzip();
546
547 pipeline_handle
548 .skip_to(PipelineStage::FinalizeTransactions)
549 .await;
550
551 self.insert_finalized_transactions(&tx_digests, sequence_number, checkpoint.timestamp_ms);
555
556 pipeline_handle.skip_to(PipelineStage::BuildDbBatch).await;
557
558 CheckpointExecutionState::new_with_global_state_hash(
559 CheckpointExecutionData {
560 checkpoint,
561 checkpoint_contents,
562 tx_digests,
563 fx_digests,
564 },
565 state_hash,
566 )
567 }
568
569 #[instrument(level = "info", skip_all)]
570 async fn execute_transactions_from_synced_checkpoint(
571 &self,
572 mut ckpt_state: CheckpointExecutionState,
573 tx_data: CheckpointTransactionData,
574 pipeline_handle: &mut PipelineHandle,
575 ) -> CheckpointExecutionState {
576 let sequence_number = ckpt_state.data.checkpoint.sequence_number;
577
578 let (unexecuted_tx_digests, unexecuted_expected_fx_digests) = {
579 let _scope = iota_metrics::monitored_scope("CheckpointExecutor::execute_transactions");
580 self.schedule_transaction_execution(&ckpt_state, &tx_data)
581 };
582
583 finish_stage!(pipeline_handle, ExecuteTransactions);
584
585 {
586 let actual_fx_digests = self
587 .transaction_cache_reader
588 .notify_read_executed_effects_digests(
589 "CheckpointExecutor::notify_read_executed_effects_digests",
590 &unexecuted_tx_digests,
591 )
592 .await;
593
594 for (tx_digest, expected, actual) in itertools::izip!(
600 unexecuted_tx_digests.iter(),
601 unexecuted_expected_fx_digests.iter(),
602 actual_fx_digests.iter()
603 ) {
604 assert_not_forked(
605 &ckpt_state.data.checkpoint,
606 tx_digest,
607 expected,
608 actual,
609 &*self.transaction_cache_reader,
610 );
611 }
612 }
613
614 finish_stage!(pipeline_handle, WaitForTransactions);
615
616 if ckpt_state.data.checkpoint.is_last_checkpoint_of_epoch() {
617 self.execute_change_epoch_tx(&tx_data).await;
618 }
619
620 let _scope = iota_metrics::monitored_scope("CheckpointExecutor::finalize_checkpoint");
621
622 if self.state.is_fullnode(&self.epoch_store) {
623 self.state.congestion_tracker.process_checkpoint_effects(
624 &*self.transaction_cache_reader,
625 &ckpt_state.data.checkpoint,
626 &tx_data.effects,
627 );
628 }
629
630 self.insert_finalized_transactions(
631 &ckpt_state.data.tx_digests,
632 sequence_number,
633 ckpt_state.data.checkpoint.timestamp_ms,
634 );
635
636 ckpt_state.state_hash = Some(
640 self.global_state_hasher
641 .accumulate_checkpoint(&tx_data.effects, sequence_number, &self.epoch_store)
642 .expect("epoch cannot have ended"),
643 );
644
645 finish_stage!(pipeline_handle, FinalizeTransactions);
646
647 ckpt_state.full_data = self.process_checkpoint_data(&ckpt_state.data, &tx_data);
648
649 finish_stage!(pipeline_handle, ProcessCheckpointData);
650
651 ckpt_state
652 }
653
654 fn checkpoint_data_enabled(&self) -> bool {
655 self.state.grpc_indexes_store.is_some()
656 || self.config.data_ingestion_dir.is_some()
657 || self.data_sender.is_some()
658 }
659
660 fn insert_finalized_transactions(
661 &self,
662 tx_digests: &[TransactionDigest],
663 sequence_number: CheckpointSequenceNumber,
664 timestamp_ms: u64,
665 ) {
666 self.epoch_store
667 .insert_finalized_transactions(tx_digests, sequence_number, timestamp_ms)
668 .expect("failed to insert finalized transactions");
669
670 if self.state.is_fullnode(&self.epoch_store) {
671 self.state
673 .get_checkpoint_cache()
674 .insert_finalized_transactions_perpetual_checkpoints(
675 tx_digests,
676 self.epoch_store.epoch(),
677 sequence_number,
678 );
679 }
680 }
681
682 #[instrument(level = "info", skip_all)]
683 fn process_checkpoint_data(
684 &self,
685 ckpt_data: &CheckpointExecutionData,
686 tx_data: &CheckpointTransactionData,
687 ) -> Option<CheckpointData> {
688 let is_checkpoint_data_enabled = self.checkpoint_data_enabled();
689 let is_last_checkpoint_of_epoch = ckpt_data.checkpoint.is_last_checkpoint_of_epoch();
692 if !is_checkpoint_data_enabled && !is_last_checkpoint_of_epoch {
693 return None;
694 }
695
696 let checkpoint_data = load_checkpoint_data(
697 ckpt_data,
698 tx_data,
699 self.state.get_object_store(),
700 &*self.transaction_cache_reader,
701 )
702 .expect("failed to load checkpoint data");
703
704 if is_last_checkpoint_of_epoch {
713 self.checkpoint_store
714 .index_epoch_boundary(&checkpoint_data)
715 .expect("failed to persist epoch info at boundary");
716 }
717
718 if !is_checkpoint_data_enabled {
719 return None;
721 }
722
723 if let Some(grpc_indexes_store) = &self.state.grpc_indexes_store {
730 grpc_indexes_store.index_checkpoint(&checkpoint_data);
731 }
732
733 if let Some(path) = &self.config.data_ingestion_dir {
734 store_checkpoint_locally(path, &checkpoint_data)
735 .expect("failed to store checkpoint locally");
736 }
737
738 Some(checkpoint_data)
739 }
740
741 #[instrument(level = "info", skip_all)]
743 fn load_checkpoint_transactions(
744 &self,
745 checkpoint: VerifiedCheckpoint,
746 ) -> (CheckpointExecutionState, CheckpointTransactionData) {
747 let seq = checkpoint.sequence_number;
748 let epoch = checkpoint.epoch;
749
750 let checkpoint_contents = self
751 .checkpoint_store
752 .get_checkpoint_contents(&checkpoint.contents_digest)
753 .expect("db error")
754 .expect("checkpoint contents not found");
755
756 if let Some(full_contents) = self
758 .checkpoint_store
759 .get_full_checkpoint_contents_by_sequence_number(seq)
760 .tap_some(|_| debug!("loaded full checkpoint contents in bulk for sequence {seq}"))
761 {
762 let num_txns = full_contents.size();
763 let mut tx_digests = Vec::with_capacity(num_txns);
764 let mut transactions = Vec::with_capacity(num_txns);
765 let mut effects = Vec::with_capacity(num_txns);
766 let mut fx_digests = Vec::with_capacity(num_txns);
767
768 full_contents
769 .iter()
770 .zip(checkpoint_contents.iter())
771 .for_each(|(execution_data, digests)| {
772 let tx_digest = digests.transaction;
773 let fx_digest = digests.effects;
774 debug_assert_eq!(tx_digest, *execution_data.transaction.digest());
775 debug_assert_eq!(fx_digest, execution_data.effects.digest());
776
777 tx_digests.push(tx_digest);
778 transactions.push(VerifiedExecutableTransaction::new_from_checkpoint(
779 VerifiedTransaction::new_unchecked(execution_data.transaction.clone()),
780 epoch,
781 seq,
782 ));
783 effects.push(execution_data.effects.clone());
784 fx_digests.push(fx_digest);
785 });
786
787 (
788 CheckpointExecutionState::new(CheckpointExecutionData {
789 checkpoint,
790 checkpoint_contents,
791 tx_digests,
792 fx_digests,
793 }),
794 CheckpointTransactionData {
795 transactions,
796 effects,
797 },
798 )
799 } else {
800 let digests = checkpoint_contents.transactions();
803
804 let (tx_digests, fx_digests): (Vec<_>, Vec<_>) =
805 digests.iter().map(|d| (d.transaction, d.effects)).unzip();
806 let verified_transactions: Vec<VerifiedTransaction> = self
807 .transaction_cache_reader
808 .multi_get_transaction_blocks(&tx_digests)
809 .into_iter()
810 .enumerate()
811 .map(|(i, tx)| {
812 let tx = tx
813 .unwrap_or_else(|| fatal!("transaction not found for {:?}", tx_digests[i]));
814 Arc::try_unwrap(tx).unwrap_or_else(|tx| (*tx).clone())
815 })
816 .collect();
817 let effects: Vec<TransactionEffects> = self
818 .transaction_cache_reader
819 .multi_get_effects(&fx_digests)
820 .into_iter()
821 .enumerate()
822 .map(|(i, effect)| {
823 effect.unwrap_or_else(|| {
824 fatal!("checkpoint effect not found for {:?}", digests[i])
825 })
826 })
827 .collect();
828
829 if self
838 .checkpoint_store
839 .should_cache_full_checkpoint_contents(seq)
840 {
841 let execution_data = verified_transactions
842 .iter()
843 .zip(effects.iter())
844 .map(|(tx, fx)| ExecutionData::new(tx.clone().into_inner(), fx.clone()));
845 let full_contents = FullCheckpointContents::from_contents_and_execution_data(
846 checkpoint_contents.clone(),
847 execution_data,
848 )
849 .expect("checkpoint has one executed transaction per pinned signature set");
852 self.checkpoint_store.cache_full_checkpoint_contents(
853 seq,
854 checkpoint.contents_digest,
855 full_contents,
856 );
857 }
858
859 let transactions = verified_transactions
860 .into_iter()
861 .map(|tx| VerifiedExecutableTransaction::new_from_checkpoint(tx, epoch, seq))
862 .collect();
863
864 (
865 CheckpointExecutionState::new(CheckpointExecutionData {
866 checkpoint,
867 checkpoint_contents,
868 tx_digests,
869 fx_digests,
870 }),
871 CheckpointTransactionData {
872 transactions,
873 effects,
874 },
875 )
876 }
877 }
878
879 #[instrument(level = "info", skip_all)]
883 fn schedule_transaction_execution(
884 &self,
885 ckpt_state: &CheckpointExecutionState,
886 tx_data: &CheckpointTransactionData,
887 ) -> (Vec<TransactionDigest>, Vec<TransactionEffectsDigest>) {
888 let executed_fx_digests = self
893 .transaction_cache_reader
894 .multi_get_executed_effects_digests(&ckpt_state.data.tx_digests);
895
896 let (unexecuted_tx_digests, unexecuted_expected_fx_digests, unexecuted_txns): (
898 Vec<_>,
899 Vec<_>,
900 Vec<_>,
901 ) = itertools::multiunzip(
902 itertools::izip!(
903 tx_data.transactions.iter(),
904 ckpt_state.data.tx_digests.iter(),
905 ckpt_state.data.fx_digests.iter(),
906 tx_data.effects.iter(),
907 executed_fx_digests.iter()
908 )
909 .filter_map(
910 |(txn, tx_digest, expected_fx_digest, effects, executed_fx_digest)| {
911 if let Some(executed_fx_digest) = executed_fx_digest {
912 assert_not_forked(
913 &ckpt_state.data.checkpoint,
914 tx_digest,
915 expected_fx_digest,
916 executed_fx_digest,
917 &*self.transaction_cache_reader,
918 );
919 None
920 } else if txn.transaction().is_end_of_epoch_tx() {
921 None
922 } else {
923 let assigned_versions = if txn.contains_shared_object()
928 || gas_object_cancellation_version_from_effects(effects).is_some()
929 {
930 self.epoch_store
931 .acquire_shared_version_assignments_from_effects(
932 txn,
933 effects,
934 &*self.object_cache_reader,
935 )
936 .expect("failed to initialize shared object versions from effects")
937 } else {
938 Default::default()
939 };
940
941 let env = ExecutionEnv::new()
942 .with_assigned_versions(assigned_versions)
943 .with_expected_effects_digest(*expected_fx_digest);
944
945 Some((tx_digest, *expected_fx_digest, (txn.clone(), env)))
946 }
947 },
948 ),
949 );
950
951 self.execution_scheduler
953 .enqueue_transactions(unexecuted_txns, &self.epoch_store);
954
955 (unexecuted_tx_digests, unexecuted_expected_fx_digests)
956 }
957
958 #[instrument(level = "error", skip_all)]
960 async fn execute_change_epoch_tx(&self, tx_data: &CheckpointTransactionData) {
961 let change_epoch_tx = tx_data.transactions.last().unwrap();
962 let change_epoch_fx = tx_data.effects.last().unwrap();
963 assert_eq!(
964 change_epoch_tx.digest(),
965 change_epoch_fx.transaction_digest()
966 );
967 assert!(
968 change_epoch_tx.transaction().is_end_of_epoch_tx(),
969 "final txn must be an end of epoch txn"
970 );
971
972 let assigned_versions = self
991 .epoch_store
992 .acquire_shared_version_assignments_from_effects(
993 change_epoch_tx,
994 change_epoch_fx,
995 self.object_cache_reader.as_ref(),
996 )
997 .expect("failed to initialize shared object versions for the change epoch transaction");
998
999 info!(
1000 "scheduling change epoch txn with digest: {:?}, expected effects digest: {:?}, \
1001 assigned versions: {:?}",
1002 change_epoch_tx.digest(),
1003 change_epoch_fx.digest(),
1004 assigned_versions
1005 );
1006 self.execution_scheduler.enqueue_transactions(
1007 vec![(
1008 change_epoch_tx.clone(),
1009 ExecutionEnv::new()
1010 .with_assigned_versions(assigned_versions)
1011 .with_expected_effects_digest(change_epoch_fx.digest()),
1012 )],
1013 &self.epoch_store,
1014 );
1015
1016 self.transaction_cache_reader
1017 .notify_read_executed_effects_digests(
1018 "CheckpointExecutor::notify_read_advance_epoch_tx",
1019 &[*change_epoch_tx.digest()],
1020 )
1021 .await;
1022 }
1023
1024 #[instrument(level = "debug", skip_all)]
1026 fn bump_highest_executed_checkpoint(&self, checkpoint: &VerifiedCheckpoint) {
1027 let seq = checkpoint.sequence_number();
1029 debug!("Bumping highest_executed_checkpoint watermark to {seq:?}");
1030 if let Some(prev_highest) = self
1031 .checkpoint_store
1032 .get_highest_executed_checkpoint_seq_number()
1033 .unwrap()
1034 {
1035 assert_eq!(prev_highest + 1, seq);
1036 } else {
1037 assert_eq!(seq, 0);
1038 }
1039 fail_point!("highest-executed-checkpoint");
1040
1041 self.checkpoint_store
1042 .update_highest_executed_checkpoint(checkpoint)
1043 .unwrap();
1044 self.metrics.last_executed_checkpoint.set(seq as i64);
1045
1046 self.metrics
1047 .last_executed_checkpoint_timestamp_ms
1048 .set(checkpoint.timestamp_ms as i64);
1049 checkpoint.report_checkpoint_age(&self.metrics.last_executed_checkpoint_age);
1050 }
1051
1052 fn broadcast_checkpoint(
1055 &self,
1056 checkpoint_exec_data: &CheckpointExecutionData,
1057 checkpoint_data: Option<&CheckpointData>,
1058 ) {
1059 if let Some(data_sender) = &self.data_sender {
1060 let checkpoint_data = if let Some(data) = checkpoint_data {
1061 data.clone()
1062 } else {
1063 let (_, tx_data) =
1066 self.load_checkpoint_transactions(checkpoint_exec_data.checkpoint.clone());
1067 load_checkpoint_data(
1068 checkpoint_exec_data,
1069 &tx_data,
1070 self.state.get_object_store(),
1071 self.transaction_cache_reader.as_ref(),
1072 )
1073 .expect("Failed to load full CheckpointData")
1074 };
1075 data_sender(&checkpoint_data);
1076 }
1077
1078 debug!(
1079 "[Fullnode] Full CheckpointData is available: seq={}",
1080 checkpoint_exec_data.checkpoint.sequence_number()
1081 );
1082 }
1083
1084 #[instrument(level = "info", skip_all)]
1087 fn commit_index_updates(&self, checkpoint: CheckpointData) {
1088 if let Some(grpc_indexes_store) = &self.state.grpc_indexes_store {
1089 grpc_indexes_store
1090 .commit_update_for_checkpoint(checkpoint.checkpoint_summary.sequence_number)
1091 .expect("failed to update gRPC indexes");
1092 }
1093 }
1094
1095 #[instrument(level = "debug", skip_all)]
1099 fn extract_randomness_rounds(
1100 &self,
1101 checkpoint: &VerifiedCheckpoint,
1102 checkpoint_contents: &CheckpointContents,
1103 ) -> Vec<RandomnessRound> {
1104 if let Some(version_specific_data) = checkpoint
1105 .parse_version_specific_data(self.epoch_store.protocol_config())
1106 .expect("unable to get version_specific_data")
1107 {
1108 version_specific_data.into_v1().randomness_rounds
1111 } else {
1112 assert_eq!(
1117 0,
1118 self.epoch_store
1119 .protocol_config()
1120 .min_checkpoint_interval_ms_as_option()
1121 .unwrap_or_default(),
1122 );
1123 if let Some(first_digest) = checkpoint_contents.transactions().first() {
1124 let maybe_randomness_tx = self.transaction_cache_reader.get_transaction_block(&first_digest.transaction)
1125 .unwrap_or_else(||
1126 fatal!(
1127 "state-sync should have ensured that transaction with digests {first_digest:?} exists for checkpoint: {}",
1128 checkpoint.sequence_number()
1129 )
1130 );
1131 if let TransactionKind::RandomnessStateUpdate(rsu) =
1132 maybe_randomness_tx.data().transaction().kind()
1133 {
1134 vec![rsu.randomness_round]
1135 } else {
1136 Vec::new()
1137 }
1138 } else {
1139 Vec::new()
1140 }
1141 }
1142 }
1143}