Skip to main content

hotshot_task_impls/quorum_vote/
handlers.rs

1// Copyright (c) 2021-2024 Espresso Systems (espressosys.com)
2// This file is part of the HotShot repository.
3
4// You should have received a copy of the MIT License
5// along with the HotShot repository. If not, see <https://mit-license.org/>.
6
7use 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
49/// Store the DRB result for the next epoch if we received it in a decided leaf.
50async 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    // Skip storing the received result if this is not the transition block.
65    if is_transition_block(decided_block_number, task_state.epoch_height) {
66        if let Some(result) = decided_leaf.next_drb_result {
67            // We don't need to check value existence and consistency because it should be
68            // impossible to decide on a block with different DRB results.
69            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/// Handles the `QuorumProposalValidated` event.
84#[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        // Skip the decide rule for the last block of the epoch.  This is so
105        // that we do not decide the block with epoch_height -2 before we enter the new epoch
106        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        // Bring in the cleanup crew. When a new decide is indeed valid, we need to clear out old memory.
189
190        let old_decided_view = consensus_writer.last_decided_view();
191        consensus_writer.collect_garbage(old_decided_view, decided_view_number);
192
193        // Set the new decided view.
194        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        // We don't need to hold this while we broadcast
221        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        // Send an update to everyone saying that we've reached a decide. The committing QC is never
243        // none if we've reached a new decide, so this is safe to unwrap.
244        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/// Updates the shared consensus state with the new voting data.
277#[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    // Try to find the validated view within the validated state map. This will be present
298    // if we have the saved leaf, but if not we'll get it when we fetch_proposal.
299    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    // Justify qc's leaf commitment should be the same as the parent's leaf commitment.
307    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    // Now that we've rounded everyone up, we need to update the shared state
374    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/// Submits the `QuorumVoteSend` event if all the dependencies are met.
400#[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    // If the proposed leaf is for the last block in the epoch and the node is part of the quorum committee
421    // in the next epoch, the node should vote to achieve the double quorum.
422    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    // Create and send the vote.
438    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    // Add to the storage.
453    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    // Make epoch root vote
465
466    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        // We are still providing LCV2 state signatures for backward compatibility
491        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}