1use std::{
2 cmp::max,
3 collections::{BTreeMap, BTreeSet},
4 marker::PhantomData,
5 sync::Arc,
6};
7
8use committable::{Commitment, CommitmentBoundsArkless, Committable};
9use hotshot::traits::BlockPayload;
10use hotshot_contract_adapter::light_client::derive_signed_state_digest;
11use hotshot_types::{
12 data::{
13 BlockNumber, EpochNumber, Leaf2, VidCommitment, VidCommitment2, VidDisperseShare2,
14 ViewChangeEvidence2, ViewNumber,
15 },
16 drb::DrbResult,
17 epoch_membership::EpochMembershipCoordinator,
18 message::{Proposal as SignedProposal, UpgradeLock},
19 simple_certificate::{
20 LightClientStateUpdateCertificateV2, QuorumCertificate2, TimeoutCertificate2,
21 check_qc_state_cert_correspondence,
22 },
23 simple_vote::{
24 HasEpoch, LightClientStateUpdateVote2, QuorumData2, SimpleVote, TimeoutData2, TimeoutVote2,
25 Vote2Data,
26 },
27 stake_table::HSStakeTable,
28 traits::{
29 block_contents::BlockHeader,
30 node_implementation::NodeType,
31 signature_key::{
32 LCV2StateSignatureKey, LCV3StateSignatureKey, SignatureKey, StateSignatureKey,
33 },
34 },
35 utils::{epoch_from_block_number, is_epoch_root, is_epoch_transition, is_last_block},
36 vote::{Certificate, HasViewNumber},
37};
38use hotshot_utils::anytrace;
39use tracing::{debug, info, instrument, warn};
40
41use crate::{
42 block::BlockAndHeaderRequest,
43 cert_verifier::ValidCert,
44 coordinator::{GcScope, VID_RECONSTRUCT_GC_MARGIN},
45 helpers::proposal_commitment,
46 logging::KeyPrefix,
47 message::{
48 CatchupEvidence, Certificate1, Certificate2, EpochChangeMessage, Proposal,
49 ProposalFetchRequest, ProposalMessage, Validated, Vote1, Vote2,
50 },
51 outbox::Outbox,
52 state::{StateRequest, StateResponse},
53 storage::{ActionKind, StorageOutput},
54};
55
56#[derive(Clone, Debug)]
64pub struct PreCutoverSeed<T: NodeType> {
65 pub decided_anchor: Leaf2<T>,
67 pub undecided: Vec<Leaf2<T>>,
69 pub high_qc: Option<QuorumCertificate2<T>>,
72 pub validated_states: BTreeMap<ViewNumber, Arc<T::ValidatedState>>,
74 pub cutover_view: ViewNumber,
78}
79
80#[derive(Eq, PartialEq, Debug, Clone)]
81#[allow(clippy::large_enum_variant)]
82pub enum ConsensusInput<T: NodeType> {
83 BlockBuilt {
84 view: ViewNumber,
85 epoch: EpochNumber,
86 payload: T::BlockPayload,
87 metadata: <T::BlockPayload as BlockPayload<T>>::Metadata,
88 payload_commitment: VidCommitment,
89 },
90 BlockReconstructed(ViewNumber, VidCommitment2),
91 Certificate1(ValidCert<Certificate1<T>>),
92 Certificate2(ValidCert<Certificate2<T>>),
93 AdvanceView(ValidCert<Certificate1<T>>),
97 EpochRootCertificates {
101 cert1: ValidCert<Certificate1<T>>,
102 state_cert: LightClientStateUpdateCertificateV2<T>,
103 },
104 EpochChange(EpochChangeMessage<T, Validated>),
105 HeaderCreated(ViewNumber, Commitment<Leaf2<T>>, T::BlockHeader),
106 Proposal(T::SignatureKey, ProposalMessage<T, Validated>),
110 VidShare(VidDisperseShare2<T>),
112 FetchedProposal(ProposalMessage<T, Validated>),
113 StateValidated(StateResponse<T>),
114 StateValidationFailed(StateResponse<T>),
115 Stored(StorageOutput<T>),
116 Timeout(ViewNumber, EpochNumber),
117 TimeoutCertificate(ValidCert<TimeoutCertificate2<T>>),
118 TimeoutOneHonest(ViewNumber, EpochNumber),
119 VidDisperseCreated(ViewNumber, VidCommitment2),
120 DrbResult(EpochNumber, DrbResult),
121}
122
123#[derive(Eq, PartialEq, Debug, Clone)]
124pub enum ConsensusOutput<T: NodeType> {
125 RequestBlockAndHeader(BlockAndHeaderRequest<T>),
126 RequestState(StateRequest<T>),
127 RequestDrbResult(EpochNumber),
128 RecordAction(ViewNumber, Option<EpochNumber>, ActionKind),
129 PersistProposal(SignedProposal<T, Proposal<T>>),
130 SendProposal(SignedProposal<T, Proposal<T>>),
131 SendTimeoutVote(TimeoutVote2<T>, Option<CatchupEvidence<T>>),
132 SendVote1(Vote1<T>),
133 SendVote2(Vote2<T>),
134 PersistHighQc(Certificate1<T>),
136 SendTimeoutCertificate(TimeoutCertificate2<T>, ViewNumber, EpochNumber),
137 SendCertificate1(Certificate1<T>),
138 SendCertificate2(Certificate2<T>),
141 SendEpochChange(EpochChangeMessage<T, Validated>),
142 RequestVidDisperse {
143 view: ViewNumber,
144 epoch: EpochNumber,
145 payload: T::BlockPayload,
146 metadata: <T::BlockPayload as BlockPayload<T>>::Metadata,
147 payload_commitment: VidCommitment2,
148 },
149 LeafDecided {
150 leaves: Vec<Leaf2<T>>,
151 cert1: Certificate1<T>,
154 cert2: Option<Certificate2<T>>,
155 vid_shares: Vec<Option<SignedProposal<T, VidDisperseShare2<T>>>>,
156 },
157 LockUpdated(Certificate2<T>),
158 ViewChanged(ViewNumber, EpochNumber),
159 ViewTimedOut(ViewNumber),
161 ProposalPaired {
163 proposal: SignedProposal<T, Proposal<T>>,
164 vid_share: VidDisperseShare2<T>,
165 },
166 ProposalValidated {
167 proposal: SignedProposal<T, Proposal<T>>,
168 sender: T::SignatureKey,
169 },
170 RequestMissingProposal {
171 view: ViewNumber,
172 leaf_commit: Commitment<Leaf2<T>>,
173 },
174 BlockPayloadReconstructed {
179 view: ViewNumber,
180 header: T::BlockHeader,
181 payload: T::BlockPayload,
182 },
183 BroadcastVidShare(VidDisperseShare2<T>),
186}
187
188type UnpairedProposals<T> = BTreeMap<
189 (ViewNumber, VidCommitment2),
190 (<T as NodeType>::SignatureKey, ProposalMessage<T, Validated>),
191>;
192
193type UnpairedVidShares<T> = BTreeMap<(ViewNumber, VidCommitment2), VidDisperseShare2<T>>;
194
195pub(crate) const DECIDE_BUFFER: u64 = 20;
198
199const _: () = assert!(DECIDE_BUFFER >= VID_RECONSTRUCT_GC_MARGIN);
201
202pub struct Consensus<T: NodeType> {
203 proposals: BTreeMap<ViewNumber, Proposal<T>>,
204 signed_proposals: BTreeMap<ViewNumber, SignedProposal<T, Proposal<T>>>,
205 proposed_views: BTreeSet<ViewNumber>,
206 vid_shares: BTreeMap<ViewNumber, VidDisperseShare2<T>>,
207 unpaired_proposals: UnpairedProposals<T>,
208 unpaired_vid_shares: UnpairedVidShares<T>,
209 states_verified: BTreeMap<ViewNumber, Commitment<Leaf2<T>>>,
210 blocks_reconstructed: BTreeSet<(ViewNumber, VidCommitment2)>,
211 blocks: BTreeMap<(ViewNumber, VidCommitment2), T::BlockPayload>,
212 certs: BTreeMap<ViewNumber, Certificate1<T>>,
213 certs2: BTreeMap<ViewNumber, Certificate2<T>>,
214 timeout_certs: BTreeMap<ViewNumber, TimeoutCertificate2<T>>,
215 locked_cert: Option<Certificate1<T>>,
216 headers: BTreeMap<(ViewNumber, Commitment<Leaf2<T>>), T::BlockHeader>,
217 leaves: BTreeMap<ViewNumber, Leaf2<T>>,
218 decided_views: BTreeSet<ViewNumber>,
221 decide_floor_view: ViewNumber,
225 last_decided_view: ViewNumber,
226 last_decided_leaf: Leaf2<T>,
227 drb_results: BTreeMap<EpochNumber, DrbResult>,
228
229 voted_1_views: BTreeSet<ViewNumber>,
230 voted_2_views: BTreeSet<ViewNumber>,
231
232 stored_proposals: BTreeMap<ViewNumber, Vec<Commitment<Leaf2<T>>>>,
234 stored_vids: BTreeSet<ViewNumber>,
235 stored_actions: BTreeSet<(ViewNumber, ActionKind)>,
236 requested_actions: BTreeSet<(ViewNumber, ActionKind)>,
237 stored_high_qc: Option<ViewNumber>,
239
240 pending_vote1: BTreeMap<ViewNumber, Vote1<T>>,
244 pending_vote2: BTreeMap<ViewNumber, (Vote2<T>, ViewNumber)>,
245 pending_proposal: BTreeMap<ViewNumber, SignedProposal<T, Proposal<T>>>,
246
247 pre_cutover_views: BTreeSet<ViewNumber>,
249
250 timeout_view: ViewNumber,
251 restart_barred_view: ViewNumber,
255 current_view: ViewNumber,
256 current_epoch: Option<EpochNumber>,
257
258 stake_table_coordinator: EpochMembershipCoordinator<T>,
261
262 public_key: T::SignatureKey,
263 private_key: <T::SignatureKey as SignatureKey>::PrivateKey,
264 state_private_key: <T::StateSignatureKey as StateSignatureKey>::StatePrivateKey,
265 stake_table_capacity: usize,
266 state_certs: BTreeMap<EpochNumber, LightClientStateUpdateCertificateV2<T>>,
267 node_id: KeyPrefix,
268 upgrade_lock: UpgradeLock<T>,
269
270 pub(crate) epoch_height: BlockNumber,
271}
272
273enum Protocol {
275 Abort,
277 Continue,
279}
280
281#[derive(Debug, thiserror::Error)]
283enum SafetyError {
284 #[error(
285 "leaf commitment at locked view does not match locked certificate \
286 locked_commit={locked_commit} proposal_commit={proposal_commit}"
287 )]
288 LockedViewCommitmentMismatch {
289 locked_commit: String,
290 proposal_commit: String,
291 },
292 #[error(
293 "justify qc neither extends nor is newer than the locked certificate \
294 locked_view={locked_view} parent_commit={parent_commit} locked_commit={locked_commit}"
295 )]
296 UnsafeProposal {
297 locked_view: ViewNumber,
298 parent_commit: String,
299 locked_commit: String,
300 },
301 #[error("failed to compute justify qc data commitment: {0}")]
302 JustifyQcCommitment(#[source] anytrace::Error),
303 #[error("failed to compute locked certificate data commitment: {0}")]
304 LockedCertCommitment(#[source] anytrace::Error),
305}
306
307impl<T: NodeType> Consensus<T> {
308 #[allow(clippy::too_many_arguments)]
309 pub fn new<B>(
310 membership_coordinator: EpochMembershipCoordinator<T>,
311 public_key: T::SignatureKey,
312 private_key: <T::SignatureKey as SignatureKey>::PrivateKey,
313 state_private_key: <T::StateSignatureKey as StateSignatureKey>::StatePrivateKey,
314 stake_table_capacity: usize,
315 upgrade_lock: UpgradeLock<T>,
316 genesis_leaf: Leaf2<T>,
317 epoch_height: B,
318 ) -> Self
319 where
320 B: Into<BlockNumber>,
321 {
322 let last_decided_view = genesis_leaf.view_number();
323 Self {
324 proposals: BTreeMap::new(),
325 signed_proposals: BTreeMap::new(),
326 proposed_views: BTreeSet::new(),
327 blocks: BTreeMap::new(),
328 states_verified: BTreeMap::new(),
329 blocks_reconstructed: BTreeSet::new(),
330 certs: BTreeMap::new(),
331 certs2: BTreeMap::new(),
332 timeout_certs: BTreeMap::new(),
333 locked_cert: None,
334 leaves: BTreeMap::new(),
335 decided_views: BTreeSet::from([last_decided_view]),
336 decide_floor_view: ViewNumber::genesis(),
337 last_decided_view,
338 last_decided_leaf: genesis_leaf,
339 headers: BTreeMap::new(),
340 drb_results: BTreeMap::new(),
341 node_id: KeyPrefix::from(&public_key),
342 public_key,
343 timeout_view: ViewNumber::genesis(),
344 restart_barred_view: ViewNumber::genesis(),
345 current_view: ViewNumber::genesis(),
346 current_epoch: None,
347 stake_table_coordinator: membership_coordinator,
348 voted_1_views: BTreeSet::new(),
349 voted_2_views: BTreeSet::new(),
350 stored_proposals: BTreeMap::new(),
351 stored_vids: BTreeSet::new(),
352 stored_actions: BTreeSet::new(),
353 requested_actions: BTreeSet::new(),
354 stored_high_qc: None,
355 pending_vote1: BTreeMap::new(),
356 pending_vote2: BTreeMap::new(),
357 pending_proposal: BTreeMap::new(),
358 pre_cutover_views: BTreeSet::new(),
359 private_key,
360 state_private_key,
361 stake_table_capacity,
362 state_certs: BTreeMap::new(),
363 upgrade_lock,
364 vid_shares: BTreeMap::new(),
365 unpaired_proposals: BTreeMap::new(),
366 unpaired_vid_shares: BTreeMap::new(),
367 epoch_height: epoch_height.into(),
368 }
369 }
370
371 pub fn seed_parent(
385 &mut self,
386 cert1: Certificate1<T>,
387 proposal: Proposal<T>,
388 reconstructed: impl IntoIterator<Item = (ViewNumber, VidCommitment2)>,
389 ) {
390 self.current_epoch = Some(proposal.epoch);
391 self.bump_stored_high_qc(cert1.view_number());
393 self.certs.insert(cert1.view_number(), cert1.clone());
394 self.locked_cert = Some(cert1);
395 self.proposals.insert(proposal.view_number, proposal);
396 for (view, commitment) in reconstructed {
397 self.blocks_reconstructed.insert((view, commitment));
398 }
399 }
400
401 pub fn seed_proposals(&mut self, proposals: impl IntoIterator<Item = Proposal<T>>) {
406 for proposal in proposals {
407 let view = proposal.view_number;
408 self.leaves.insert(view, proposal.clone().into());
409 self.proposals.insert(view, proposal);
410 }
411 }
412
413 pub fn seed_locked_cert(&mut self, cert1: Certificate1<T>) {
419 let view = cert1.view_number();
420 self.bump_stored_high_qc(view);
421 self.certs.entry(view).or_insert_with(|| cert1.clone());
422 if self
423 .locked_cert
424 .as_ref()
425 .is_none_or(|locked| locked.view_number() < view)
426 {
427 self.locked_cert = Some(cert1);
428 }
429 }
430
431 fn bump_stored_high_qc(&mut self, view: ViewNumber) {
433 if self.stored_high_qc.is_none_or(|cur| cur < view) {
434 self.stored_high_qc = Some(view);
435 }
436 }
437
438 fn high_qc_persisted(&self, required: ViewNumber) -> bool {
440 self.stored_high_qc.is_some_and(|stored| stored >= required)
441 }
442
443 fn parent_reconstructed(
449 &self,
450 parent_view: ViewNumber,
451 parent_block_commitment: VidCommitment2,
452 parent_leaf: Commitment<Leaf2<T>>,
453 ) -> bool {
454 self.blocks_reconstructed
455 .contains(&(parent_view, parent_block_commitment))
456 || self.locked_cert.as_ref().is_some_and(|lock| {
457 lock.view_number() == parent_view && lock.data().leaf_commit == parent_leaf
458 })
459 }
460
461 pub fn seed_state_cert(&mut self, state_cert: LightClientStateUpdateCertificateV2<T>) {
465 self.state_certs.insert(state_cert.epoch, state_cert);
466 }
467
468 pub fn apply_pre_cutover_seed(&mut self, seed: PreCutoverSeed<T>) {
478 let view = seed.decided_anchor.view_number();
479 if view > self.last_decided_view {
480 self.last_decided_view = view;
481 self.last_decided_leaf = seed.decided_anchor.clone();
482 self.decided_views.insert(view);
483 }
484 if view > self.decide_floor_view {
485 self.decide_floor_view = view;
486 }
487
488 let mut highest_seeded_block: u64 = seed.decided_anchor.block_header().block_number();
489
490 for leaf in seed.undecided {
491 let view = leaf.view_number();
492 let justify_qc = leaf.justify_qc().clone();
493 self.register_legacy_qc(&justify_qc);
494
495 let block_number = leaf.block_header().block_number();
496 let epoch = EpochNumber::new(epoch_from_block_number(block_number, *self.epoch_height));
497 if block_number > highest_seeded_block {
498 highest_seeded_block = block_number;
499 }
500
501 let view_change_evidence = leaf.view_change_evidence.clone().and_then(|e| match e {
502 ViewChangeEvidence2::Timeout(tc) => Some(tc),
503 ViewChangeEvidence2::ViewSync(_) => None,
504 });
505 let proposal = Proposal {
506 block_header: leaf.block_header().clone(),
507 view_number: view,
508 epoch,
509 justify_qc,
510 next_epoch_justify_qc: None,
511 upgrade_certificate: leaf.upgrade_certificate().clone(),
512 view_change_evidence,
513 next_drb_result: leaf.next_drb_result,
514 state_cert: None,
515 };
516
517 self.leaves.insert(view, leaf);
518 self.proposals.insert(view, proposal);
519 self.pre_cutover_views.insert(view);
520
521 self.proposed_views.insert(view);
522 self.voted_1_views.insert(view);
523 self.voted_2_views.insert(view);
524 }
525
526 if let Some(high_qc) = &seed.high_qc {
527 self.register_legacy_qc(high_qc);
528 }
529
530 let cutover_view = seed.cutover_view;
531 if cutover_view == ViewNumber::genesis() {
532 return;
533 }
534 let last_pre_cutover = cutover_view - 1;
535 if last_pre_cutover > self.timeout_view {
536 self.timeout_view = last_pre_cutover;
537 }
538 if last_pre_cutover > self.current_view {
539 self.current_view = last_pre_cutover;
540 }
541 let seeded_epoch = EpochNumber::new(epoch_from_block_number(
542 highest_seeded_block,
543 *self.epoch_height,
544 ));
545 if self.current_epoch.is_none_or(|cur| cur < seeded_epoch) {
546 self.current_epoch = Some(seeded_epoch);
547 }
548 }
549
550 pub(crate) fn register_legacy_qc(&mut self, justify_qc: &Certificate1<T>) {
553 let parent_view = justify_qc.view_number();
554 self.certs
555 .entry(parent_view)
556 .or_insert_with(|| justify_qc.clone());
557 if self
558 .locked_cert
559 .as_ref()
560 .is_none_or(|locked| locked.view_number() < parent_view)
561 {
562 self.locked_cert = Some(justify_qc.clone());
563 }
564 }
565
566 pub fn proposal_at(&self, view: ViewNumber) -> Option<&Proposal<T>> {
568 self.proposals.get(&view)
569 }
570
571 pub fn cert1_at(&self, view: ViewNumber) -> Option<&Certificate1<T>> {
573 self.certs.get(&view)
574 }
575
576 pub fn catchup_evidence(&self) -> Option<CatchupEvidence<T>> {
580 let tc = self.timeout_certs.last_key_value().map(|(_, tc)| tc);
581 match (tc, self.locked_cert.as_ref()) {
582 (Some(tc), Some(qc)) if qc.view_number() > tc.view_number() => {
583 Some(CatchupEvidence::Qc(qc.clone()))
584 },
585 (Some(tc), _) => Some(CatchupEvidence::Tc(tc.clone())),
586 (None, Some(qc)) => Some(CatchupEvidence::Qc(qc.clone())),
587 (None, None) => None,
588 }
589 }
590
591 fn signed_vid_share(
592 &self,
593 view: ViewNumber,
594 ) -> Option<SignedProposal<T, VidDisperseShare2<T>>> {
595 self.vid_shares
596 .get(&view)?
597 .clone()
598 .to_proposal(&self.private_key)
599 }
600
601 pub fn signed_proposal_fetch_request(
602 &self,
603 view: ViewNumber,
604 ) -> Result<ProposalFetchRequest<T>, <T::SignatureKey as SignatureKey>::SignError> {
605 ProposalFetchRequest::new(view, self.public_key.clone(), &self.private_key)
606 }
607
608 pub fn cert2_at(&self, view: ViewNumber) -> Option<&Certificate2<T>> {
610 self.certs2.get(&view)
611 }
612
613 pub fn timeout_cert_at(&self, view: ViewNumber) -> Option<&TimeoutCertificate2<T>> {
617 self.timeout_certs.get(&view)
618 }
619
620 pub fn locked_view(&self) -> Option<ViewNumber> {
622 self.locked_cert.as_ref().map(|c| c.view_number())
623 }
624
625 pub(crate) fn decide_floor(&self) -> ViewNumber {
629 max(
630 self.last_decided_view.saturating_sub(DECIDE_BUFFER).into(),
631 self.decide_floor_view,
632 )
633 }
634
635 #[instrument(level = "debug", skip_all, fields(node = %self.node_id, view = %input.view_number()))]
637 pub fn apply(&mut self, input: ConsensusInput<T>, outbox: &mut Outbox<ConsensusOutput<T>>) {
638 let view = if matches!(&input, ConsensusInput::DrbResult(..)) {
644 self.current_view
645 } else {
646 input.view_number()
647 };
648 let proto = match input {
649 ConsensusInput::Proposal(sender, proposal) => {
650 debug!(
651 sender = %KeyPrefix::from(&sender),
652 block = %proposal.proposal.data.block_header.block_number(),
653 epoch = %proposal.proposal.data.epoch,
654 "apply: proposal"
655 );
656 self.pair_proposal(sender, proposal, outbox)
657 },
658 ConsensusInput::VidShare(vid_share) => {
659 debug!("apply: vid share");
660 self.pair_vid_share(vid_share, outbox)
661 },
662 ConsensusInput::FetchedProposal(message) => {
663 debug!(
664 view = %message.proposal.data.view_number,
665 "apply: fetched proposal"
666 );
667 self.handle_fetched_proposal(message, outbox);
668 self.maybe_decide(view, outbox);
671 let views_extending_fetched: Vec<ViewNumber> = self
673 .proposals
674 .range(view + 1..)
675 .filter(|(_, proposal)| proposal.justify_qc.view_number() == view)
676 .map(|(extending_view, _)| *extending_view)
677 .collect();
678 for extending_view in views_extending_fetched {
679 self.maybe_vote_1(extending_view, outbox);
680 self.maybe_vote_2_and_update_lock(extending_view, outbox);
681 self.maybe_decide(extending_view, outbox);
682 }
683 self.maybe_propose(view + 1, outbox);
684 return;
685 },
686 ConsensusInput::Certificate1(certificate) => {
687 debug!(epoch = %certificate.epoch(), "apply: certificate1");
688 self.handle_certificate1(certificate)
689 },
690 ConsensusInput::Certificate2(certificate) => {
691 debug!(epoch = %certificate.epoch(), "apply: certificate2");
692 self.handle_certificate2(certificate, outbox)
693 },
694 ConsensusInput::AdvanceView(certificate) => {
695 debug!(
696 view = %certificate.view_number(),
697 epoch = %certificate.epoch(),
698 "apply: advance view"
699 );
700 self.handle_advance_view(certificate, outbox)
701 },
702 ConsensusInput::EpochRootCertificates { cert1, state_cert } => {
703 info!(
704 epoch = %state_cert.epoch,
705 "apply: epoch root certificates"
706 );
707 self.state_certs.insert(state_cert.epoch, state_cert);
711 self.handle_certificate1(cert1)
712 },
713 ConsensusInput::TimeoutCertificate(certificate) => {
714 let timed_out_view = certificate.view_number();
715 let leader = self.leader_label(timed_out_view, certificate.epoch());
716 warn!(
717 view = %timed_out_view,
718 epoch = %certificate.epoch(),
719 %leader,
720 "apply: timeout certificate"
721 );
722 self.handle_timeout_certificate(certificate, outbox)
723 },
724 ConsensusInput::BlockReconstructed(view, vid_commitment) => {
725 debug!(%view, "apply: block reconstructed");
726 self.blocks_reconstructed.insert((view, vid_commitment));
727 self.maybe_vote_1(view + 1, outbox);
729 Protocol::Continue
730 },
731 ConsensusInput::StateValidated(state_response) => {
732 debug!(view = %state_response.view, "apply: state validated");
733 self.states_verified
734 .insert(state_response.view, state_response.commitment);
735 Protocol::Continue
736 },
737 ConsensusInput::HeaderCreated(view, commitment, header) => {
738 debug!(%view, block = %header.block_number(), "apply: header created");
739 self.headers.insert((view, commitment), header);
740 Protocol::Continue
741 },
742 ConsensusInput::Stored(stored) => {
743 debug!(?stored, "apply: stored");
744 self.handle_stored(stored, outbox);
745 Protocol::Continue
746 },
747 ConsensusInput::StateValidationFailed(state_response) => {
748 let view = state_response.view;
749 let stored_proposal = self.proposals.get(&view);
750 if let Some(proposal) = stored_proposal {
751 let matches = proposal_commitment(proposal) == state_response.commitment;
752 warn!(
753 %view,
754 block = %proposal.block_header.block_number(),
755 epoch = %proposal.epoch,
756 qc_view = %proposal.justify_qc.view_number(),
757 qc_epoch = ?proposal.justify_qc.epoch(),
758 commitment_matches = matches,
759 "apply: state validation failed"
760 );
761 if !matches {
762 return;
763 }
764 } else {
765 warn!(%view, "apply: state validation failed (no stored proposal)");
766 }
767 self.proposals.remove(&view);
768 self.leaves.remove(&view);
769 self.vid_shares.remove(&view);
770 return;
771 },
772 ConsensusInput::Timeout(view, epoch) => {
773 let leader = self.leader_label(view, epoch);
774 warn!(%view, %epoch, %leader, "apply: timeout");
775 self.handle_timeout(view, epoch, outbox)
776 },
777 ConsensusInput::TimeoutOneHonest(view, epoch) => {
778 let leader = self.leader_label(view, epoch);
779 warn!(%view, %epoch, %leader, "apply: timeout (one honest)");
780 self.handle_timeout(view, epoch, outbox)
781 },
782 ConsensusInput::BlockBuilt {
783 view,
784 epoch,
785 payload,
786 metadata,
787 payload_commitment,
788 } => {
789 debug!(%view, %epoch, "apply: block built");
790 if let VidCommitment::V2(payload_commitment) = payload_commitment {
791 outbox.push_back(ConsensusOutput::RequestVidDisperse {
792 view,
793 epoch,
794 payload: payload.clone(),
795 metadata,
796 payload_commitment,
797 });
798 self.blocks.insert((view, payload_commitment), payload);
799 } else {
800 warn!(%view, %epoch, "block built with non-V2 payload commitment; ignoring");
801 }
802 Protocol::Continue
803 },
804 ConsensusInput::VidDisperseCreated(view, payload_commitment) => {
805 debug!(%view, "apply: vid disperse created");
806 self.blocks_reconstructed.insert((view, payload_commitment));
807 Protocol::Continue
808 },
809 ConsensusInput::DrbResult(epoch, drb_result) => {
810 info!(%epoch, "apply: drb result");
811 self.drb_results.insert(epoch, drb_result);
812 Protocol::Continue
813 },
814 ConsensusInput::EpochChange(epoch_change) => {
815 info!(
816 view = %epoch_change.cert1.view_number(),
817 epoch = ?epoch_change.cert1.epoch().map(|e| *e),
818 "apply: epoch change"
819 );
820 self.handle_epoch_change(epoch_change, outbox)
821 },
822 };
823
824 if matches!(proto, Protocol::Abort) {
825 debug!("aborting protocol");
826 return;
827 }
828
829 self.maybe_vote_1(view, outbox);
830 self.maybe_vote_2_and_update_lock(view, outbox);
831 self.maybe_decide(view, outbox);
832 self.maybe_propose(view, outbox);
833 self.maybe_propose(view + 1, outbox);
835 }
836
837 pub fn last_decided_view(&self) -> ViewNumber {
838 self.last_decided_view
839 }
840
841 pub fn last_decided_leaf(&self) -> &Leaf2<T> {
842 &self.last_decided_leaf
843 }
844
845 pub fn undecided_leaves(&self) -> impl Iterator<Item = &Leaf2<T>> {
846 self.leaves
847 .range((
848 std::ops::Bound::Excluded(self.last_decided_view),
849 std::ops::Bound::Unbounded,
850 ))
851 .map(|(_, leaf)| leaf)
852 }
853
854 pub fn current_view(&self) -> ViewNumber {
855 self.current_view
856 }
857
858 pub fn current_epoch(&self) -> Option<EpochNumber> {
859 self.current_epoch
860 }
861
862 #[cfg(test)]
863 pub fn set_view(&mut self, view: ViewNumber, epoch: EpochNumber) {
864 self.current_view = view;
865 self.current_epoch = Some(epoch);
866 }
867
868 pub fn resume_from_restart(
874 &mut self,
875 anchor_view: ViewNumber,
876 restart_view: ViewNumber,
877 last_actioned_view: ViewNumber,
878 ) {
879 let first_allowed = max(anchor_view + 1, max(restart_view, last_actioned_view + 1));
880 let last_barred = first_allowed - 1;
881 if last_barred > self.timeout_view {
882 self.timeout_view = last_barred;
883 }
884 if last_barred > self.restart_barred_view {
885 self.restart_barred_view = last_barred;
886 }
887 if anchor_view > self.decide_floor_view {
888 self.decide_floor_view = anchor_view;
889 }
890 let resume_view = self.stored_high_qc.unwrap_or(anchor_view + 1);
893 if resume_view > self.current_view {
894 self.current_view = resume_view;
895 }
896 }
897
898 pub fn wants_proposal_for_view(&self, view: &ViewNumber) -> bool {
899 let locked_too_new = self
900 .locked_cert
901 .as_ref()
902 .is_some_and(|l| l.view_number() > *view);
903 let fully_processed =
909 self.proposals.contains_key(view) && self.vid_shares.contains_key(view);
910 !(locked_too_new || fully_processed)
911 }
912
913 pub fn signed_proposal(&self, view: &ViewNumber) -> Option<&SignedProposal<T, Proposal<T>>> {
914 self.signed_proposals.get(view)
915 }
916
917 pub fn gc(&mut self, scope: GcScope) {
923 match scope {
924 GcScope::Local(view) => {
925 let c = Commitment::default_commitment_no_preimage();
926 let vc = VidCommitment2::default();
927 self.headers = self.headers.split_off(&(view, c));
928 self.unpaired_proposals = self.unpaired_proposals.split_off(&(view, vc));
929 self.unpaired_vid_shares = self.unpaired_vid_shares.split_off(&(view, vc));
930 self.proposed_views = self.proposed_views.split_off(&view);
931 self.states_verified = self.states_verified.split_off(&view);
932 self.timeout_certs = self.timeout_certs.split_off(&view);
933 self.voted_1_views = self.voted_1_views.split_off(&view);
934 self.voted_2_views = self.voted_2_views.split_off(&view);
935 },
936 GcScope::Decided(view) => {
937 let vc = VidCommitment2::default();
938 self.blocks = self.blocks.split_off(&(view, vc));
939 self.blocks_reconstructed = self.blocks_reconstructed.split_off(&(view, vc));
940 let keep_from = self.decide_floor();
941 self.certs = self.certs.split_off(&keep_from);
942 self.certs2 = self.certs2.split_off(&keep_from);
943 self.decided_views = self.decided_views.split_off(&keep_from);
944 self.proposals = self.proposals.split_off(&keep_from);
945 self.leaves = self.leaves.split_off(&view);
946 self.signed_proposals = self.signed_proposals.split_off(&view);
947 self.vid_shares = self.vid_shares.split_off(&view);
948 self.stored_proposals = self.stored_proposals.split_off(&view);
949 self.stored_vids = self.stored_vids.split_off(&view);
950 self.stored_actions = self.stored_actions.split_off(&(view, ActionKind::Vote));
951 self.requested_actions =
952 self.requested_actions.split_off(&(view, ActionKind::Vote));
953 self.pending_vote1 = self.pending_vote1.split_off(&view);
954 self.pending_vote2 = self.pending_vote2.split_off(&view);
955 self.pending_proposal = self.pending_proposal.split_off(&view);
956 if let Some(epoch) = self.current_epoch {
957 let epoch = EpochNumber::new(epoch.saturating_sub(1));
958 self.drb_results = self.drb_results.split_off(&epoch);
959 self.state_certs = self.state_certs.split_off(&epoch);
960 }
961 },
962 GcScope::Timeout(view) => {
963 if self.certs.contains_key(&view) || self.certs2.contains_key(&view) {
966 return;
967 }
968 self.vid_shares.remove(&view);
969 let vc = VidCommitment2::default();
970 self.blocks
971 .extract_if((view, vc)..(view + 1, vc), |_, _| true)
972 .for_each(drop);
973 },
974 }
975 }
976
977 #[cfg(test)]
985 pub(crate) fn force_set_proposal(&mut self, view: ViewNumber, proposal: Proposal<T>) {
986 self.proposals.insert(view, proposal);
987 }
988
989 fn pair_proposal(
993 &mut self,
994 sender: T::SignatureKey,
995 proposal: ProposalMessage<T, Validated>,
996 outbox: &mut Outbox<ConsensusOutput<T>>,
997 ) -> Protocol {
998 let view = proposal.view_number();
999 let VidCommitment::V2(commit) = proposal.proposal.data.block_header.payload_commitment()
1000 else {
1001 warn!(%view, "proposal payload commitment is not V2, discarding");
1002 return Protocol::Abort;
1003 };
1004 let Some(vid_share) = self.unpaired_vid_shares.remove(&(view, commit)) else {
1005 self.unpaired_proposals
1006 .insert((view, commit), (sender, proposal));
1007 return Protocol::Abort;
1008 };
1009 self.on_proposal_paired(sender, proposal, vid_share, outbox)
1010 }
1011
1012 fn pair_vid_share(
1016 &mut self,
1017 vid_share: VidDisperseShare2<T>,
1018 outbox: &mut Outbox<ConsensusOutput<T>>,
1019 ) -> Protocol {
1020 let key = (vid_share.view_number(), vid_share.payload_commitment);
1021 let Some((sender, proposal)) = self.unpaired_proposals.remove(&key) else {
1022 self.unpaired_vid_shares.insert(key, vid_share);
1023 return Protocol::Abort;
1024 };
1025 self.on_proposal_paired(sender, proposal, vid_share, outbox)
1026 }
1027
1028 fn on_proposal_paired(
1029 &mut self,
1030 sender: T::SignatureKey,
1031 proposal: ProposalMessage<T, Validated>,
1032 vid_share: VidDisperseShare2<T>,
1033 outbox: &mut Outbox<ConsensusOutput<T>>,
1034 ) -> Protocol {
1035 let view = proposal.view_number();
1037 let vc = VidCommitment2::default();
1038 self.unpaired_proposals = self.unpaired_proposals.split_off(&(view + 1, vc));
1039 self.unpaired_vid_shares = self.unpaired_vid_shares.split_off(&(view + 1, vc));
1040 outbox.push_back(ConsensusOutput::ProposalPaired {
1041 proposal: proposal.proposal.clone(),
1042 vid_share: vid_share.clone(),
1043 });
1044 self.handle_proposal_with_vid_share(sender, proposal, vid_share, outbox)
1045 }
1046
1047 #[instrument(level = "debug", skip_all)]
1048 fn handle_proposal_with_vid_share(
1049 &mut self,
1050 sender: T::SignatureKey,
1051 proposal: ProposalMessage<T, Validated>,
1052 vid_share: VidDisperseShare2<T>,
1053 outbox: &mut Outbox<ConsensusOutput<T>>,
1054 ) -> Protocol {
1055 let view = proposal.view_number();
1056 let proposer = KeyPrefix::from(&sender);
1057 let block_number = proposal.proposal.data.block_header.block_number();
1058 let qc_view = proposal.proposal.data.justify_qc.view_number();
1059
1060 if !self.wants_proposal_for_view(&view) {
1061 warn!(
1062 %view, %proposer, block = %block_number,
1063 epoch = %proposal.proposal.data.epoch, %qc_view,
1064 "proposal too old"
1065 );
1066 return Protocol::Abort;
1067 }
1068
1069 let signed_proposal = proposal.proposal.clone();
1070 let proposal = proposal.proposal.data;
1071 let epoch = proposal.epoch;
1072 let Some(qc_epoch) = proposal.justify_qc.epoch() else {
1074 warn!(
1075 %view, %proposer, block = %block_number, %epoch, %qc_view,
1076 "proposal has no epoch number"
1077 );
1078 return Protocol::Abort;
1079 };
1080
1081 if let Err(err) = self.is_safe(&proposal) {
1082 warn!(
1083 %view, %proposer, block = %block_number, %epoch, %qc_view, %qc_epoch, %err,
1084 "proposal not safe"
1085 );
1086 return Protocol::Abort;
1087 }
1088
1089 let payload_size = vid_share.payload_byte_len();
1090
1091 self.proposals.insert(view, proposal.clone());
1096 self.signed_proposals.insert(view, signed_proposal.clone());
1097 self.leaves.insert(view, proposal.clone().into());
1098 self.vid_shares.insert(view, vid_share);
1099 self.adopt_certified_drb(view);
1100
1101 self.request_parent_proposal_if_missing(&proposal, outbox);
1102
1103 if let Some(state_cert) = &proposal.state_cert {
1104 self.state_certs
1105 .entry(state_cert.epoch)
1106 .or_insert_with(|| state_cert.clone());
1107 }
1108
1109 if proposal.epoch > EpochNumber::genesis()
1117 && is_epoch_transition(block_number, *self.epoch_height)
1118 {
1119 if let Some(drb) = self.drb_results.get(&(epoch + 1)) {
1120 if proposal
1121 .next_drb_result
1122 .is_none_or(|proposed_drb| drb != &proposed_drb)
1123 {
1124 warn!(
1125 %view, %proposer, block = %block_number, %epoch, %qc_view, %qc_epoch,
1126 "DRB result does not match proposal"
1127 );
1128 return Protocol::Abort;
1129 }
1130 } else {
1131 outbox.push_back(ConsensusOutput::RequestDrbResult(epoch + 1));
1132 }
1133 }
1134
1135 outbox.push_back(ConsensusOutput::RequestState(StateRequest {
1136 view: proposal.view_number(),
1137 parent_view: proposal.justify_qc.view_number(),
1138 epoch,
1139 block: proposal.block_header.block_number().into(),
1140 proposal: proposal.clone(),
1141 parent_commitment: proposal.justify_qc.data().leaf_commit,
1142 payload_size,
1143 }));
1144
1145 let epoch = if is_last_block(block_number, *self.epoch_height) {
1146 epoch + 1
1147 } else {
1148 epoch
1149 };
1150
1151 outbox.push_back(ConsensusOutput::ProposalValidated {
1152 proposal: signed_proposal,
1153 sender,
1154 });
1155
1156 if self.is_leader(view + 1, epoch) {
1157 outbox.push_back(ConsensusOutput::RequestBlockAndHeader(
1158 BlockAndHeaderRequest {
1159 view: view + 1,
1160 epoch,
1161 parent_proposal: proposal,
1162 },
1163 ));
1164 }
1165
1166 Protocol::Continue
1167 }
1168
1169 #[instrument(level = "debug", skip_all)]
1170 fn handle_fetched_proposal(
1171 &mut self,
1172 message: ProposalMessage<T, Validated>,
1173 outbox: &mut Outbox<ConsensusOutput<T>>,
1174 ) {
1175 let signed_proposal = message.proposal;
1176 let proposal = signed_proposal.data.clone();
1177 let view = proposal.view_number;
1178 if view <= self.last_decided_view {
1179 debug!(%view, "fetched proposal at or below decided view; discarding");
1180 return;
1181 }
1182 if self.proposals.contains_key(&view) {
1183 debug!(%view, "fetched proposal already present; discarding");
1184 return;
1185 }
1186 self.leaves.insert(view, proposal.clone().into());
1187 self.signed_proposals.insert(view, signed_proposal);
1188 self.request_parent_proposal_if_missing(&proposal, outbox);
1189 self.proposals.insert(view, proposal);
1190 self.adopt_certified_drb(view);
1191 }
1192
1193 fn request_parent_proposal_if_missing(
1194 &self,
1195 proposal: &Proposal<T>,
1196 outbox: &mut Outbox<ConsensusOutput<T>>,
1197 ) {
1198 let parent_view = proposal.justify_qc.view_number();
1199 if parent_view > self.last_decided_view && !self.proposals.contains_key(&parent_view) {
1200 warn!(
1201 view = %proposal.view_number,
1202 %parent_view,
1203 "parent proposal missing; requesting fetch"
1204 );
1205 outbox.push_back(ConsensusOutput::RequestMissingProposal {
1206 view: parent_view,
1207 leaf_commit: proposal.justify_qc.data().leaf_commit,
1208 });
1209 }
1210 }
1211
1212 #[instrument(level = "debug", skip_all)]
1213 fn handle_certificate1(&mut self, certificate: ValidCert<Certificate1<T>>) -> Protocol {
1214 let view = certificate.view_number();
1215 if view <= self.decide_floor() {
1216 return Protocol::Continue;
1217 }
1218 self.certs.entry(view).or_insert(certificate.into_cert());
1219 self.adopt_certified_drb(view);
1220 Protocol::Continue
1221 }
1222
1223 #[instrument(level = "debug", skip_all)]
1224 fn handle_certificate2(
1225 &mut self,
1226 certificate: ValidCert<Certificate2<T>>,
1227 outbox: &mut Outbox<ConsensusOutput<T>>,
1228 ) -> Protocol {
1229 let view = certificate.view_number();
1230 if view <= self.decide_floor() {
1231 return Protocol::Continue;
1232 }
1233 if self.certs2.contains_key(&view) {
1234 return Protocol::Continue;
1235 }
1236 if !self.decided_views.contains(&view) {
1240 outbox.push_back(ConsensusOutput::SendCertificate2(
1241 certificate.cert().clone(),
1242 ));
1243 }
1244 if view > self.last_decided_view && !self.proposals.contains_key(&view) {
1245 warn!(%view, "have certificate2 but no proposal; requesting fetch");
1246 outbox.push_back(ConsensusOutput::RequestMissingProposal {
1247 view,
1248 leaf_commit: certificate.data.leaf_commit,
1249 });
1250 }
1251 self.certs2.insert(view, certificate.into_cert());
1252 Protocol::Continue
1253 }
1254
1255 fn adopt_certified_drb(&mut self, view: ViewNumber) {
1260 let Some(proposal) = self.proposals.get(&view) else {
1261 return;
1262 };
1263 if proposal.epoch <= EpochNumber::genesis()
1265 || !is_epoch_transition(proposal.block_header.block_number(), *self.epoch_height)
1266 {
1267 return;
1268 }
1269 let Some(drb) = proposal.next_drb_result else {
1270 return;
1271 };
1272 let next_epoch = proposal.epoch + 1;
1273 if self.drb_results.contains_key(&next_epoch) {
1274 return;
1275 }
1276 let Some(cert) = self.certs.get(&view) else {
1278 return;
1279 };
1280 if proposal_commitment(proposal) != cert.data.leaf_commit {
1281 return;
1282 }
1283 self.drb_results.insert(next_epoch, drb);
1284 debug!(%view, %next_epoch, "adopted quorum-certified next_drb_result");
1285 }
1286
1287 #[instrument(level = "debug", skip_all)]
1289 fn handle_advance_view(
1290 &mut self,
1291 cert1: ValidCert<Certificate1<T>>,
1292 outbox: &mut Outbox<ConsensusOutput<T>>,
1293 ) -> Protocol {
1294 let view = cert1.view_number();
1295
1296 if view < self.current_view {
1297 return Protocol::Continue;
1298 }
1299
1300 let epoch = cert1.epoch();
1301
1302 self.certs.entry(view).or_insert(cert1.into_cert());
1303 self.adopt_certified_drb(view);
1304
1305 self.maybe_vote_2_and_update_lock(view, outbox);
1307
1308 let next_view = view + 1;
1309
1310 if next_view > self.current_view {
1311 self.current_view = next_view;
1312 self.current_epoch = Some(epoch);
1313 outbox.push_back(ConsensusOutput::ViewChanged(next_view, epoch));
1314 }
1315
1316 Protocol::Continue
1317 }
1318
1319 #[instrument(level = "debug", skip_all)]
1320 fn handle_timeout(
1321 &mut self,
1322 view: ViewNumber,
1323 epoch: EpochNumber,
1324 outbox: &mut Outbox<ConsensusOutput<T>>,
1325 ) -> Protocol {
1326 if view < self.current_view {
1327 debug!(
1328 %view,
1329 current_view = %self.current_view,
1330 "ignoring timeout for stale view"
1331 );
1332 return Protocol::Abort;
1333 }
1334 let we_were_leader = self.is_leader(view, epoch);
1335 if we_were_leader {
1336 if self.proposed_views.contains(&view) {
1337 warn!(%view, %epoch, "timeout: we were the leader and did propose for this view");
1338 } else {
1339 let missing = self.missing_for_propose(view);
1340 let missing_str = if missing.is_empty() {
1341 "none".to_string()
1342 } else {
1343 missing.join(",")
1344 };
1345 warn!(
1346 %view, %epoch, missing = %missing_str,
1347 "timeout: we were the leader but did not propose"
1348 );
1349 }
1350 }
1351
1352 if !we_were_leader || self.proposed_views.contains(&view) {
1356 if self.voted_1_views.contains(&view) {
1357 warn!(%view, %epoch, "timeout: we did vote1 for this view");
1358 } else {
1359 let missing = self.missing_for_vote1(view);
1360 let missing_str = if missing.is_empty() {
1361 "none".to_string()
1362 } else {
1363 missing.join(",")
1364 };
1365 warn!(
1366 %view, %epoch, missing = %missing_str,
1367 "timeout: we did not vote1 for this view"
1368 );
1369 }
1370 }
1371
1372 if self.certs.contains_key(&view) {
1376 let proposal_commit = self
1377 .proposals
1378 .get(&view)
1379 .map(|p| p.block_header.payload_commitment());
1380 match proposal_commit {
1381 Some(VidCommitment::V2(prop))
1382 if self.blocks_reconstructed.contains(&(view, prop)) =>
1383 {
1384 warn!(%view, %epoch, "timeout: have cert1 and matching reconstructed block");
1385 },
1386 Some(VidCommitment::V2(_)) => {
1387 warn!(
1388 %view, %epoch,
1389 "timeout: have cert1, but no reconstructed block matching the proposal"
1390 );
1391 },
1392 Some(_) => {
1393 warn!(
1396 %view, %epoch,
1397 "timeout: have cert1 but proposal payload commitment is not V2"
1398 );
1399 },
1400 None => {
1401 warn!(%view, %epoch, "timeout: have cert1 but no proposal stored");
1402 },
1403 }
1404 }
1405 self.timeout_view = max(self.timeout_view, view);
1406 let data = TimeoutData2 {
1407 view,
1408 epoch: Some(epoch),
1409 };
1410 let vote = match SimpleVote::create_signed_vote(
1411 data,
1412 view,
1413 &self.public_key,
1414 &self.private_key,
1415 &self.upgrade_lock,
1416 ) {
1417 Ok(vote) => vote,
1418 Err(err) => {
1419 warn!(%view, %err, "failed to create timeout vote");
1420 return Protocol::Abort;
1421 },
1422 };
1423 outbox.push_back(ConsensusOutput::SendTimeoutVote(
1424 vote,
1425 self.catchup_evidence(),
1426 ));
1427 Protocol::Abort
1428 }
1429
1430 #[instrument(level = "debug", skip_all)]
1431 fn handle_timeout_certificate(
1432 &mut self,
1433 certificate: ValidCert<TimeoutCertificate2<T>>,
1434 outbox: &mut Outbox<ConsensusOutput<T>>,
1435 ) -> Protocol {
1436 let view = certificate.view_number() + 1;
1437 if view < self.current_view {
1438 debug!(
1439 %view,
1440 current_view = %self.current_view,
1441 "ignoring stale timeout certificate"
1442 );
1443 return Protocol::Abort;
1444 }
1445 if self.timeout_certs.contains_key(&view) {
1446 return Protocol::Continue;
1447 }
1448 let epoch = certificate.epoch();
1449 self.timeout_certs.insert(view, certificate.cert().clone());
1450 self.current_view = self.current_view.max(view);
1451 self.current_epoch = Some(epoch);
1452 outbox.push_back(ConsensusOutput::ViewChanged(view, epoch));
1453 outbox.push_back(ConsensusOutput::ViewTimedOut(certificate.view_number()));
1454 outbox.push_back(ConsensusOutput::SendTimeoutCertificate(
1455 certificate.into_cert(),
1456 view,
1457 epoch,
1458 ));
1459 if !self.is_leader(view, epoch) {
1460 debug!(%epoch, "not leader");
1461 return Protocol::Abort;
1462 }
1463
1464 let Some(locked_view) = self.locked_cert.as_ref().map(|cert| cert.view_number()) else {
1467 debug!("locked certificate not available");
1468 return Protocol::Abort;
1469 };
1470 let Some(proposal) = self.proposals.get(&locked_view) else {
1471 debug!(%locked_view, "proposal not available");
1472 return Protocol::Abort;
1473 };
1474 outbox.push_back(ConsensusOutput::RequestBlockAndHeader(
1477 BlockAndHeaderRequest {
1478 view,
1479 epoch,
1480 parent_proposal: proposal.clone(),
1481 },
1482 ));
1483 Protocol::Continue
1484 }
1485
1486 #[instrument(level = "debug", skip_all)]
1487 fn handle_epoch_change(
1488 &mut self,
1489 epoch_change: EpochChangeMessage<T, Validated>,
1490 outbox: &mut Outbox<ConsensusOutput<T>>,
1491 ) -> Protocol {
1492 let EpochChangeMessage {
1493 cert1,
1494 cert2,
1495 proposal,
1496 ..
1497 } = epoch_change;
1498 if self
1501 .current_epoch
1502 .is_some_and(|current| cert2.data.epoch < current)
1503 {
1504 debug!(
1505 view = %cert2.view_number(),
1506 epoch = %cert2.data.epoch,
1507 current_epoch = ?self.current_epoch.map(|e| *e),
1508 "ignoring stale epoch change for an epoch we have already entered"
1509 );
1510 return Protocol::Abort;
1511 }
1512 if self
1514 .locked_cert
1515 .as_ref()
1516 .is_some_and(|locked_cert| locked_cert.view_number() > cert1.view_number())
1517 {
1518 warn!("locked certificate is newer than epoch change certificate1");
1519 return Protocol::Abort;
1520 }
1521 if cert1.view_number() != cert2.view_number()
1523 || cert1.epoch() != cert2.epoch()
1524 || cert1.data.leaf_commit != cert2.data.leaf_commit
1525 {
1526 warn!("epoch change certificates do not match");
1527 return Protocol::Abort;
1528 }
1529 if !is_last_block(cert2.data.block_number, *self.epoch_height) {
1531 warn!("epoch change certificate2 is not the last block of the epoch");
1532 return Protocol::Abort;
1533 }
1534 if cert2.data.block_number / *self.epoch_height != *cert2.data.epoch {
1535 warn!("epoch change certificate2 is not for the correct epoch");
1536 return Protocol::Abort;
1537 }
1538 if proposal_commitment(&proposal) != cert1.data.leaf_commit {
1540 warn!("epoch change proposal commitment does not match certificate1 leaf commitment");
1541 return Protocol::Abort;
1542 }
1543 let next_view = cert2.view_number() + 1;
1544 let next_epoch = cert2.data.epoch + 1;
1545 self.current_view = self.current_view.max(next_view);
1547 self.current_epoch = Some(next_epoch);
1548 outbox.push_back(ConsensusOutput::ViewChanged(next_view, next_epoch));
1549
1550 if self.is_leader(next_view, next_epoch) {
1552 outbox.push_back(ConsensusOutput::RequestBlockAndHeader(
1553 BlockAndHeaderRequest {
1554 view: next_view,
1555 epoch: next_epoch,
1556 parent_proposal: proposal.clone(),
1557 },
1558 ));
1559 }
1560
1561 let boundary_view = cert2.view_number();
1562 self.proposals.insert(boundary_view, proposal);
1563 self.certs.insert(cert1.view_number(), cert1);
1564 self.certs2.insert(boundary_view, cert2);
1565 self.adopt_certified_drb(boundary_view);
1566 Protocol::Continue
1567 }
1568
1569 #[instrument(level = "debug", skip_all)]
1570 fn maybe_propose(&mut self, view: ViewNumber, outbox: &mut Outbox<ConsensusOutput<T>>) {
1571 if view <= self.timeout_view {
1572 return;
1573 }
1574 if self.proposed_views.contains(&view) {
1575 return;
1576 }
1577
1578 let view_change_evidence = self.timeout_certs.get(&view).cloned();
1579 let parent_cert = if view_change_evidence.is_some() {
1580 let Some(cert) = &self.locked_cert else {
1581 debug!("no locked qc");
1582 return;
1583 };
1584 cert
1585 } else {
1586 let Some(cert) = self.certs.get(&ViewNumber::from(view.saturating_sub(1))) else {
1587 debug!("no parent certificate");
1588 return;
1589 };
1590 cert
1591 };
1592 let parent_view = parent_cert.view_number();
1593 let Some(proposal) = self.proposals.get(&parent_view) else {
1594 debug!(parent = %parent_view, "no proposal for parent view");
1595 return;
1596 };
1597
1598 let parent_commitment = if parent_view == ViewNumber::genesis() {
1612 proposal_commitment(proposal)
1613 } else if proposal_commitment(proposal) != parent_cert.data.leaf_commit {
1614 warn!(
1615 %parent_view,
1616 "stored proposal at parent_view does not match parent cert's leaf_commit; \
1617 refusing to propose with mismatched parent"
1618 );
1619 return;
1620 } else {
1621 parent_cert.data.leaf_commit
1622 };
1623 let Some(header) = self.headers.get(&(view, parent_commitment)) else {
1624 if view_change_evidence.is_some() {
1628 let request_epoch =
1629 if is_last_block(proposal.block_header.block_number(), *self.epoch_height) {
1630 proposal.epoch + 1
1631 } else {
1632 proposal.epoch
1633 };
1634 if self.is_leader(view, request_epoch) {
1635 outbox.push_back(ConsensusOutput::RequestBlockAndHeader(
1636 BlockAndHeaderRequest {
1637 view,
1638 epoch: request_epoch,
1639 parent_proposal: proposal.clone(),
1640 },
1641 ));
1642 }
1643 }
1644 debug!("no block header");
1645 return;
1646 };
1647 let VidCommitment::V2(block_commitment) = header.payload_commitment() else {
1648 debug!("header payload commitment is not V2");
1649 return;
1650 };
1651 if !self.blocks.contains_key(&(view, block_commitment)) {
1652 debug!("no block");
1653 return;
1654 };
1655
1656 let first_proposal_of_epoch =
1657 is_last_block(header.block_number().saturating_sub(1), *self.epoch_height);
1658 let proposal_epoch = if first_proposal_of_epoch {
1659 proposal.epoch + 1
1660 } else {
1661 proposal.epoch
1662 };
1663 if !self.is_leader(view, proposal_epoch) {
1664 warn!(epoch = %proposal_epoch, "not the leader for this view, we should not have a header");
1665 return;
1666 }
1667
1668 let next_drb_result = if proposal.epoch > EpochNumber::genesis()
1676 && is_epoch_transition(header.block_number(), *self.epoch_height)
1677 {
1678 let Some(drb) = self.drb_results.get(&EpochNumber::new(*proposal.epoch + 1)) else {
1679 debug!(%proposal.epoch, "no DRB result for epoch");
1680 outbox.push_back(ConsensusOutput::RequestDrbResult(proposal.epoch + 1));
1684 return;
1685 };
1686 Some(*drb)
1687 } else {
1688 None
1689 };
1690 let next_epoch_justify_qc = if first_proposal_of_epoch {
1691 let Some(next_epoch_justify_qc) = self.certs2.get(&parent_view) else {
1692 debug!("no next epoch justify QC");
1693 return;
1694 };
1695 Some(next_epoch_justify_qc.clone())
1696 } else {
1697 None
1698 };
1699
1700 let parent_block_number = parent_cert.data.block_number.unwrap_or(0);
1705 let state_cert = if is_epoch_root(parent_block_number, *self.epoch_height) {
1706 let Some(parent_epoch) = parent_cert.data.epoch() else {
1707 warn!("epoch-root parent QC has no epoch; cannot propose");
1708 return;
1709 };
1710 let Some(sc) = self.state_certs.get(&parent_epoch).cloned() else {
1711 warn!(
1712 %view,
1713 "epoch-root parent QC without state_cert — atomicity invariant broken; skipping propose"
1714 );
1715 return;
1716 };
1717 if !check_qc_state_cert_correspondence(parent_cert, &sc, *self.epoch_height) {
1718 warn!(%view, "state_cert does not correspond to parent QC; skipping propose");
1719 return;
1720 }
1721 Some(sc)
1722 } else {
1723 None
1724 };
1725
1726 let proposal = Proposal::<T> {
1727 block_header: header.clone(),
1728 view_number: view,
1729 epoch: proposal_epoch,
1730 justify_qc: parent_cert.clone(),
1731 next_epoch_justify_qc,
1732 upgrade_certificate: None,
1733 view_change_evidence,
1734 next_drb_result,
1735 state_cert,
1736 };
1737
1738 let proposed_leaf: Leaf2<T> = proposal.clone().into();
1740 let signature =
1741 match T::SignatureKey::sign(&self.private_key, proposed_leaf.commit().as_ref()) {
1742 Ok(sig) => sig,
1743 Err(err) => {
1744 warn!(%view, %err, "failed to sign proposal");
1745 return;
1746 },
1747 };
1748
1749 let message = SignedProposal {
1750 data: proposal,
1751 signature,
1752 _pd: PhantomData,
1753 };
1754
1755 self.proposed_views.insert(view);
1756 outbox.push_back(ConsensusOutput::PersistProposal(message.clone()));
1757 self.request_action(view, Some(proposal_epoch), ActionKind::Propose, outbox);
1758 self.pending_proposal.insert(view, message);
1759 }
1760
1761 #[instrument(level = "debug", skip_all)]
1762 fn maybe_decide(&mut self, view: ViewNumber, outbox: &mut Outbox<ConsensusOutput<T>>) {
1763 let floor = self.decide_floor();
1766 if view <= floor || self.decided_views.contains(&view) {
1767 return;
1768 }
1769 let Some(cert2) = self.certs2.get(&view) else {
1770 debug!(%view, "cert2 not available");
1771 return;
1772 };
1773 let Some(proposal) = self.proposals.get(&view) else {
1774 debug!(%view, "proposal not available");
1775 return;
1776 };
1777 let block = proposal.block_header.block_number();
1778 let epoch = proposal.epoch;
1779 let qc_view = proposal.justify_qc.view_number();
1780 let qc_epoch = proposal.justify_qc.epoch();
1781 let proposal_commit = proposal_commitment(proposal);
1782 if cert2.data.leaf_commit != proposal_commit {
1783 debug!(
1784 %view, %block, %epoch, %qc_view, ?qc_epoch,
1785 "cert2 commitment does not match proposal commitment"
1786 );
1787 return;
1788 }
1789 let Some(cert1) = self.certs.get(&view).cloned() else {
1792 debug!(%view, "cert1 missing");
1793 return;
1794 };
1795 if is_last_block(proposal.block_header.block_number(), *self.epoch_height)
1798 && cert1.data.leaf_commit == proposal_commit
1799 {
1800 let epoch_change =
1801 EpochChangeMessage::validated(cert1.clone(), cert2.clone(), proposal.clone());
1802 outbox.push_back(ConsensusOutput::SendEpochChange(epoch_change));
1803 }
1804 let mut leaf: Leaf2<T> = proposal.clone().into();
1806 if let VidCommitment::V2(pc) = proposal.block_header.payload_commitment()
1807 && let Some(payload) = self.blocks.get(&(view, pc))
1808 {
1809 leaf.fill_block_payload_unchecked(payload.clone());
1810 }
1811 let mut decided = vec![leaf];
1812 let mut vid_shares = vec![self.signed_vid_share(view)];
1813
1814 let mut parent_view = proposal.justify_qc.view_number();
1815 let mut parent_commit = proposal.justify_qc.data.leaf_commit;
1816
1817 while parent_view > floor
1819 && !self.decided_views.contains(&parent_view)
1820 && let Some(proposal) = self.proposals.get(&parent_view)
1821 {
1822 let proposal_commit = proposal_commitment(proposal);
1823 if proposal_commit != parent_commit {
1824 break;
1825 }
1826 let mut leaf: Leaf2<T> = proposal.clone().into();
1827 if let VidCommitment::V2(pc) = proposal.block_header.payload_commitment()
1828 && let Some(payload) = self.blocks.get(&(parent_view, pc))
1829 {
1830 leaf.fill_block_payload_unchecked(payload.clone());
1831 }
1832 vid_shares.push(self.signed_vid_share(parent_view));
1833 decided.push(leaf);
1834 parent_view = proposal.justify_qc.view_number();
1835 parent_commit = proposal.justify_qc.data.leaf_commit;
1836 }
1837 self.decided_views
1838 .extend(decided.iter().map(|l| l.view_number()));
1839 if view > self.last_decided_view {
1841 self.last_decided_view = view;
1842 self.last_decided_leaf = decided[0].clone();
1843 }
1844 outbox.push_back(ConsensusOutput::LeafDecided {
1845 leaves: decided,
1846 cert1,
1847 cert2: Some(cert2.clone()),
1848 vid_shares,
1849 });
1850 }
1851
1852 fn build_state_vote(
1859 &self,
1860 proposal: &Proposal<T>,
1861 ) -> anyhow::Result<LightClientStateUpdateVote2<T>> {
1862 let view_number = proposal.view_number;
1863 let light_client_state = proposal
1864 .block_header
1865 .get_light_client_state(view_number)
1866 .map_err(|e| anyhow::anyhow!("failed to generate light client state: {e}"))?;
1867 let auth_root = proposal
1868 .block_header
1869 .auth_root()
1870 .map_err(|e| anyhow::anyhow!("failed to fetch auth root: {e}"))?;
1871 let membership = self
1872 .stake_table_coordinator
1873 .membership_for_epoch(Some(proposal.epoch))
1874 .map_err(|e| anyhow::anyhow!("membership lookup failed: {e}"))?;
1875 let next_stake_table = membership
1876 .next_epoch_stake_table()
1877 .map_err(|e| anyhow::anyhow!("next-epoch stake table lookup failed: {e}"))?;
1878 let next_stake_table_state = HSStakeTable::from_iter(next_stake_table.stake_table())
1879 .commitment(self.stake_table_capacity)
1880 .map_err(|e| anyhow::anyhow!("failed to compute stake table commitment: {e}"))?;
1881 let v2_signature = <T::StateSignatureKey as LCV2StateSignatureKey>::sign_state(
1882 &self.state_private_key,
1883 &light_client_state,
1884 &next_stake_table_state,
1885 )
1886 .map_err(|e| anyhow::anyhow!("failed to sign LCV2 state: {e}"))?;
1887 let signed_state_digest =
1888 derive_signed_state_digest(&light_client_state, &next_stake_table_state, &auth_root);
1889 let signature = <T::StateSignatureKey as LCV3StateSignatureKey>::sign_state(
1890 &self.state_private_key,
1891 signed_state_digest,
1892 )
1893 .map_err(|e| anyhow::anyhow!("failed to sign LCV3 state: {e}"))?;
1894 Ok(LightClientStateUpdateVote2 {
1895 epoch: proposal.epoch,
1896 light_client_state,
1897 next_stake_table_state,
1898 signature,
1899 v2_signature,
1900 auth_root,
1901 signed_state_digest,
1902 })
1903 }
1904
1905 fn handle_stored(&mut self, stored: StorageOutput<T>, outbox: &mut Outbox<ConsensusOutput<T>>) {
1906 let view = stored.view_number();
1907 match stored {
1908 StorageOutput::Proposal(view, commitment) => {
1909 self.stored_proposals
1910 .entry(view)
1911 .or_default()
1912 .push(commitment);
1913 },
1914 StorageOutput::Vid(view) => {
1915 self.stored_vids.insert(view);
1916 },
1917 StorageOutput::Action(view, kind) => {
1918 self.stored_actions.insert((view, kind));
1919 },
1920 StorageOutput::HighQc(view) => {
1921 self.bump_stored_high_qc(view);
1922 let pending: Vec<ViewNumber> = self.pending_vote2.keys().copied().collect();
1924 for view in pending {
1925 self.release_vote2(view, outbox);
1926 }
1927 return;
1928 },
1929 }
1930 self.release_vote1(view, outbox);
1931 self.release_vote2(view, outbox);
1932 self.release_proposal(view, outbox);
1933 }
1934
1935 fn release_vote1(&mut self, view: ViewNumber, outbox: &mut Outbox<ConsensusOutput<T>>) {
1936 let Some(vote1) = self.pending_vote1.get(&view) else {
1937 return;
1938 };
1939 if !self.stored_actions.contains(&(view, ActionKind::Vote))
1940 || !self.is_proposal_stored(view, &vote1.vote.data.leaf_commit)
1941 {
1942 return;
1943 }
1944 let vote1 = self.pending_vote1.remove(&view).expect("checked above");
1945 if view <= self.timeout_view {
1946 debug!(%view, "dropping pending vote1 for timed-out view");
1947 return;
1948 }
1949 let vid_share = self.vid_shares.get(&view).cloned();
1950 outbox.push_back(ConsensusOutput::SendVote1(vote1));
1951 if let Some(vid_share) = vid_share {
1952 outbox.push_back(ConsensusOutput::BroadcastVidShare(vid_share));
1953 } else {
1954 debug!(%view, "vid share gone for released vote1; skipping broadcast");
1955 }
1956 }
1957
1958 fn release_vote2(&mut self, view: ViewNumber, outbox: &mut Outbox<ConsensusOutput<T>>) {
1959 let Some((_, required)) = self.pending_vote2.get(&view) else {
1960 return;
1961 };
1962 if self.certs2.contains_key(&view) {
1963 self.pending_vote2.remove(&view);
1964 return;
1965 }
1966 let required = *required;
1967 if !self.vote2_persisted(view) || !self.high_qc_persisted(required) {
1968 return;
1969 }
1970 let (vote2, _) = self.pending_vote2.remove(&view).expect("checked above");
1971 outbox.push_back(ConsensusOutput::SendVote2(vote2));
1972 }
1973
1974 fn vote2_persisted(&self, view: ViewNumber) -> bool {
1978 view <= self.restart_barred_view
1979 || (self.stored_actions.contains(&(view, ActionKind::Vote))
1980 && self.stored_vids.contains(&view))
1981 }
1982
1983 fn release_proposal(&mut self, view: ViewNumber, outbox: &mut Outbox<ConsensusOutput<T>>) {
1984 let Some(message) = self.pending_proposal.get(&view) else {
1985 return;
1986 };
1987 if !self.stored_actions.contains(&(view, ActionKind::Propose))
1988 || !self.is_proposal_stored(view, &proposal_commitment(&message.data))
1989 {
1990 return;
1991 }
1992 let message = self.pending_proposal.remove(&view).expect("checked above");
1993 outbox.push_back(ConsensusOutput::SendProposal(message));
1994 }
1995
1996 fn is_proposal_stored(&self, view: ViewNumber, commitment: &Commitment<Leaf2<T>>) -> bool {
1997 self.stored_proposals
1998 .get(&view)
1999 .is_some_and(|commitments| commitments.contains(commitment))
2000 }
2001
2002 fn request_action(
2003 &mut self,
2004 view: ViewNumber,
2005 epoch: Option<EpochNumber>,
2006 kind: ActionKind,
2007 outbox: &mut Outbox<ConsensusOutput<T>>,
2008 ) {
2009 if self.requested_actions.insert((view, kind)) {
2010 outbox.push_back(ConsensusOutput::RecordAction(view, epoch, kind));
2011 }
2012 }
2013
2014 #[instrument(level = "debug", skip_all)]
2015 fn maybe_vote_1(&mut self, view: ViewNumber, outbox: &mut Outbox<ConsensusOutput<T>>) {
2016 if view <= self.timeout_view {
2017 return;
2018 }
2019 if self.voted_1_views.contains(&view) {
2020 return;
2021 }
2022
2023 let Some(state_commitment) = self.states_verified.get(&view) else {
2024 debug!(%view, "state commitment not available");
2025 return;
2026 };
2027 let Some(proposal) = self.proposals.get(&view) else {
2028 debug!(%view, "proposal not available");
2029 return;
2030 };
2031 let Some(vid_share) = self.vid_shares.get(&view) else {
2032 debug!(%view, "vid share not available");
2033 return;
2034 };
2035
2036 let block_number = proposal.block_header.block_number();
2037 let epoch = proposal.epoch;
2038 let qc_view = proposal.justify_qc.view_number();
2039 let qc_epoch = proposal.justify_qc.epoch();
2040
2041 if proposal.epoch > EpochNumber::genesis()
2045 && is_epoch_transition(block_number, *self.epoch_height)
2046 {
2047 let Some(drb) = self.drb_results.get(&(proposal.epoch + 1)) else {
2048 debug!(%view, block = %block_number, %epoch, "DRB result not yet available, deferring vote");
2049 return;
2050 };
2051 if proposal
2052 .next_drb_result
2053 .is_none_or(|proposed_drb| drb != &proposed_drb)
2054 {
2055 warn!(
2056 %view, block = %block_number, %epoch, %qc_view, ?qc_epoch,
2057 "DRB result does not match proposal, refusing to vote"
2058 );
2059 return;
2060 }
2061 }
2062
2063 if !self.staked_in_epoch(proposal.epoch) {
2064 return;
2065 }
2066
2067 let parent_view = proposal.justify_qc.view_number();
2069
2070 let parent_is_pre_cutover = self.pre_cutover_views.contains(&parent_view);
2072 if parent_view != ViewNumber::genesis()
2073 && !is_last_block(
2074 proposal.block_header.block_number().saturating_sub(1),
2075 *self.epoch_height,
2076 )
2077 {
2078 let Some(prev_proposal) = self.proposals.get(&parent_view) else {
2079 debug!(%view, %parent_view, "proposal not available");
2080 return;
2081 };
2082 let parent_block = prev_proposal.block_header.block_number();
2083 let parent_epoch = prev_proposal.epoch;
2084
2085 if !parent_is_pre_cutover {
2086 let VidCommitment::V2(prev_block_commitment) =
2087 prev_proposal.block_header.payload_commitment()
2088 else {
2089 warn! {
2090 %view, block = %block_number, %epoch,
2091 %parent_view, %parent_block, %parent_epoch,
2092 "prev. proposal payload commitment is not a V2 VID commitment"
2093 }
2094 return;
2095 };
2096 if !self.parent_reconstructed(
2098 parent_view,
2099 prev_block_commitment,
2100 proposal_commitment(prev_proposal),
2101 ) {
2102 debug!(
2103 %view, block = %block_number, %epoch,
2104 %parent_view, %parent_block, %parent_epoch,
2105 "no reconstructed block matching the parent block commitment"
2106 );
2107 return;
2108 }
2109 }
2110
2111 if proposal.justify_qc.data().leaf_commit != proposal_commitment(prev_proposal) {
2112 debug!(
2113 %view, block = %block_number, %epoch,
2114 %parent_view, %parent_block, %parent_epoch,
2115 "justify qc commitment does not match proposal commitment"
2116 );
2117 return;
2118 }
2119 }
2120
2121 let proposal_commit = proposal_commitment(proposal);
2122
2123 if state_commitment != &proposal_commit {
2125 debug!(
2126 %view, block = %block_number, %epoch, %qc_view, ?qc_epoch,
2127 "state commitment does not match proposal commitment"
2128 );
2129 return;
2130 }
2131
2132 let inner_vote = match SimpleVote::create_signed_vote(
2133 QuorumData2 {
2134 leaf_commit: proposal_commit,
2135 epoch: proposal.epoch(),
2136 block_number: Some(proposal.block_header.block_number()),
2137 },
2138 view,
2139 &self.public_key,
2140 &self.private_key,
2141 &self.upgrade_lock,
2142 ) {
2143 Ok(vote) => vote,
2144 Err(err) => {
2145 warn!(%view, %err, "failed to created signed vote for proposal");
2146 return;
2147 },
2148 };
2149
2150 let state_vote = if is_epoch_root(proposal.block_header.block_number(), *self.epoch_height)
2151 {
2152 match self.build_state_vote(proposal) {
2153 Ok(sv) => Some(sv),
2154 Err(err) => {
2155 warn!(%view, %err, "failed to build state vote for epoch-root leaf; skipping vote1");
2156 return;
2157 },
2158 }
2159 } else {
2160 None
2161 };
2162
2163 let vote = Vote1 {
2164 vote: inner_vote,
2165 state_vote,
2166 };
2167 let can_send = self.stored_actions.contains(&(view, ActionKind::Vote))
2168 && self.is_proposal_stored(view, &proposal_commit);
2169 let vid_share = can_send.then(|| vid_share.clone());
2170 self.voted_1_views.insert(view);
2171 if let Some(vid_share) = vid_share {
2172 outbox.push_back(ConsensusOutput::SendVote1(vote));
2173 outbox.push_back(ConsensusOutput::BroadcastVidShare(vid_share));
2174 } else {
2175 self.request_action(view, Some(epoch), ActionKind::Vote, outbox);
2176 self.pending_vote1.insert(view, vote);
2177 }
2178 }
2179
2180 #[instrument(level = "debug", skip_all)]
2181 fn maybe_vote_2_and_update_lock(
2182 &mut self,
2183 view: ViewNumber,
2184 outbox: &mut Outbox<ConsensusOutput<T>>,
2185 ) {
2186 if self.pre_cutover_views.contains(&view) {
2188 return;
2189 }
2190 if self.voted_2_views.contains(&view) {
2191 return;
2192 }
2193 let Some(cert1) = self.certs.get(&view) else {
2194 debug!(%view, "cert1 not available");
2195 return;
2196 };
2197 let Some(proposal) = self.proposals.get(&view) else {
2198 debug!(%view, "proposal not available");
2199 return;
2200 };
2201 let proposal_epoch = proposal.epoch;
2202 let block = proposal.block_header.block_number();
2203 let qc_view = proposal.justify_qc.view_number();
2204 let qc_epoch = proposal.justify_qc.epoch();
2205
2206 let proposal_commit = proposal_commitment(proposal);
2207
2208 if cert1.data.leaf_commit != proposal_commit {
2210 warn!(
2211 %view, %block, epoch = %proposal_epoch, %qc_view, ?qc_epoch,
2212 "cert1 commitment does not match proposal commitment"
2213 );
2214 return;
2215 }
2216 let VidCommitment::V2(proposal_block_commitment) =
2217 proposal.block_header.payload_commitment()
2218 else {
2219 warn!(
2220 %view, %block, epoch = %proposal_epoch, %qc_view, ?qc_epoch,
2221 "proposal payload commitment is not a V2 VID commitment"
2222 );
2223 return;
2224 };
2225 if !self
2226 .blocks_reconstructed
2227 .contains(&(view, proposal_block_commitment))
2228 {
2229 debug!(
2230 %view, %block, epoch = %proposal_epoch, %qc_view, ?qc_epoch,
2231 "no reconstructed block matching the proposal commitment"
2232 );
2233 return;
2234 }
2235
2236 if self
2239 .locked_cert
2240 .as_mut()
2241 .is_none_or(|locked_cert| locked_cert.view_number() < cert1.view_number())
2242 {
2243 self.locked_cert = Some(cert1.clone());
2244 self.current_view = self.current_view.max(view + 1);
2245 self.current_epoch = Some(proposal_epoch);
2246 outbox.push_back(ConsensusOutput::ViewChanged(view + 1, proposal_epoch));
2247 outbox.push_back(ConsensusOutput::SendCertificate1(cert1.clone()));
2248 outbox.push_back(ConsensusOutput::PersistHighQc(cert1.clone()));
2250 }
2251
2252 if self.certs2.contains_key(&view)
2253 || self.decided_views.contains(&view)
2254 || view <= self.decide_floor()
2255 {
2256 return;
2257 }
2258
2259 if !self.staked_in_epoch(proposal_epoch) {
2260 return;
2261 }
2262
2263 let vote = match SimpleVote::create_signed_vote(
2264 Vote2Data {
2265 leaf_commit: proposal_commit,
2266 epoch: proposal_epoch,
2267 block_number: proposal.block_header.block_number(),
2268 },
2269 view,
2270 &self.public_key,
2271 &self.private_key,
2272 &self.upgrade_lock,
2273 ) {
2274 Ok(vote) => vote,
2275 Err(err) => {
2276 warn!(%view, %err, "failed to created signed vote2");
2277 return;
2278 },
2279 };
2280 self.voted_2_views.insert(view);
2281 let required = self
2283 .locked_view()
2284 .expect("locked_cert is set before voting in phase 2");
2285 if self.vote2_persisted(view) && self.high_qc_persisted(required) {
2286 outbox.push_back(ConsensusOutput::SendVote2(vote));
2287 } else {
2288 if view > self.restart_barred_view {
2289 self.request_action(view, Some(proposal_epoch), ActionKind::Vote, outbox);
2290 }
2291 self.pending_vote2.insert(view, (vote, required));
2292 }
2293 }
2294
2295 #[instrument(level = "trace", skip_all)]
2296 fn is_safe(&self, proposal: &Proposal<T>) -> Result<(), SafetyError> {
2297 let Some(locked_cert) = self.locked_cert.as_ref() else {
2298 debug!("at genesis");
2300 return Ok(());
2301 };
2302
2303 if locked_cert.view_number() == proposal.view_number() {
2305 let locked_commit = locked_cert.data.leaf_commit;
2306 let proposal_commit = proposal_commitment(proposal);
2307 if locked_commit != proposal_commit {
2308 return Err(SafetyError::LockedViewCommitmentMismatch {
2309 locked_commit: locked_commit.to_string(),
2310 proposal_commit: proposal_commit.to_string(),
2311 });
2312 }
2313 return Ok(());
2314 }
2315
2316 let parent_commit = proposal
2317 .justify_qc
2318 .data_commitment(&self.upgrade_lock)
2319 .map_err(SafetyError::JustifyQcCommitment)?;
2320 let locked_commit = locked_cert
2321 .data_commitment(&self.upgrade_lock)
2322 .map_err(SafetyError::LockedCertCommitment)?;
2323
2324 let safety = parent_commit == locked_commit;
2325 let liveness = proposal.justify_qc.view_number() > locked_cert.view_number();
2326 if safety || liveness {
2327 return Ok(());
2328 }
2329
2330 Err(SafetyError::UnsafeProposal {
2331 locked_view: locked_cert.view_number(),
2332 parent_commit: parent_commit.to_string(),
2333 locked_commit: locked_commit.to_string(),
2334 })
2335 }
2336
2337 fn leader_label(&self, view: ViewNumber, epoch: EpochNumber) -> String {
2341 match self
2342 .stake_table_coordinator
2343 .membership_for_epoch(Some(epoch))
2344 {
2345 Ok(stake_table) => match stake_table.leader(view) {
2346 Ok(leader) => KeyPrefix::from(&leader).to_string(),
2347 Err(_) => "unknown".to_string(),
2348 },
2349 Err(_) => "unknown".to_string(),
2350 }
2351 }
2352
2353 #[instrument(level = "trace", skip_all)]
2354 fn is_leader(&self, view: ViewNumber, epoch: EpochNumber) -> bool {
2355 match self
2356 .stake_table_coordinator
2357 .membership_for_epoch(Some(epoch))
2358 {
2359 Ok(stake_table) => match stake_table.leader(view) {
2360 Ok(leader) => leader == self.public_key,
2361 Err(err) => {
2362 warn!(%view, %epoch, %err, "failed to get leader from stake table");
2363 false
2364 },
2365 },
2366 Err(err) => {
2367 warn!(%view, %epoch, %err, "failed to get stake table");
2368 false
2369 },
2370 }
2371 }
2372
2373 fn staked_in_epoch(&self, epoch: EpochNumber) -> bool {
2374 match self
2375 .stake_table_coordinator
2376 .membership_for_epoch(Some(epoch))
2377 {
2378 Ok(stake_table) => stake_table.has_stake(&self.public_key),
2379 Err(err) => {
2380 warn!(%epoch, %err, "failed to get stake table");
2381 false
2382 },
2383 }
2384 }
2385
2386 fn missing_for_vote1(&self, view: ViewNumber) -> Vec<&'static str> {
2388 let mut missing = Vec::new();
2389 if !self.states_verified.contains_key(&view) {
2390 missing.push("state_validation");
2391 }
2392 let proposal = self.proposals.get(&view);
2393 if proposal.is_none() {
2394 missing.push("proposal");
2395 }
2396 if !self.vid_shares.contains_key(&view) {
2397 missing.push("vid_share");
2398 }
2399 if let Some(proposal) = proposal {
2400 let block_number = proposal.block_header.block_number();
2401 if proposal.epoch > EpochNumber::genesis()
2402 && is_epoch_transition(block_number, *self.epoch_height)
2403 && !self.drb_results.contains_key(&(proposal.epoch + 1))
2404 {
2405 missing.push("drb_result_for_next_epoch");
2406 }
2407 let parent_view = proposal.justify_qc.view_number();
2411 if parent_view != ViewNumber::genesis()
2412 && !is_last_block(block_number.saturating_sub(1), *self.epoch_height)
2413 {
2414 let reconstructed = self.proposals.get(&parent_view).is_some_and(|p| {
2415 let VidCommitment::V2(c) = p.block_header.payload_commitment() else {
2416 return false;
2417 };
2418 self.parent_reconstructed(parent_view, c, proposal_commitment(p))
2419 });
2420 if !reconstructed {
2421 missing.push("parent_block_reconstructed");
2422 }
2423 }
2424 }
2425 missing
2426 }
2427
2428 fn missing_for_propose(&self, view: ViewNumber) -> Vec<&'static str> {
2430 let mut missing = Vec::new();
2431
2432 let view_change_evidence = self.timeout_certs.get(&view);
2433 let parent_cert = if view_change_evidence.is_some() {
2434 match self.locked_cert.as_ref() {
2435 Some(c) => c,
2436 None => {
2437 missing.push("locked_cert");
2438 return missing;
2439 },
2440 }
2441 } else {
2442 match self.certs.get(&ViewNumber::from(view.saturating_sub(1))) {
2443 Some(c) => c,
2444 None => {
2445 missing.push("parent_cert");
2446 return missing;
2447 },
2448 }
2449 };
2450 let parent_view = parent_cert.view_number();
2451 let Some(parent_proposal) = self.proposals.get(&parent_view) else {
2452 missing.push("parent_proposal");
2453 return missing;
2454 };
2455
2456 let parent_commitment = proposal_commitment(parent_proposal);
2457 let header = self.headers.get(&(view, parent_commitment));
2458 if header.is_none() {
2459 missing.push("block_header");
2460 }
2461 let block_present = header
2462 .and_then(|h| {
2463 if let VidCommitment::V2(c) = h.payload_commitment() {
2464 Some(c)
2465 } else {
2466 None
2467 }
2468 })
2469 .is_some_and(|c| self.blocks.contains_key(&(view, c)));
2470 if !block_present {
2471 missing.push("block_payload");
2472 }
2473
2474 if let Some(header) = header {
2477 let first_proposal_of_epoch =
2478 is_last_block(header.block_number().saturating_sub(1), *self.epoch_height);
2479
2480 if parent_proposal.epoch > EpochNumber::genesis()
2481 && is_epoch_transition(header.block_number(), *self.epoch_height)
2482 && !self
2483 .drb_results
2484 .contains_key(&EpochNumber::new(*parent_proposal.epoch + 1))
2485 {
2486 missing.push("drb_result_for_next_epoch");
2487 }
2488
2489 if first_proposal_of_epoch && !self.certs2.contains_key(&parent_view) {
2490 missing.push("next_epoch_justify_qc");
2491 }
2492 }
2493
2494 let parent_block_number = parent_cert.data.block_number.unwrap_or(0);
2495 if is_epoch_root(parent_block_number, *self.epoch_height) {
2496 match parent_cert.data.epoch() {
2497 None => missing.push("parent_cert_epoch"),
2498 Some(parent_epoch) => {
2499 if !self.state_certs.contains_key(&parent_epoch) {
2500 missing.push("state_cert");
2501 }
2502 },
2503 }
2504 }
2505
2506 missing
2507 }
2508}
2509
2510impl<T: NodeType> ConsensusInput<T> {
2511 fn view_number(&self) -> ViewNumber {
2512 match self {
2513 ConsensusInput::BlockBuilt { view, .. } => *view,
2514 ConsensusInput::BlockReconstructed(view, _) => *view,
2515 ConsensusInput::Certificate1(cert) => cert.view_number(),
2516 ConsensusInput::Certificate2(cert) => cert.view_number(),
2517 ConsensusInput::AdvanceView(cert) => cert.view_number() + 1,
2519 ConsensusInput::EpochRootCertificates { cert1, .. } => cert1.view_number(),
2520 ConsensusInput::HeaderCreated(view, ..) => *view,
2521 ConsensusInput::Proposal(_, prop) => prop.view_number(),
2522 ConsensusInput::VidShare(share) => share.view_number(),
2523 ConsensusInput::FetchedProposal(prop) => prop.view_number(),
2524 ConsensusInput::StateValidated(response) => response.view,
2525 ConsensusInput::StateValidationFailed(request) => request.view,
2526 ConsensusInput::Stored(stored) => stored.view_number(),
2527 ConsensusInput::Timeout(view, _) => *view,
2528 ConsensusInput::TimeoutOneHonest(view, _) => *view,
2529 ConsensusInput::TimeoutCertificate(cert) => {
2530 cert.view_number() + 1
2533 },
2534 ConsensusInput::VidDisperseCreated(view, _) => *view,
2535 ConsensusInput::DrbResult(..) => ViewNumber::genesis(),
2539 ConsensusInput::EpochChange(epoch_change) => epoch_change.cert1.view_number(),
2540 }
2541 }
2542}