1use std::{sync::Arc, time::Instant};
8
9use async_broadcast::{InactiveReceiver, Sender};
10use chrono::Utc;
11use committable::Committable;
12use hotshot_contract_adapter::light_client::derive_signed_state_digest;
13use hotshot_types::{
14 consensus::OuterConsensus,
15 data::{EpochNumber, Leaf2, QuorumProposalWrapper, VidDisperseShare, ViewNumber},
16 drb::INITIAL_DRB_RESULT,
17 epoch_membership::{EpochMembership, EpochMembershipCoordinator},
18 event::{Event, EventType},
19 message::{Proposal, UpgradeLock},
20 simple_vote::{EpochRootQuorumVote2, LightClientStateUpdateVote2, QuorumData2, QuorumVote2},
21 stake_table::HSStakeTable,
22 storage_metrics::StorageMetricsValue,
23 traits::{
24 ValidatedState,
25 block_contents::BlockHeader,
26 election::Membership,
27 node_implementation::{NodeImplementation, NodeType},
28 signature_key::{
29 LCV2StateSignatureKey, LCV3StateSignatureKey, SignatureKey, StateSignatureKey,
30 },
31 storage::Storage,
32 },
33 utils::{epoch_from_block_number, is_epoch_transition, is_last_block, is_transition_block},
34 vote::HasViewNumber,
35};
36use hotshot_utils::anytrace::*;
37use tracing::instrument;
38use versions::EPOCH_VERSION;
39
40use super::QuorumVoteTaskState;
41use crate::{
42 events::HotShotEvent,
43 helpers::{
44 LeafChainTraversalOutcome, broadcast_event, decide_from_proposal, decide_from_proposal_2,
45 fetch_proposal, handle_drb_result,
46 },
47};
48
49async fn store_drb_result<TYPES: NodeType, I: NodeImplementation<TYPES>>(
51 task_state: &mut QuorumVoteTaskState<TYPES, I>,
52 decided_leaf: &Leaf2<TYPES>,
53) -> Result<()> {
54 if task_state.epoch_height == 0 {
55 tracing::info!("Epoch height is 0, skipping DRB storage.");
56 return Ok(());
57 }
58
59 let decided_block_number = decided_leaf.block_header().block_number();
60 let current_epoch_number = EpochNumber::new(epoch_from_block_number(
61 decided_block_number,
62 task_state.epoch_height,
63 ));
64 if is_transition_block(decided_block_number, task_state.epoch_height) {
66 if let Some(result) = decided_leaf.next_drb_result {
67 handle_drb_result::<TYPES, I>(
70 task_state.membership.membership(),
71 current_epoch_number + 1,
72 &task_state.storage,
73 result,
74 )
75 .await;
76 } else {
77 bail!("The last block of the epoch is decided but doesn't contain a DRB result.");
78 }
79 }
80 Ok(())
81}
82
83#[instrument(skip_all, fields(id = task_state.id, view = *proposal.view_number()))]
85pub(crate) async fn handle_quorum_proposal_validated<
86 TYPES: NodeType,
87 I: NodeImplementation<TYPES>,
88>(
89 proposal: &QuorumProposalWrapper<TYPES>,
90 task_state: &mut QuorumVoteTaskState<TYPES, I>,
91 event_sender: &Sender<Arc<HotShotEvent<TYPES>>>,
92) -> Result<()> {
93 let version = task_state.upgrade_lock.version(proposal.view_number())?;
94
95 let LeafChainTraversalOutcome {
96 new_locked_view_number,
97 new_decided_view_number,
98 committing_qc,
99 deciding_qc,
100 leaf_views,
101 included_txns,
102 decided_upgrade_cert,
103 } = if version >= EPOCH_VERSION {
104 if !is_last_block(
107 proposal.block_header().block_number(),
108 task_state.epoch_height,
109 ) {
110 decide_from_proposal_2::<TYPES, I>(
111 proposal,
112 OuterConsensus::new(Arc::clone(&task_state.consensus.inner_consensus)),
113 &task_state.upgrade_lock,
114 &task_state.public_key,
115 version >= EPOCH_VERSION,
116 &task_state.membership,
117 &task_state.storage,
118 )
119 .await
120 } else {
121 LeafChainTraversalOutcome::default()
122 }
123 } else {
124 decide_from_proposal::<TYPES, I>(
125 proposal,
126 OuterConsensus::new(Arc::clone(&task_state.consensus.inner_consensus)),
127 &task_state.upgrade_lock,
128 &task_state.public_key,
129 version >= EPOCH_VERSION,
130 &task_state.membership,
131 &task_state.storage,
132 task_state.epoch_height,
133 )
134 .await
135 };
136
137 if let (Some(cert), Some(_)) = (decided_upgrade_cert.clone(), new_decided_view_number) {
138 task_state
139 .upgrade_lock
140 .set_decided_upgrade_cert(cert.clone());
141 if cert.data.new_version >= EPOCH_VERSION
142 && task_state.upgrade_lock.upgrade().base < EPOCH_VERSION
143 {
144 let epoch_height = task_state.consensus.read().await.epoch_height;
145 let first_epoch_number = EpochNumber::new(epoch_from_block_number(
146 proposal.block_header().block_number(),
147 epoch_height,
148 ));
149
150 tracing::debug!("Calling set_first_epoch for epoch {first_epoch_number:?}");
151 task_state
152 .membership
153 .membership()
154 .set_first_epoch(first_epoch_number, INITIAL_DRB_RESULT);
155
156 broadcast_event(
157 Arc::new(HotShotEvent::SetFirstEpoch(
158 cert.data.new_version_first_view,
159 first_epoch_number,
160 )),
161 event_sender,
162 )
163 .await;
164 }
165
166 for da_committee in &task_state.da_committees {
167 if cert.data.new_version >= da_committee.start_version {
168 task_state.membership.membership().add_da_committee(
169 da_committee.start_epoch.into(),
170 da_committee.committee.clone(),
171 );
172 }
173 }
174
175 let _ = task_state
176 .storage
177 .update_decided_upgrade_certificate(Some(cert.clone()))
178 .await;
179 }
180
181 let mut consensus_writer = task_state.consensus.write().await;
182 if let Some(locked_view_number) = new_locked_view_number {
183 let _ = consensus_writer.update_locked_view(locked_view_number);
184 }
185
186 #[allow(clippy::cast_precision_loss)]
187 if let Some(decided_view_number) = new_decided_view_number {
188 let old_decided_view = consensus_writer.last_decided_view();
191 consensus_writer.collect_garbage(old_decided_view, decided_view_number);
192
193 consensus_writer
195 .update_last_decided_view(decided_view_number)
196 .context(|e| {
197 warn!("`update_last_decided_view` failed; this should never happen. Error: {e}")
198 })?;
199
200 consensus_writer
201 .metrics
202 .last_decided_time
203 .set(Utc::now().timestamp().try_into().unwrap());
204 consensus_writer.metrics.invalid_qc.set(0);
205 consensus_writer
206 .metrics
207 .last_decided_view
208 .set(usize::try_from(consensus_writer.last_decided_view().u64()).unwrap());
209 let cur_number_of_views_per_decide_event =
210 *proposal.view_number() - consensus_writer.last_decided_view().u64();
211 consensus_writer
212 .metrics
213 .number_of_views_per_decide_event
214 .add_point(cur_number_of_views_per_decide_event as f64);
215 for leaf in leaf_views.iter().rev() {
216 consensus_writer
217 .update_participation_from_qc(&leaf.leaf.justify_qc(), &task_state.membership)?;
218 }
219
220 drop(consensus_writer);
222
223 for leaf_info in &leaf_views {
224 tracing::info!(
225 "Sending decide for view {:?} at height {:?}",
226 leaf_info.leaf.view_number(),
227 leaf_info.leaf.block_header().block_number(),
228 );
229 }
230
231 broadcast_event(
232 Arc::new(HotShotEvent::LeavesDecided(
233 leaf_views
234 .iter()
235 .map(|leaf_info| leaf_info.leaf.clone())
236 .collect(),
237 )),
238 event_sender,
239 )
240 .await;
241
242 let committing_qc = Arc::new(committing_qc.unwrap());
245 broadcast_event(
246 Event {
247 view_number: decided_view_number,
248 event: EventType::Decide {
249 leaf_chain: Arc::new(leaf_views.clone()),
250 committing_qc: committing_qc.clone(),
251 deciding_qc: deciding_qc.map(Arc::new),
252 block_size: included_txns.map(|txns| txns.len().try_into().unwrap()),
253 },
254 },
255 &task_state.output_event_stream,
256 )
257 .await;
258
259 tracing::debug!(
260 "Successfully sent decide event, leaf views: {:?}, leaf views len: {:?}, qc view: {:?}",
261 decided_view_number,
262 leaf_views.len(),
263 committing_qc.view_number()
264 );
265
266 if version >= EPOCH_VERSION {
267 for leaf_view in leaf_views {
268 store_drb_result(task_state, &leaf_view.leaf).await?;
269 }
270 }
271 }
272
273 Ok(())
274}
275
276#[instrument(skip_all, target = "VoteDependencyHandle", fields(view = *view_number))]
278#[allow(clippy::too_many_arguments)]
279pub(crate) async fn update_shared_state<TYPES: NodeType>(
280 consensus: OuterConsensus<TYPES>,
281 sender: Sender<Arc<HotShotEvent<TYPES>>>,
282 receiver: InactiveReceiver<Arc<HotShotEvent<TYPES>>>,
283 membership: EpochMembershipCoordinator<TYPES>,
284 public_key: TYPES::SignatureKey,
285 private_key: <TYPES::SignatureKey as SignatureKey>::PrivateKey,
286 upgrade_lock: UpgradeLock<TYPES>,
287 view_number: ViewNumber,
288 instance_state: Arc<TYPES::InstanceState>,
289 proposed_leaf: &Leaf2<TYPES>,
290 vid_share: &Proposal<TYPES, VidDisperseShare<TYPES>>,
291 parent_view_number: Option<ViewNumber>,
292 epoch_height: u64,
293) -> Result<()> {
294 let justify_qc = &proposed_leaf.justify_qc();
295
296 let consensus_reader = consensus.read().await;
297 let mut maybe_validated_view = parent_view_number.and_then(|view_number| {
300 consensus_reader
301 .validated_state_map()
302 .get(&view_number)
303 .cloned()
304 });
305
306 let mut maybe_parent = consensus_reader
308 .saved_leaves()
309 .get(&justify_qc.data.leaf_commit)
310 .cloned();
311
312 drop(consensus_reader);
313
314 maybe_parent = match maybe_parent {
315 Some(p) => Some(p),
316 None => {
317 match fetch_proposal(
318 justify_qc,
319 sender.clone(),
320 receiver.activate_cloned(),
321 membership.clone(),
322 OuterConsensus::new(Arc::clone(&consensus.inner_consensus)),
323 public_key.clone(),
324 private_key.clone(),
325 &upgrade_lock,
326 epoch_height,
327 )
328 .await
329 .ok()
330 {
331 Some((leaf, view)) => {
332 maybe_validated_view = Some(view);
333 Some(leaf)
334 },
335 None => None,
336 }
337 },
338 };
339
340 let parent = maybe_parent.context(info!(
341 "Proposal's parent missing from storage with commitment: {:?}, proposal view {}",
342 justify_qc.data.leaf_commit,
343 proposed_leaf.view_number(),
344 ))?;
345
346 let Some(validated_view) = maybe_validated_view else {
347 bail!("Failed to fetch view for parent, parent view {parent_view_number:?}");
348 };
349
350 let (Some(parent_state), _) = validated_view.state_and_delta() else {
351 bail!("Parent state not found! Consensus internally inconsistent");
352 };
353
354 let version = upgrade_lock.version(view_number)?;
355
356 let now = Instant::now();
357 let (validated_state, state_delta) = parent_state
358 .validate_and_apply_header(
359 &instance_state,
360 &parent,
361 &proposed_leaf.block_header().clone(),
362 vid_share.data.payload_byte_len(),
363 version,
364 *view_number,
365 )
366 .await
367 .wrap()
368 .context(warn!("Block header doesn't extend the proposal!"))?;
369 let validation_duration = now.elapsed();
370 tracing::debug!("Validation time: {validation_duration:?}");
371
372 let now = Instant::now();
373 let mut consensus_writer = consensus.write().await;
375
376 if let Err(e) = consensus_writer.update_leaf(
377 proposed_leaf.clone(),
378 Arc::new(validated_state),
379 Some(Arc::new(state_delta)),
380 ) {
381 tracing::trace!("{e:?}");
382 }
383 let update_leaf_duration = now.elapsed();
384
385 consensus_writer
386 .metrics
387 .validate_and_apply_header_duration
388 .add_point(validation_duration.as_secs_f64());
389 consensus_writer
390 .metrics
391 .update_leaf_duration
392 .add_point(update_leaf_duration.as_secs_f64());
393 drop(consensus_writer);
394 tracing::debug!("update_leaf time: {update_leaf_duration:?}");
395
396 Ok(())
397}
398
399#[instrument(skip_all, fields(name = "Submit quorum vote", level = "error"))]
401#[allow(clippy::too_many_arguments)]
402pub(crate) async fn submit_vote<TYPES: NodeType, I: NodeImplementation<TYPES>>(
403 sender: Sender<Arc<HotShotEvent<TYPES>>>,
404 membership: EpochMembership<TYPES>,
405 public_key: TYPES::SignatureKey,
406 private_key: <TYPES::SignatureKey as SignatureKey>::PrivateKey,
407 upgrade_lock: UpgradeLock<TYPES>,
408 view_number: ViewNumber,
409 storage: I::Storage,
410 storage_metrics: Arc<StorageMetricsValue>,
411 leaf: Leaf2<TYPES>,
412 vid_share: Proposal<TYPES, VidDisperseShare<TYPES>>,
413 extended_vote: bool,
414 epoch_root_vote: bool,
415 epoch_height: u64,
416 state_private_key: &<TYPES::StateSignatureKey as StateSignatureKey>::StatePrivateKey,
417 stake_table_capacity: usize,
418) -> Result<()> {
419 let committee_member_in_current_epoch = membership.has_stake(&public_key);
420 let committee_member_in_next_epoch = leaf.with_epoch
423 && is_epoch_transition(leaf.height(), epoch_height)
424 && membership.next_epoch_stake_table()?.has_stake(&public_key);
425
426 ensure!(
427 committee_member_in_current_epoch || committee_member_in_next_epoch,
428 info!("We were not chosen for quorum committee on {view_number}")
429 );
430
431 let height = if membership.epoch().is_some() {
432 Some(leaf.height())
433 } else {
434 None
435 };
436
437 let vote = QuorumVote2::<TYPES>::create_signed_vote(
439 QuorumData2 {
440 leaf_commit: leaf.commit(),
441 epoch: membership.epoch(),
442 block_number: height,
443 },
444 view_number,
445 &public_key,
446 &private_key,
447 &upgrade_lock,
448 )
449 .wrap()
450 .context(error!("Failed to sign vote. This should never happen."))?;
451 let now = Instant::now();
452 storage
454 .append_vid(&vid_share)
455 .await
456 .wrap()
457 .context(error!("Failed to store VID share"))?;
458 let append_vid_duration = now.elapsed();
459 storage_metrics
460 .append_vid_duration
461 .add_point(append_vid_duration.as_secs_f64());
462 tracing::debug!("append_vid_general time: {append_vid_duration:?}");
463
464 let epoch_enabled = upgrade_lock.epochs_enabled(view_number);
467 if extended_vote && epoch_enabled {
468 tracing::debug!("sending extended vote to everybody",);
469 broadcast_event(
470 Arc::new(HotShotEvent::ExtendedQuorumVoteSend(vote)),
471 &sender,
472 )
473 .await;
474 } else if epoch_root_vote && epoch_enabled {
475 tracing::debug!(
476 "sending epoch root vote to next quorum leader {:?}",
477 vote.view_number() + 1
478 );
479 let light_client_state = leaf
480 .block_header()
481 .get_light_client_state(view_number)
482 .wrap()
483 .context(error!("Failed to generate light client state"))?;
484 let next_stake_table =
485 HSStakeTable::from_iter(membership.next_epoch_stake_table()?.stake_table());
486 let next_stake_table_state = next_stake_table
487 .commitment(stake_table_capacity)
488 .wrap()
489 .context(error!("Failed to compute stake table commitment"))?;
490 let v2_signature = <TYPES::StateSignatureKey as LCV2StateSignatureKey>::sign_state(
492 state_private_key,
493 &light_client_state,
494 &next_stake_table_state,
495 )
496 .wrap()
497 .context(error!("Failed to sign the light client state"))?;
498 let auth_root = leaf
499 .block_header()
500 .auth_root()
501 .wrap()
502 .context(error!(format!(
503 "Failed to get auth root for light client state certificate. view={view_number}"
504 )))?;
505 let signed_state_digest =
506 derive_signed_state_digest(&light_client_state, &next_stake_table_state, &auth_root);
507 let signature = <TYPES::StateSignatureKey as LCV3StateSignatureKey>::sign_state(
508 state_private_key,
509 signed_state_digest,
510 )
511 .wrap()
512 .context(error!("Failed to sign the light client state"))?;
513 let state_vote = LightClientStateUpdateVote2 {
514 epoch: EpochNumber::new(epoch_from_block_number(leaf.height(), epoch_height)),
515 light_client_state,
516 next_stake_table_state,
517 signature,
518 v2_signature,
519 auth_root,
520 signed_state_digest,
521 };
522 broadcast_event(
523 Arc::new(HotShotEvent::EpochRootQuorumVoteSend(
524 EpochRootQuorumVote2 { vote, state_vote },
525 )),
526 &sender,
527 )
528 .await;
529 } else {
530 tracing::debug!(
531 "sending vote to next quorum leader {:?}",
532 vote.view_number() + 1
533 );
534 broadcast_event(Arc::new(HotShotEvent::QuorumVoteSend(vote)), &sender).await;
535 }
536
537 Ok(())
538}