iota_core/checkpoints/
checkpoint_output.rs1use 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 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 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}