Skip to main content

iota_core/checkpoints/
checkpoint_output.rs

1// Copyright (c) Mysten Labs, Inc.
2// Modifications Copyright (c) 2024 IOTA Stiftung
3// SPDX-License-Identifier: Apache-2.0
4
5use std::sync::Arc;
6
7use async_trait::async_trait;
8use iota_sdk_types::{CheckpointContents, CheckpointSummary};
9use iota_types::{
10    base_types::AuthorityName,
11    error::IotaResult,
12    messages_checkpoint::{
13        CertifiedCheckpointSummary, CheckpointSignatureMessage, CheckpointSummaryExt,
14        SignedCheckpointSummary, VerifiedCheckpoint,
15    },
16    messages_consensus::ConsensusTransaction,
17};
18use tracing::{debug, info, instrument, trace};
19
20use super::{CheckpointMetrics, CheckpointStore};
21use crate::{
22    authority::{StableSyncAuthoritySigner, authority_per_epoch_store::AuthorityPerEpochStore},
23    consensus_adapter::SubmitToConsensus,
24    epoch::reconfiguration::ReconfigurationInitiator,
25};
26
27const REPORT_END_OF_EPOCH_MARGIN_MS: u64 = 2000;
28const MIN_CHECKPOINTS_BETWEEN_REPORTS: u64 = 1000;
29const MAX_CHECKPOINT_LAG_FOR_REPORT: u64 = 100;
30#[async_trait]
31pub trait CheckpointOutput: Sync + Send + 'static {
32    async fn checkpoint_created(
33        &self,
34        summary: &CheckpointSummary,
35        contents: &CheckpointContents,
36        epoch_store: &Arc<AuthorityPerEpochStore>,
37        checkpoint_store: &Arc<CheckpointStore>,
38    ) -> IotaResult;
39}
40
41#[async_trait]
42pub trait CertifiedCheckpointOutput: Sync + Send + 'static {
43    async fn certified_checkpoint_created(
44        &self,
45        summary: &CertifiedCheckpointSummary,
46    ) -> IotaResult;
47}
48
49pub struct SubmitCheckpointToConsensus<T> {
50    pub sender: T,
51    pub signer: StableSyncAuthoritySigner,
52    pub authority: AuthorityName,
53    pub next_reconfiguration_timestamp_ms: u64,
54    pub metrics: Arc<CheckpointMetrics>,
55}
56
57pub struct LogCheckpointOutput;
58
59impl LogCheckpointOutput {
60    pub fn boxed() -> Box<dyn CheckpointOutput> {
61        Box::new(Self)
62    }
63
64    pub fn boxed_certified() -> Box<dyn CertifiedCheckpointOutput> {
65        Box::new(Self)
66    }
67}
68
69#[async_trait]
70impl<T: SubmitToConsensus + ReconfigurationInitiator> CheckpointOutput
71    for SubmitCheckpointToConsensus<T>
72{
73    #[instrument(level = "debug", skip_all)]
74    async fn checkpoint_created(
75        &self,
76        summary: &CheckpointSummary,
77        contents: &CheckpointContents,
78        epoch_store: &Arc<AuthorityPerEpochStore>,
79        checkpoint_store: &Arc<CheckpointStore>,
80    ) -> IotaResult {
81        LogCheckpointOutput
82            .checkpoint_created(summary, contents, epoch_store, checkpoint_store)
83            .await?;
84
85        let checkpoint_timestamp = summary.timestamp_ms;
86        let checkpoint_seq = summary.sequence_number;
87        self.metrics.checkpoint_creation_latency.observe(
88            summary
89                .timestamp()
90                .elapsed()
91                .unwrap_or_default()
92                .as_secs_f64(),
93        );
94
95        let highest_verified_checkpoint =
96            checkpoint_store.get_highest_verified_checkpoint_seq_number()?;
97
98        if Some(checkpoint_seq) > highest_verified_checkpoint {
99            debug!(
100                "Sending checkpoint signature at sequence {checkpoint_seq} to consensus, timestamp {checkpoint_timestamp}.
101                {}ms left till end of epoch at timestamp {}",
102                self.next_reconfiguration_timestamp_ms.saturating_sub(checkpoint_timestamp), self.next_reconfiguration_timestamp_ms
103            );
104
105            let summary = SignedCheckpointSummary::new(
106                epoch_store.epoch(),
107                summary.clone(),
108                &*self.signer,
109                self.authority,
110            );
111
112            let message = CheckpointSignatureMessage { summary };
113            let transaction = ConsensusTransaction::new_checkpoint_signature_message(message);
114            self.sender
115                .submit_to_consensus(&[transaction], epoch_store)?;
116            self.metrics
117                .last_sent_checkpoint_signature
118                .set(checkpoint_seq as i64);
119        } else {
120            debug!(
121                "Checkpoint at sequence {checkpoint_seq} is already certified, skipping signature submission to consensus",
122            );
123            self.metrics
124                .last_skipped_checkpoint_signature_submission
125                .set(checkpoint_seq as i64);
126        }
127
128        // If `calculate_validator_scores` is enabled in protocol config, we also send
129        // misbehavior reports to consensus at this point. Misbehavior reports
130        // containing proofs of misbehaviour can be sent whenever the misbehavior is
131        // detected, but we choose to send the ones that include only unprovable counts
132        // at this point, due to periodicity reasons and to ensure a (approximate)
133        // synchronization with the score updates.
134        //
135        // Reports are rate-limited: only sent when metrics have changed (different
136        // summaries) and at least 1000 checkpoints have passed since the last report.
137        // We also require that the checkpoint for which we want to send the report is
138        // at most 100 checkpoints behind the highest verified checkpoint, to avoid
139        // sending reports during resync.
140        //
141        // Additionally to these periodic reports, we also send a report when the epoch
142        // is coming to an end. Since `close_epoch` is called according to local clocks,
143        // we use an analogous rule for the last reports, requiring that the checkpoint
144        // is close to the next reconfiguration timestamp.
145        if epoch_store.protocol_config().calculate_validator_scores() {
146            let mut rl = epoch_store.misbehavior_monitor.rate_limit();
147
148            let should_send_last_report = checkpoint_timestamp
149                >= self
150                    .next_reconfiguration_timestamp_ms
151                    .saturating_sub(REPORT_END_OF_EPOCH_MARGIN_MS)
152                && !rl.has_sent_end_of_epoch_report;
153
154            if (checkpoint_seq.saturating_sub(rl.last_report_checkpoint_seq)
155                >= MIN_CHECKPOINTS_BETWEEN_REPORTS
156                && Some(checkpoint_seq + MAX_CHECKPOINT_LAG_FOR_REPORT)
157                    >= highest_verified_checkpoint)
158                || should_send_last_report
159            {
160                let misbehavior_report = epoch_store
161                    .misbehavior_monitor
162                    .generate_report(checkpoint_seq);
163                let new_report_summary = misbehavior_report.summary();
164                if new_report_summary != rl.last_report_summary || should_send_last_report {
165                    let transaction =
166                        ConsensusTransaction::new_misbehavior_report(misbehavior_report);
167                    info!(?transaction, "submitting misbehavior report to consensus");
168                    self.sender
169                        .submit_to_consensus(&[transaction], epoch_store)?;
170                    rl.last_report_summary = new_report_summary;
171                    rl.last_report_checkpoint_seq = checkpoint_seq;
172                    if should_send_last_report {
173                        rl.has_sent_end_of_epoch_report = true;
174                    }
175                }
176            }
177        }
178
179        if checkpoint_timestamp >= self.next_reconfiguration_timestamp_ms {
180            // close_epoch is ok if called multiple times
181            self.sender.close_epoch(epoch_store);
182        }
183        Ok(())
184    }
185}
186
187#[async_trait]
188impl CheckpointOutput for LogCheckpointOutput {
189    async fn checkpoint_created(
190        &self,
191        summary: &CheckpointSummary,
192        contents: &CheckpointContents,
193        _epoch_store: &Arc<AuthorityPerEpochStore>,
194        _checkpoint_store: &Arc<CheckpointStore>,
195    ) -> IotaResult {
196        trace!(
197            "Including following transactions in checkpoint {}: {:?}",
198            summary.sequence_number, contents
199        );
200        info!(
201            "Creating checkpoint {:?} at epoch {}, sequence {}, previous digest {:?}, transactions count {}, content digest {:?}, end_of_epoch_data {:?}",
202            summary.digest(),
203            summary.epoch,
204            summary.sequence_number,
205            summary.previous_digest,
206            contents.len(),
207            summary.contents_digest,
208            summary.end_of_epoch_data,
209        );
210
211        Ok(())
212    }
213}
214
215#[async_trait]
216impl CertifiedCheckpointOutput for LogCheckpointOutput {
217    async fn certified_checkpoint_created(
218        &self,
219        summary: &CertifiedCheckpointSummary,
220    ) -> IotaResult {
221        debug!(
222            "Certified checkpoint with sequence {} and digest {}",
223            summary.sequence_number,
224            summary.digest()
225        );
226        Ok(())
227    }
228}
229
230pub struct SendCheckpointToStateSync {
231    handle: iota_network::state_sync::Handle,
232}
233
234impl SendCheckpointToStateSync {
235    pub fn new(handle: iota_network::state_sync::Handle) -> Self {
236        Self { handle }
237    }
238}
239
240#[async_trait]
241impl CertifiedCheckpointOutput for SendCheckpointToStateSync {
242    #[instrument(level = "trace", name = "checkpoint_created_from_consensus", skip_all)]
243    async fn certified_checkpoint_created(
244        &self,
245        summary: &CertifiedCheckpointSummary,
246    ) -> IotaResult {
247        debug!(
248            "Certified checkpoint with sequence {} and digest {}",
249            summary.sequence_number,
250            summary.digest()
251        );
252        self.handle
253            .send_checkpoint(VerifiedCheckpoint::new_unchecked(summary.to_owned()))
254            .await;
255
256        Ok(())
257    }
258}