Skip to main content

iota_core/checkpoints/checkpoint_executor/
mod.rs

1// Copyright (c) Mysten Labs, Inc.
2// Modifications Copyright (c) 2024 IOTA Stiftung
3// SPDX-License-Identifier: Apache-2.0
4
5//! CheckpointExecutor is a Node component that executes all checkpoints for the
6//! given epoch. It acts as a Consumer to StateSync
7//! for newly synced checkpoints, taking these checkpoints and
8//! scheduling and monitoring their execution. Its primary goal is to allow
9//! for catching up to the current checkpoint sequence number of the network
10//! as quickly as possible so that a newly joined, or recovering Node can
11//! participate in a timely manner. To that end, CheckpointExecutor attempts
12//! to saturate the CPU with executor tasks (one per checkpoint), each of which
13//! handle scheduling and awaiting checkpoint transaction execution.
14//!
15//! CheckpointExecutor is made recoverable in the event of Node shutdown by way
16//! of a watermark, highest_executed_checkpoint, which is guaranteed to be
17//! updated sequentially in order, despite checkpoints themselves potentially
18//! being executed nonsequentially and in parallel. CheckpointExecutor
19//! parallelizes checkpoints of the same epoch as much as possible.
20//! CheckpointExecutor enforces the invariant that if `run` returns
21//! successfully, we have reached the end of epoch. This allows us to use it as
22//! a signal for reconfig.
23
24use 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, // No callback for data
187            None, // No progress tracker for tests
188        )
189    }
190
191    // Gets the next checkpoint to schedule for execution. If the epoch is already
192    // completed, returns None.
193    fn get_next_to_schedule(&self) -> Option<CheckpointSequenceNumber> {
194        // Decide the first checkpoint to schedule for execution.
195        // If we haven't executed anything in the past, we schedule checkpoint 0.
196        // Otherwise we schedule the one after highest executed.
197        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                // We can arrive at this point if we bump the highest_executed_checkpoint
207                // watermark, and then crash before completing reconfiguration.
208                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                    // TODO this invariant may no longer hold once we introduce snapshots
219                    assert_eq!(self.epoch_store.epoch(), 0);
220                    // we need to execute the genesis checkpoint
221                    0
222                }),
223        )
224    }
225
226    /// Execute all checkpoints for the current epoch, ensuring that the node
227    /// has not forked, and return when finished.
228    /// If `run_with_range` is set, execution will stop early.
229    #[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        // check if we want to run this epoch based on RunWithRange condition value
239        // we want to be inclusive of the defined RunWithRangeEpoch::Epoch
240        // i.e Epoch(N) means we will execute epoch N and stop when reaching N+1
241        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        // Checkpoint loading and execution is parallelized: each
273        // checkpoint's task loads its transaction data before awaiting
274        // admission to the ordered pipeline stages.
275        .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        // Take the last value from the stream to determine if we completed the epoch
288        .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    /// Load all data for a checkpoint, ensure all transactions are executed,
304    /// and check for forks.
305    #[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 /* is final checkpoint */ {
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        // Load the transaction data before awaiting admission to the
320        // pipeline: it reads only state-sync-written data (contents,
321        // transactions, and expected effects), so it can run in parallel
322        // across the buffered checkpoints instead of inside the ordered
323        // ExecuteTransactions stage. The clone is needed because the
324        // locally-built-checkpoint path below wants the checkpoint too.
325        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        // Note: only `execute_transactions_from_synced_checkpoint` has end-of-epoch
351        // logic.
352        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                // Commit all transaction effects to disk
393                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        // Once the checkpoint is finalized, we know that any randomness contained in
418        // this checkpoint has been successfully included in a checkpoint
419        // certified by quorum of validators. (RandomnessManager/
420        // RandomnessReporter is only present on validators.)
421        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        // Nudge the pruner now that this checkpoint is executed and available;
467        // pruning of aged-out data runs off the propagation path.
468        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        // Important: code after the last pipeline stage is finished can run out of
479        // checkpoint order.
480
481        ckpt_state.data.checkpoint.is_last_checkpoint_of_epoch()
482    }
483
484    // On validators, checkpoints have often already been constructed locally, in
485    // which case we can skip many steps of the checkpoint execution process.
486    #[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            // fall back to tx-by-tx execution path if we are catching up.
505            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        // Check for fork
514        assert_checkpoint_not_forked(
515            &locally_built_checkpoint,
516            &checkpoint,
517            &self.checkpoint_store,
518        );
519
520        // Checkpoint builder triggers accumulation of the checkpoint, so this is
521        // guaranteed to finish.
522        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        // Checkpoint builder triggers accumulation of the checkpoint, so this is
534        // guaranteed to finish.
535
536        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        // Currently this code only runs on validators, where this method call does
552        // nothing. But in the future, fullnodes may follow the consensus dag
553        // and build their own checkpoints.
554        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            // A transaction scheduled here can also be scheduled from consensus,
595            // and whichever enqueue arrives second is dropped along with its
596            // environment — so the expected digest passed above may never have
597            // reached execution. Compare here, where the checkpoint's digests are
598            // known, so a fork is caught regardless of which enqueue won.
599            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        // The early versions of the hasher (prior to effectsv2) rely on db
637        // state, so we must wait until all transactions have been executed
638        // before accumulating the checkpoint.
639        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            // TODO remove once we no longer need to support this table for read RPC
672            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        // Boundaries always need full `CheckpointData` to persist `epoch_info`,
690        // even when no other consumer is configured.
691        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        // Persist the boundary's `epoch_info` row eagerly. Two properties make
705        // that safe. Boundaries run in epoch order: the boundary waits for every
706        // earlier checkpoint (`notify_read_executed_checkpoint(seq - 1)` in
707        // `execute_checkpoint`) and the stream stops at epoch end, so two
708        // boundaries are never in flight. And `index_epoch_boundary` is
709        // idempotent — a row upsert plus a contiguous +1 watermark guard — so a
710        // crash before the executed watermark advances just re-applies it on the
711        // next run.
712        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            // Data was built solely to seed `epoch_info`; skip the other consumers.
720            return None;
721        }
722
723        // Index the checkpoint. The grpc indexes accumulate non-idempotent state
724        // (owner indexes, live-object sets), so each update must land exactly
725        // once. Indexing runs here out of order (checkpoints execute
726        // concurrently), so the write is only staged now and committed later, in
727        // sequence order, via `commit_update_for_checkpoint` — keeping the grpc
728        // watermark consistent with the executed checkpoint and crash-safe.
729        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    // Load all required transaction and effects data for the checkpoint.
742    #[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        // attempt to load full checkpoint contents in bulk
757        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            // load items one-by-one
801
802            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            // Reached when the contents were not already cached: a fullnode
830            // syncing from a peer, or a node catching up. Cache the assembled
831            // contents so this node can in turn serve state-sync peers without
832            // reconstruction. (Validators' locally built checkpoints are cached
833            // by the checkpoint builder and take the bulk path above instead.)
834            //
835            // The assembly clones every transaction and effect, so skip it
836            // when the cache wouldn't retain the entry (see `should_cache`).
837            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                // Both were loaded for this checkpoint's own contents above,
850                // one transaction and one effects entry per digest.
851                .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    /// Enqueues the checkpoint's not-yet-executed transactions, and returns
880    /// their digests together with the effects digest the checkpoint expects
881    /// for each, so the caller can check for a fork once they execute.
882    #[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        // Which transactions have already been executed must be read here, in
889        // the ordered pipeline stage, and not when the checkpoint data is
890        // loaded: execution progresses in between, and a stale answer would
891        // skip the fork check for a transaction that executed since loading.
892        let executed_fx_digests = self
893            .transaction_cache_reader
894            .multi_get_executed_effects_digests(&ckpt_state.data.tx_digests);
895
896        // Find unexecuted transactions and their expected effects digests
897        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                        // A transaction without shared inputs normally needs no version
924                        // assignments, but one cancelled for congestion carries its
925                        // cancellation version on the gas object, which re-execution
926                        // must reproduce.
927                        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        // Enqueue unexecuted transactions with their expected effects digests
952        self.execution_scheduler
953            .enqueue_transactions(unexecuted_txns, &self.epoch_store);
954
955        (unexecuted_tx_digests, unexecuted_expected_fx_digests)
956    }
957
958    // Execute the change epoch txn
959    #[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        // Ordinarily we would assert that the change epoch txn has not been executed
973        // yet. However, during crash recovery, it is possible that we already
974        // passed this point and the txn has been executed. You can uncomment
975        // this assert if you are debugging a problem related to reconfig. If
976        // you hit this assert and it is not because of crash-recovery,
977        // it may indicate a bug in the checkpoint executor.
978        //
979        //     if self
980        //         .transaction_cache_reader
981        //         .get_executed_effects(change_epoch_tx.digest())
982        //         .is_some()
983        //     {
984        //         fatal!(
985        //             "end of epoch txn must not have been executed: {:?}",
986        //             change_epoch_tx.digest()
987        //         );
988        //     }
989
990        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    // Increment the highest executed checkpoint watermark
1025    #[instrument(level = "debug", skip_all)]
1026    fn bump_highest_executed_checkpoint(&self, checkpoint: &VerifiedCheckpoint) {
1027        // Ensure that we are not skipping checkpoints at any point
1028        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    /// Helper to broadcast checkpoint summary and data if the
1053    /// channels are set.
1054    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                // Reconstruct checkpoint data if needed (rare case: data_sender configured but
1064                // checkpoint_data_enabled is false)
1065                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    /// If configured, commit the pending index updates for the provided
1085    /// checkpoint
1086    #[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    // Extract randomness rounds from the checkpoint version-specific data (if
1096    // available). Otherwise, extract randomness rounds from the first
1097    // transaction in the checkpoint
1098    #[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            // With version-specific data, randomness rounds are stored in checkpoint
1109            // summary.
1110            version_specific_data.into_v1().randomness_rounds
1111        } else {
1112            // Before version-specific data, checkpoint batching must be disabled. In this
1113            // case, randomness state update tx must be first if it exists,
1114            // because all other transactions in a checkpoint that includes a
1115            // randomness state update are causally dependent on it.
1116            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}