1use std::{
9 sync::{
10 atomic::{AtomicBool, Ordering as AtomicOrdering},
11 Arc,
12 },
13 time::Duration,
14};
15
16use parking_lot::RwLock;
17
18use anyhow::{anyhow, bail, ensure, Context, Result};
19use fail::fail_point;
20use futures::{
21 channel::{mpsc, oneshot},
22 SinkExt,
23};
24use serde::Serialize;
25
26use consensus_types::{
27 block::Block,
28 block_retrieval::{BlockRetrievalResponse, BlockRetrievalStatus},
29 common::{Author, Round},
30 proposal_msg::ProposalMsg,
31 quorum_cert::QuorumCert,
32 sync_info::SyncInfo,
33 timeout_certificate::TimeoutCertificate,
34 vote::Vote,
35 vote_msg::VoteMsg,
36};
37use diem_config::keys::ConfigKey;
38use diem_crypto::{hash::CryptoHash, HashValue, SigningKey, VRFPrivateKey};
39use diem_logger::prelude::*;
40use diem_types::{
41 account_address::{from_consensus_public_key, AccountAddress},
42 block_info::PivotBlockDecision,
43 chain_id::ChainId,
44 epoch_state::EpochState,
45 ledger_info::LedgerInfoWithSignatures,
46 transaction::{
47 ConflictSignature, DisputePayload, ElectionPayload, RawTransaction,
48 SignedTransaction, TransactionPayload,
49 },
50 validator_config::{ConsensusPrivateKey, ConsensusVRFPrivateKey},
51 validator_verifier::ValidatorVerifier,
52};
53use safety_rules::SafetyRules;
54
55use crate::pos::{
56 mempool::SubmissionStatus,
57 protocol::message::block_retrieval_response::BlockRetrievalRpcResponse,
58};
59
60use super::{
61 block_storage::{BlockReader, BlockRetriever, BlockStore},
62 error::VerifyError,
63 liveness::{
64 proposal_generator::ProposalGenerator,
65 proposer_election::ProposerElection,
66 round_state::{NewRoundEvent, RoundState},
67 },
68 logging::{LogEvent, LogSchema},
69 network::{
70 ConsensusMsg, ConsensusNetworkSender, IncomingBlockRetrievalRequest,
71 },
72 pending_votes::VoteReceptionResult,
73 persistent_liveness_storage::PersistentLivenessStorage,
74};
75
76#[derive(Serialize, Clone)]
77pub enum UnverifiedEvent {
78 ProposalMsg(Box<ProposalMsg>),
79 VoteMsg(Box<VoteMsg>),
80 SyncInfo(Box<SyncInfo>),
81}
82
83impl UnverifiedEvent {
84 pub fn verify(
85 self, validator: &ValidatorVerifier, epoch_vrf_seed: &[u8],
86 ) -> Result<VerifiedEvent, VerifyError> {
87 Ok(match self {
88 UnverifiedEvent::ProposalMsg(p) => {
89 p.verify(validator, epoch_vrf_seed)?;
90 VerifiedEvent::ProposalMsg(p)
91 }
92 UnverifiedEvent::VoteMsg(v) => {
93 v.verify(validator)?;
94 VerifiedEvent::VoteMsg(v)
95 }
96 UnverifiedEvent::SyncInfo(s) => {
97 s.verify(validator)?;
98 VerifiedEvent::SyncInfo(s)
99 }
100 })
101 }
102
103 pub fn epoch(&self) -> u64 {
104 match self {
105 UnverifiedEvent::ProposalMsg(p) => p.epoch(),
106 UnverifiedEvent::VoteMsg(v) => v.epoch(),
107 UnverifiedEvent::SyncInfo(s) => s.epoch(),
108 }
109 }
110}
111
112impl From<ConsensusMsg> for UnverifiedEvent {
113 fn from(value: ConsensusMsg) -> Self {
114 match value {
115 ConsensusMsg::ProposalMsg(m) => UnverifiedEvent::ProposalMsg(m),
116 ConsensusMsg::VoteMsg(m) => UnverifiedEvent::VoteMsg(m),
117 ConsensusMsg::SyncInfo(m) => UnverifiedEvent::SyncInfo(m),
118 _ => unreachable!("Unexpected conversion"),
119 }
120 }
121}
122
123pub enum VerifiedEvent {
124 ProposalMsg(Box<ProposalMsg>),
125 VoteMsg(Box<VoteMsg>),
126 SyncInfo(Box<SyncInfo>),
127}
128
129pub struct RoundManager {
135 epoch_state: EpochState,
136 block_store: Arc<BlockStore>,
137 round_state: RoundState,
138 proposer_election: Box<dyn ProposerElection + Send + Sync>,
139 proposal_generator: Option<ProposalGenerator>,
141 safety_rules: Arc<RwLock<SafetyRules>>,
142 network: ConsensusNetworkSender,
143 storage: Arc<dyn PersistentLivenessStorage>,
144 sync_only: bool,
145 tx_sender: mpsc::Sender<(
146 SignedTransaction,
147 oneshot::Sender<anyhow::Result<SubmissionStatus>>,
148 )>,
149 chain_id: ChainId,
150
151 is_voting: bool,
152 election_control: Arc<AtomicBool>,
153 consensus_private_key: Option<ConfigKey<ConsensusPrivateKey>>,
154 vrf_private_key: Option<ConfigKey<ConsensusVRFPrivateKey>>,
155}
156
157impl RoundManager {
158 pub fn new(
159 epoch_state: EpochState, block_store: Arc<BlockStore>,
160 round_state: RoundState,
161 proposer_election: Box<dyn ProposerElection + Send + Sync>,
162 proposal_generator: Option<ProposalGenerator>,
163 safety_rules: Arc<RwLock<SafetyRules>>,
164 network: ConsensusNetworkSender,
165 storage: Arc<dyn PersistentLivenessStorage>, sync_only: bool,
166 tx_sender: mpsc::Sender<(
167 SignedTransaction,
168 oneshot::Sender<anyhow::Result<SubmissionStatus>>,
169 )>,
170 chain_id: ChainId, is_voting: bool, election_control: Arc<AtomicBool>,
171 consensus_private_key: Option<ConfigKey<ConsensusPrivateKey>>,
172 vrf_private_key: Option<ConfigKey<ConsensusVRFPrivateKey>>,
173 ) -> Self {
174 Self {
175 epoch_state,
176 block_store,
177 round_state,
178 proposer_election,
179 proposal_generator,
180 is_voting,
181 safety_rules,
182 network,
183 storage,
184 sync_only,
185 tx_sender,
186 chain_id,
187 election_control,
188 consensus_private_key,
189 vrf_private_key,
190 }
191 }
192
193 fn create_block_retriever(&self, author: Author) -> BlockRetriever {
194 BlockRetriever::new(self.network.clone(), author)
195 }
196
197 async fn process_new_round_event(
212 &mut self, new_round_event: NewRoundEvent,
213 ) -> anyhow::Result<()> {
214 diem_debug!(
215 self.new_log(LogEvent::NewRound),
216 reason = new_round_event.reason
217 );
218 if self.proposer_election.is_random_election() {
219 self.proposer_election.next_round(
220 new_round_event.round,
221 self.epoch_state.vrf_seed.clone(),
222 );
223 self.round_state
224 .setup_proposal_timeout(self.epoch_state.epoch);
225 }
226
227 if let Err(e) = self.broadcast_pivot_decision().await {
228 diem_error!("error in broadcasting pivot decision tx: {:?}", e);
229 }
230 if let Err(e) = self.broadcast_election().await {
233 diem_error!("error in broadcasting election tx: {:?}", e);
234 }
235
236 if self.is_validator() {
237 if let Some(ref proposal_generator) = self.proposal_generator {
238 let author = proposal_generator.author();
239
240 if self
241 .proposer_election
242 .is_valid_proposer(author, new_round_event.round)
243 {
244 let proposal_msg = ConsensusMsg::ProposalMsg(Box::new(
245 self.generate_proposal(new_round_event).await?,
246 ));
247 let mut network = self.network.clone();
248 network.broadcast(proposal_msg, vec![]).await;
249 }
250 }
251 }
252 Ok(())
253 }
254
255 pub async fn broadcast_pivot_decision(&mut self) -> anyhow::Result<()> {
256 if !self.is_validator() {
257 return Ok(());
259 }
260 diem_debug!("broadcast_pivot_decision starts");
261
262 let hqc = self.block_store.highest_quorum_cert();
263 let parent_block = hqc.certified_block();
264 if self.block_store.path_from_root(parent_block.id()).is_none() {
266 bail!("HQC {} already pruned", parent_block);
267 }
268
269 let parent_decision = parent_block
272 .pivot_decision()
273 .map(|d| d.block_hash)
274 .unwrap_or_default();
275 let pivot_decision = match self
276 .block_store
277 .pow_handler
278 .next_pivot_decision(parent_decision)
279 .await
280 {
281 Some(res) => res,
282 None => {
283 diem_debug!("No new pivot decision");
285 return Ok(());
286 }
287 };
288
289 let proposal_generator =
290 self.proposal_generator.as_ref().expect("checked");
291 diem_info!("Broadcast new pivot decision: {:?}", pivot_decision);
292 let raw_tx = RawTransaction::new_pivot_decision(
295 proposal_generator.author(),
296 PivotBlockDecision {
297 block_hash: pivot_decision.1,
298 height: pivot_decision.0,
299 },
300 self.chain_id,
301 );
302 let signed_tx =
303 raw_tx.sign(&proposal_generator.private_key)?.into_inner();
304 let (tx, rx) = oneshot::channel();
305 self.tx_sender.send((signed_tx, tx)).await?;
306 rx.await??;
308 diem_debug!("broadcast_pivot_decision sends");
309 Ok(())
310 }
311
312 pub async fn broadcast_election(&mut self) -> anyhow::Result<()> {
313 if !self.election_control.load(AtomicOrdering::Relaxed) {
314 diem_debug!("Skip election for election_control");
315 return Ok(());
316 }
317 if !self.is_voting {
318 return Ok(());
320 }
321 if !self.block_store.pow_handler.is_normal_phase() {
322 return Ok(());
326 }
327 if self.vrf_private_key.is_none()
328 || self.consensus_private_key.is_none()
329 {
330 diem_warn!("broadcast_election without keys");
331 return Ok(());
332 }
333 let private_key = self.consensus_private_key.as_ref().unwrap();
334 let vrf_private_key = self.vrf_private_key.as_ref().unwrap();
335 let author = from_consensus_public_key(
336 &private_key.public_key(),
337 &vrf_private_key.public_key(),
338 );
339 diem_debug!("broadcast_election starts");
340 let pos_state = self.storage.pos_ledger_db().get_latest_pos_state();
341 if let Some(target_term) = pos_state.next_elect_term(&author) {
342 let epoch_vrf_seed = pos_state.target_term_seed(target_term);
343 let election_payload = ElectionPayload {
344 public_key: private_key.public_key(),
345 vrf_public_key: vrf_private_key.public_key(),
346 target_term,
347 vrf_proof: vrf_private_key
348 .private_key()
349 .compute(epoch_vrf_seed.as_slice())
350 .unwrap(),
351 };
352 let raw_tx = RawTransaction::new_election(
353 author,
354 election_payload,
355 self.chain_id,
356 );
357 let signed_tx =
358 raw_tx.sign(&private_key.private_key())?.into_inner();
359 let (tx, rx) = oneshot::channel();
360 self.tx_sender.send((signed_tx, tx)).await?;
361 rx.await??;
363 diem_debug!(
364 "broadcast_election sends: target_term={}",
365 target_term
366 );
367 } else {
368 diem_debug!("Skip election for elected");
369 if let Some(node_data) = pos_state.account_node_data(author) {
370 if node_data.lock_status().force_retired().is_some() {
371 warn!("The node stops elections for force retire!");
372 }
373 }
374 }
375 Ok(())
376 }
377
378 async fn generate_proposal(
379 &mut self, new_round_event: NewRoundEvent,
380 ) -> anyhow::Result<ProposalMsg> {
381 let proposal = self
384 .proposal_generator
385 .as_mut()
386 .expect("checked by process_new_round_event")
387 .generate_proposal(
388 new_round_event.round,
389 self.epoch_state.verifier().clone(),
390 )
391 .await?;
392 let mut signed_proposal =
393 self.safety_rules.write().sign_proposal(proposal)?;
394 if self.proposer_election.is_random_election() {
395 signed_proposal.set_vrf_nonce_and_proof(
396 self.proposer_election
397 .gen_vrf_nonce_and_proof(signed_proposal.block_data())
398 .expect("threshold checked in is_valid_proposer"),
399 )
400 }
401 diem_debug!(self.new_log(LogEvent::Propose), "{}", signed_proposal);
402 Ok(ProposalMsg::new(
404 signed_proposal,
405 self.block_store.sync_info(),
406 ))
407 }
408
409 pub async fn process_proposal_msg(
413 &mut self, proposal_msg: ProposalMsg,
414 ) -> anyhow::Result<()> {
415 fail_point!("consensus::process_proposal_msg", |_| {
416 Err(anyhow::anyhow!("Injected error in process_proposal_msg"))
417 });
418
419 let proposer = proposal_msg.proposer().ok_or_else(|| {
420 anyhow::anyhow!(
421 "ProposalMsg without author reached process_proposal_msg \
422 — verify_well_formed should have rejected it"
423 )
424 })?;
425 if self
426 .ensure_round_and_sync_up(
427 proposal_msg.proposal().round(),
428 proposal_msg.sync_info(),
429 proposer,
430 true,
431 )
432 .await
433 .context("[RoundManager] Process proposal")?
434 {
435 if self
436 .process_proposal(proposal_msg.clone().take_proposal())
437 .await?
438 {
439 let exclude = vec![proposer, self.network.author];
455 self.network
456 .broadcast(
457 ConsensusMsg::ProposalMsg(Box::new(proposal_msg)),
458 exclude,
459 )
460 .await;
461 }
462 Ok(())
463 } else {
464 bail!(
465 "Stale proposal {}, current round {}",
466 proposal_msg.proposal(),
467 self.round_state.current_round()
468 );
469 }
470 }
471
472 pub async fn sync_up(
476 &mut self, sync_info: &SyncInfo, author: Author, help_remote: bool,
477 ) -> anyhow::Result<()> {
478 let local_sync_info = self.block_store.sync_info();
479 if help_remote && local_sync_info.has_newer_certificates(&sync_info) {
480 diem_debug!(
481 self.new_log(LogEvent::HelpPeerSync).remote_peer(author),
482 "Remote peer has stale state {}, send it back {}",
483 sync_info,
484 local_sync_info,
485 );
486 self.network.send_sync_info(local_sync_info.clone(), author);
487 }
488 if sync_info.has_newer_certificates(&local_sync_info) {
489 diem_debug!(
490 self.new_log(LogEvent::SyncToPeer).remote_peer(author),
491 "Local state {} is stale than remote state {}",
492 local_sync_info,
493 sync_info
494 );
495 sync_info
498 .verify(&self.epoch_state().verifier())
499 .map_err(|e| {
500 diem_error!(
501 SecurityEvent::InvalidSyncInfoMsg,
502 sync_info = sync_info,
503 remote_peer = author,
504 error = ?e,
505 );
506 VerifyError::from(e)
507 })?;
508 let mut retriever = self.create_block_retriever(author);
516 self.block_store
517 .insert_quorum_cert(
518 &sync_info.highest_commit_cert(),
519 &mut retriever,
520 )
521 .await?;
522 self.block_store
523 .insert_quorum_cert(
524 &sync_info.highest_quorum_cert(),
525 &mut retriever,
526 )
527 .await?;
528 if let Some(tc) = sync_info.highest_timeout_certificate() {
529 self.block_store
530 .insert_timeout_certificate(Arc::new(tc.clone()))?;
531 }
532 self.process_certificates().await
533 } else {
534 Ok(())
535 }
536 }
537
538 pub async fn sync_to_ledger_info(
540 &mut self, ledger_info: &LedgerInfoWithSignatures,
541 peer_id: AccountAddress,
542 ) -> Result<()> {
543 diem_debug!("sync_to_ledger_info: {:?}", ledger_info);
544 let mut retriever = self.create_block_retriever(peer_id);
545 if !self
546 .block_store
547 .block_exists(ledger_info.ledger_info().consensus_block_id())
548 {
549 let block_for_ledger_info = retriever
550 .retrieve_block_for_ledger_info(ledger_info)
551 .await?;
552 self.block_store
553 .insert_quorum_cert(
554 block_for_ledger_info.quorum_cert(),
555 &mut retriever,
556 )
557 .await?;
558 self.block_store.execute_and_insert_block(
559 block_for_ledger_info,
560 true,
561 false,
562 )?;
563 };
564 self.block_store.commit(ledger_info.clone()).await
565 }
566
567 pub async fn ensure_round_and_sync_up(
576 &mut self, message_round: Round, sync_info: &SyncInfo, author: Author,
577 help_remote: bool,
578 ) -> anyhow::Result<bool> {
579 if message_round < self.round_state.current_round() {
580 return Ok(false);
581 }
582 self.sync_up(sync_info, author, help_remote).await?;
583 ensure!(
584 message_round == self.round_state.current_round(),
585 "After sync, round {} doesn't match local {}",
586 message_round,
587 self.round_state.current_round()
588 );
589 Ok(true)
590 }
591
592 pub async fn process_sync_info_msg(
594 &mut self, sync_info: SyncInfo, peer: Author,
595 ) -> anyhow::Result<()> {
596 fail_point!("consensus::process_sync_info_msg", |_| {
597 Err(anyhow::anyhow!("Injected error in process_sync_info_msg"))
598 });
599 diem_debug!(
600 self.new_log(LogEvent::ReceiveSyncInfo).remote_peer(peer),
601 "{}",
602 sync_info
603 );
604 self.ensure_round_and_sync_up(
607 sync_info
608 .highest_round()
609 .checked_add(1)
610 .ok_or_else(|| anyhow!("round overflow"))?,
611 &sync_info,
612 peer,
613 false,
614 )
615 .await
616 .context("[RoundManager] Failed to process sync info msg")?;
617 Ok(())
618 }
619
620 pub async fn process_new_round_timeout(
621 &mut self, epoch_round: (u64, Round),
622 ) -> anyhow::Result<()> {
623 diem_debug!("process_new_round_timeout: round={:?}", epoch_round);
624 if epoch_round
625 != (self.epoch_state.epoch, self.round_state.current_round())
626 {
627 return Ok(());
628 }
629 let round = epoch_round.1;
630
631 match self
632 .round_state
633 .get_round_certificate(&self.epoch_state.verifier())
634 {
635 VoteReceptionResult::NewQuorumCertificate(qc) => {
636 self.new_qc_aggregated(
637 qc.clone(),
638 qc.ledger_info()
639 .signatures()
640 .keys()
641 .next()
642 .expect("qc formed")
643 .clone(),
644 )
645 .await?;
646 }
647 VoteReceptionResult::NewTimeoutCertificate(tc) => {
648 self.new_tc_aggregated(tc).await?;
649 }
650 _ => {
651 anyhow::bail!(
653 "New round timeout without new certificate! round={}",
654 round
655 );
656 }
657 }
658 Ok(())
659 }
660
661 pub async fn process_local_timeout(
671 &mut self, epoch_round: (u64, Round),
672 ) -> anyhow::Result<()> {
673 diem_debug!("process_local_timeout: round={:?}", epoch_round);
674 if epoch_round
675 != (self.epoch_state.epoch, self.round_state.current_round())
676 {
677 return Ok(());
678 }
679 let round = epoch_round.1;
680
681 if !self.round_state.process_local_timeout(epoch_round) {
682 return Ok(());
683 }
684
685 self.network
686 .broadcast(
687 ConsensusMsg::SyncInfo(Box::new(self.block_store.sync_info())),
688 vec![],
689 )
690 .await;
691
692 match self
693 .round_state
694 .get_round_certificate(&self.epoch_state.verifier())
695 {
696 VoteReceptionResult::NewQuorumCertificate(_)
697 | VoteReceptionResult::NewTimeoutCertificate(_) => {
698 return Ok(());
700 }
701 _ => {
702 }
704 }
705
706 if !self.is_validator() {
707 return Ok(());
708 }
709
710 let (use_last_vote, mut timeout_vote) =
711 match self.round_state.vote_sent() {
712 Some(vote) if vote.vote_data().proposed().round() == round => {
713 (true, vote)
714 }
715 _ => {
716 let nil_block = self
718 .proposal_generator
719 .as_ref()
720 .expect("checked in is_validator")
721 .generate_nil_block(round)?;
722 diem_debug!(
723 self.new_log(LogEvent::VoteNIL),
724 "Planning to vote for a NIL block {}",
725 nil_block
726 );
727 let nil_vote = self.execute_and_vote(nil_block).await?;
728 (false, nil_vote)
729 }
730 };
731
732 if !timeout_vote.is_timeout() {
733 let timeout = timeout_vote.timeout();
734 let signature = self
735 .safety_rules
736 .write()
737 .sign_timeout(&timeout)
738 .context("[RoundManager] SafetyRules signs timeout")?;
739 timeout_vote.add_timeout_signature(signature);
740 }
741
742 self.round_state.record_vote(timeout_vote.clone());
743 let timeout_vote_msg = ConsensusMsg::VoteMsg(Box::new(VoteMsg::new(
744 timeout_vote,
745 self.block_store.sync_info(),
746 )));
747 self.network.broadcast(timeout_vote_msg, vec![]).await;
748 diem_error!(
749 round = round,
750 voted = use_last_vote,
751 event = LogEvent::Timeout,
752 );
753 bail!("Round {} timeout, broadcast to all peers", round);
754 }
755
756 pub async fn process_proposal_timeout(
757 &mut self, epoch_round: (u64, Round),
758 ) -> anyhow::Result<()> {
759 diem_debug!("process_proposal_timeout: round={:?}", epoch_round);
760 if epoch_round
761 != (self.epoch_state.epoch, self.round_state.current_round())
762 {
763 return Ok(());
764 }
765 let round = epoch_round.1;
766
767 if let Some(proposal) = self.proposer_election.choose_proposal_to_vote()
768 {
769 if self.is_validator() {
770 let vote = self
772 .execute_and_vote(proposal)
773 .await
774 .context("[RoundManager] Process proposal")?;
775 diem_debug!(self.new_log(LogEvent::Vote), "{}", vote);
776
777 self.round_state.record_vote(vote.clone());
778 let vote_msg = VoteMsg::new(vote, self.block_store.sync_info());
779 self.network
780 .broadcast(
781 ConsensusMsg::VoteMsg(Box::new(vote_msg)),
782 vec![],
783 )
784 .await;
785 Ok(())
786 } else {
787 self.block_store
789 .execute_and_insert_block(proposal, false, false)
790 .context(
791 "[RoundManager] Failed to execute_and_insert the block",
792 )?;
793 Ok(())
794 }
795 } else {
796 debug!("No proposal to vote: round={}", round);
797 self.process_local_timeout(epoch_round).await
799 }
800 }
801
802 async fn process_certificates(&mut self) -> anyhow::Result<()> {
805 let sync_info = self.block_store.sync_info();
806 if let Some(new_round_event) =
807 self.round_state.process_certificates(sync_info)
808 {
809 self.process_new_round_event(new_round_event).await?;
810 }
811 Ok(())
812 }
813
814 async fn process_proposal(&mut self, proposal: Block) -> Result<bool> {
823 let author = proposal
824 .author()
825 .expect("Proposal should be verified having an author");
826
827 diem_info!(
828 self.new_log(LogEvent::ReceiveProposal).remote_peer(author),
829 block_hash = proposal.id(),
830 block_parent_hash = proposal.quorum_cert().certified_block().id(),
831 );
832
833 ensure!(
834 self.proposer_election.is_valid_proposal(&proposal),
835 "[RoundManager] Proposer {} for block {} is not a valid proposer for this round",
836 author,
837 proposal,
838 );
839
840 let block_time_since_epoch =
841 Duration::from_micros(proposal.timestamp_usecs());
842
843 ensure!(
844 block_time_since_epoch < self.round_state.current_round_deadline(),
845 "[RoundManager] Waiting until proposal block timestamp usecs {:?} \
846 would exceed the round duration {:?}, hence will not vote for this round",
847 block_time_since_epoch,
848 self.round_state.current_round_deadline(),
849 );
850
851 if self.proposer_election.is_random_election() {
852 if self
853 .proposer_election
854 .receive_proposal_candidate(&proposal)?
855 {
856 self.block_store.execute_and_insert_block(
857 proposal.clone(),
858 false,
859 false,
860 )?;
861 self.proposer_election.set_proposal_candidate(proposal);
862 Ok(true)
863 } else {
864 Ok(false)
869 }
870 } else {
871 bail!("unsupported election rules")
872 }
873 }
874
875 async fn execute_and_vote(
882 &mut self, proposed_block: Block,
883 ) -> anyhow::Result<Vote> {
884 let executed_block = self
885 .block_store
886 .execute_and_insert_block(proposed_block, false, false)
887 .context("[RoundManager] Failed to execute_and_insert the block")?;
888
889 ensure!(
891 self.round_state.vote_sent().is_none(),
892 "[RoundManager] Already vote on this round {}",
893 self.round_state.current_round()
894 );
895
896 ensure!(
897 !self.sync_only,
898 "[RoundManager] sync_only flag is set, stop voting"
899 );
900
901 let maybe_signed_vote_proposal =
902 executed_block.maybe_signed_vote_proposal();
903 let vote = self
904 .safety_rules
905 .write()
906 .construct_and_sign_vote(&maybe_signed_vote_proposal)
907 .context(format!(
908 "[RoundManager] SafetyRules Rejected {}",
909 executed_block.block()
910 ))?;
911 self.storage
912 .save_vote(&vote)
913 .context("[RoundManager] Fail to persist last vote")?;
914
915 Ok(vote)
916 }
917
918 pub async fn process_vote_msg(
925 &mut self, vote_msg: VoteMsg,
926 ) -> anyhow::Result<()> {
927 fail_point!("consensus::process_vote_msg", |_| {
928 Err(anyhow::anyhow!("Injected error in process_vote_msg"))
929 });
930 if self
931 .ensure_round_and_sync_up(
932 vote_msg.vote().vote_data().proposed().round(),
933 vote_msg.sync_info(),
934 vote_msg.vote().author(),
935 true,
936 )
937 .await
938 .context("[RoundManager] Stop processing vote")?
939 {
940 let relay = self
941 .process_vote(vote_msg.vote())
942 .await
943 .context("[RoundManager] Add a new vote")?;
944 if relay {
945 let exclude =
946 vec![vote_msg.vote().author(), self.network.author];
947 self.network
948 .broadcast(
949 ConsensusMsg::VoteMsg(Box::new(vote_msg)),
950 exclude,
951 )
952 .await;
953 }
954 }
955 Ok(())
956 }
957
958 async fn process_vote(&mut self, vote: &Vote) -> anyhow::Result<bool> {
965 diem_info!(
966 self.new_log(LogEvent::ReceiveVote)
967 .remote_peer(vote.author()),
968 vote = %vote,
969 vote_epoch = vote.vote_data().proposed().epoch(),
970 vote_round = vote.vote_data().proposed().round(),
971 vote_id = vote.vote_data().proposed().id(),
972 vote_state = vote.vote_data().proposed().executed_state_id(),
973 );
974
975 let mut relay = true;
977 match self
978 .round_state
979 .insert_vote(vote, &self.epoch_state.verifier())
980 {
981 VoteReceptionResult::NewQuorumCertificate(_)
982 | VoteReceptionResult::NewTimeoutCertificate(_) => {
983 self.round_state
986 .setup_new_round_timeout(self.epoch_state.epoch);
987 }
988 VoteReceptionResult::VoteAdded(_) => {}
989 VoteReceptionResult::DuplicateVote => {
990 relay = false;
993 }
994 VoteReceptionResult::EquivocateVote((vote1, vote2)) => {
995 match &self.proposal_generator {
999 Some(proposal_generator) => {
1000 ensure!(
1001 vote1.author() == vote2.author(),
1002 "incorrect author"
1003 );
1004 ensure!(
1005 vote1.vote_data().proposed().round()
1006 == vote2.vote_data().proposed().round(),
1007 "incorrect round"
1008 );
1009 diem_warn!("Find Equivocate Vote!!! author={}, vote1={:?}, vote2={:?}", vote.author(), vote1, vote2);
1010 let dispute_payload = DisputePayload {
1011 address: vote1.author(),
1012 bls_pub_key: self
1013 .epoch_state
1014 .verifier()
1015 .get_public_key(&vote1.author())
1016 .expect("checked in verify"),
1017 vrf_pub_key: self
1018 .epoch_state
1019 .verifier()
1020 .get_vrf_public_key(&vote1.author())
1021 .expect("checked in verify")
1022 .unwrap(),
1023 conflicting_votes: ConflictSignature::Vote((
1024 bcs::to_bytes(&vote1).expect("encoding error"),
1025 bcs::to_bytes(&vote2).expect("encoding error"),
1026 )),
1027 };
1028 let raw_tx = RawTransaction::new_dispute(
1029 proposal_generator.author(),
1030 dispute_payload,
1031 );
1032 let signed_tx = raw_tx
1033 .sign(&proposal_generator.private_key)?
1034 .into_inner();
1035 let (tx, rx) = oneshot::channel();
1038 self.tx_sender.send((signed_tx, tx)).await?;
1039 rx.await??;
1040 }
1041 None => {}
1042 }
1043 bail!("EquivocateVote!")
1044 }
1045 r => bail!("vote not added with result {:?}", r),
1048 }
1049 Ok(relay)
1050 }
1051
1052 async fn new_qc_aggregated(
1053 &mut self, qc: Arc<QuorumCert>, preferred_peer: Author,
1054 ) -> anyhow::Result<()> {
1055 let result = self
1056 .block_store
1057 .insert_quorum_cert(
1058 &qc,
1059 &mut self.create_block_retriever(preferred_peer),
1060 )
1061 .await
1062 .context("[RoundManager] Failed to process a newly aggregated QC");
1063 self.process_certificates().await?;
1064 result
1065 }
1066
1067 async fn new_tc_aggregated(
1068 &mut self, tc: Arc<TimeoutCertificate>,
1069 ) -> anyhow::Result<()> {
1070 let result = self
1071 .block_store
1072 .insert_timeout_certificate(tc.clone())
1073 .context("[RoundManager] Failed to process a newly aggregated TC");
1074 self.process_certificates().await?;
1075 result
1076 }
1077
1078 pub async fn process_block_retrieval(
1085 &self, request: IncomingBlockRetrievalRequest,
1086 ) -> anyhow::Result<()> {
1087 fail_point!("consensus::process_block_retrieval", |_| {
1088 Err(anyhow::anyhow!("Injected error in process_block_retrieval"))
1089 });
1090 let mut blocks = vec![];
1091 let mut status = BlockRetrievalStatus::Succeeded;
1092 let mut id = request.req.block_id();
1093 while (blocks.len() as u64) < request.req.num_blocks() {
1094 if let Some(executed_block) = self.block_store.get_block(id) {
1095 if executed_block.block().is_genesis_block() {
1096 break;
1097 }
1098 id = executed_block.parent_id();
1099 blocks.push(executed_block.block().clone());
1100 } else if let Ok(Some(block)) =
1101 self.block_store.get_ledger_block(&id)
1102 {
1103 if block.is_genesis_block() {
1104 break;
1105 }
1106 id = block.parent_id();
1107 blocks.push(block);
1108 } else {
1109 break;
1112 }
1113 }
1114
1115 if blocks.is_empty() {
1116 status = BlockRetrievalStatus::IdNotFound;
1117 }
1118
1119 let response = BlockRetrievalRpcResponse {
1120 request_id: request.request_id,
1121 response: BlockRetrievalResponse::new(status, blocks),
1122 };
1123 self.network
1124 .network_sender()
1125 .send_message_with_peer_id(&request.peer_id, &response)?;
1126 Ok(())
1127 }
1128
1129 pub async fn start(&mut self, last_vote_sent: Option<Vote>) {
1131 let new_round_event = self
1132 .round_state
1133 .process_certificates(self.block_store.sync_info())
1134 .expect(
1135 "Can not jump start a round_state from existing certificates.",
1136 );
1137 if let Some(vote) = last_vote_sent {
1138 self.round_state.record_vote(vote);
1139 }
1140 if let Err(e) = self.process_new_round_event(new_round_event).await {
1141 diem_error!(error = ?e, "[RoundManager] Error during start");
1142 }
1143 }
1144
1145 pub fn epoch_state(&self) -> &EpochState { &self.epoch_state }
1146
1147 pub fn round_state(&self) -> &RoundState { &self.round_state }
1148
1149 fn new_log(&self, event: LogEvent) -> LogSchema {
1150 LogSchema::new(event)
1151 .round(self.round_state.current_round())
1152 .epoch(self.epoch_state.epoch)
1153 }
1154
1155 fn is_validator(&self) -> bool {
1156 let r = self.proposal_generator.is_some();
1157 diem_debug!("Check validator: r={} is_voting={}", r, self.is_voting);
1158 r && self.is_voting
1159 }
1160
1161 pub fn filter_proposal(&self, p: &ProposalMsg) -> bool {
1163 self.proposer_election
1164 .receive_proposal_candidate(p.proposal())
1165 .unwrap_or(false)
1166 }
1167
1168 pub fn filter_vote(&self, v: &VoteMsg) -> bool {
1170 !self.round_state.vote_received(v.vote())
1171 }
1172}
1173
1174impl RoundManager {
1176 pub async fn force_vote_proposal(
1180 &mut self, block_id: HashValue, author: Author,
1181 private_key: &ConsensusPrivateKey,
1182 ) -> Result<()> {
1183 let proposal = self
1184 .block_store
1185 .get_block(block_id)
1186 .ok_or(anyhow::anyhow!("force sign block not received"))?;
1187 let vote_proposal = proposal.maybe_signed_vote_proposal().vote_proposal;
1188 let vote_data =
1189 SafetyRules::extension_check(&vote_proposal).map_err(|e| {
1190 anyhow::anyhow!("extension_check error: err={:?}", e)
1191 })?;
1192 let ledger_info = SafetyRules::construct_ledger_info(
1193 vote_proposal.block(),
1194 vote_data.hash(),
1195 )
1196 .map_err(|e| anyhow::anyhow!("extension_check error: err={:?}", e))?;
1197 let signature = private_key.sign(&ledger_info);
1198 let vote =
1199 Vote::new_with_signature(vote_data, author, ledger_info, signature);
1200 let vote_msg = VoteMsg::new(vote, self.block_store.sync_info());
1201 diem_debug!("force_vote_proposal: broadcast {:?}", vote_msg);
1202 self.network
1203 .broadcast(ConsensusMsg::VoteMsg(Box::new(vote_msg)), vec![])
1204 .await;
1205 Ok(())
1206 }
1207
1208 pub async fn force_propose(
1212 &mut self, round: Round, parent_block_id: HashValue,
1213 payload: Vec<TransactionPayload>, private_key: &ConsensusPrivateKey,
1214 ) -> Result<()> {
1215 let parent_qc = self
1216 .block_store
1217 .get_quorum_cert_for_block(parent_block_id)
1218 .ok_or(anyhow::anyhow!(
1219 "no QC for parent: {:?}",
1220 parent_block_id
1221 ))?;
1222 let block_data = self
1223 .proposal_generator
1224 .as_ref()
1225 .ok_or(anyhow::anyhow!("proposal generator is None"))?
1226 .force_propose(round, parent_qc, payload)?;
1227 let signature = private_key.sign(&block_data);
1228 let mut signed_proposal =
1229 Block::new_proposal_from_block_data_and_signature(
1230 block_data, signature, None,
1231 );
1232 signed_proposal.set_vrf_nonce_and_proof(
1235 self.proposer_election
1236 .gen_vrf_nonce_and_proof(signed_proposal.block_data())
1237 .ok_or(anyhow::anyhow!(
1238 "The proposer should not propose in this round"
1239 ))?,
1240 );
1241 let proposal_msg =
1244 ProposalMsg::new(signed_proposal, self.block_store.sync_info());
1245 diem_debug!("force_propose: broadcast {:?}", proposal_msg);
1246 self.network
1247 .broadcast(
1248 ConsensusMsg::ProposalMsg(Box::new(proposal_msg)),
1249 vec![],
1250 )
1251 .await;
1252 Ok(())
1253 }
1254
1255 pub async fn force_sign_pivot_decision(
1256 &mut self, pivot_decision: PivotBlockDecision,
1257 ) -> anyhow::Result<()> {
1258 let proposal_generator = self.proposal_generator.as_ref().ok_or(
1259 anyhow::anyhow!("Non-validator cannot sign pivot decision"),
1260 )?;
1261 diem_info!("force_sign_pivot_decision: {:?}", pivot_decision);
1262 let raw_tx = RawTransaction::new_pivot_decision(
1265 proposal_generator.author(),
1266 pivot_decision,
1267 self.chain_id,
1268 );
1269 let signed_tx =
1270 raw_tx.sign(&proposal_generator.private_key)?.into_inner();
1271 let (tx, rx) = oneshot::channel();
1272 self.tx_sender.send((signed_tx, tx)).await?;
1273 rx.await??;
1275 diem_debug!("force_sign_pivot_decision sends");
1276 Ok(())
1277 }
1278
1279 pub fn get_chosen_proposal(&self) -> anyhow::Result<Option<Block>> {
1280 let chosen = self.proposer_election.choose_proposal_to_vote();
1283 Ok(chosen)
1284 }
1285
1286 pub fn start_voting(&mut self, initialize: bool) -> anyhow::Result<()> {
1287 if !initialize {
1288 self.safety_rules
1289 .write()
1290 .start_voting(initialize)
1291 .map_err(anyhow::Error::from)?;
1292 }
1293 self.is_voting = true;
1294 Ok(())
1295 }
1296
1297 pub fn stop_voting(&mut self) -> anyhow::Result<()> {
1298 self.safety_rules
1299 .write()
1300 .stop_voting()
1301 .map_err(anyhow::Error::from)?;
1302 self.is_voting = false;
1303 Ok(())
1304 }
1305}