cfxcore/pos/consensus/
round_manager.rs

1// Copyright (c) The Diem Core Contributors
2// SPDX-License-Identifier: Apache-2.0
3
4// Copyright 2021 Conflux Foundation. All rights reserved.
5// Conflux is free software and distributed under GNU General Public License.
6// See http://www.gnu.org/licenses/
7
8use 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
129/// Consensus SMR is working in an event based fashion: RoundManager is
130/// responsible for processing the individual events (e.g., process_new_round,
131/// process_proposal, process_vote, etc.). It is exposing the async processing
132/// functions for each event type. The caller is responsible for running the
133/// event loops and driving the execution via some executors.
134pub struct RoundManager {
135    epoch_state: EpochState,
136    block_store: Arc<BlockStore>,
137    round_state: RoundState,
138    proposer_election: Box<dyn ProposerElection + Send + Sync>,
139    // None if this is not a validator.
140    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    /// Leader:
198    ///
199    /// This event is triggered by a new quorum certificate at the previous
200    /// round or a timeout certificate at the previous round.  In either
201    /// case, if this replica is the new proposer for this round, it is
202    /// ready to propose and guarantee that it can create a proposal
203    /// that all honest replicas can vote for.  While this method should only be
204    /// invoked at most once per round, we ensure that only at most one
205    /// proposal can get generated per round to avoid accidental
206    /// equivocation of proposals.
207    ///
208    /// Replica:
209    ///
210    /// Do nothing
211    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        // After the election transaction has been packed and executed,
231        // `broadcast_election` will be a no-op.
232        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            // Not an active validator, so do not need to sign pivot decision.
258            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        // TODO(lpl): Check if this may happen.
265        if self.block_store.path_from_root(parent_block.id()).is_none() {
266            bail!("HQC {} already pruned", parent_block);
267        }
268
269        // Sending non-existent H256 (default) will return the latest pivot
270        // decision.
271        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                // No new pivot decision.
284                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        // It's allowed for a node to sign conflict pivot decision,
293        // so we do not need to persist this signing event.
294        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        // TODO(lpl): Check if we want to wait here.
307        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            // This node does not participate in any signing or voting.
319            return Ok(());
320        }
321        if !self.block_store.pow_handler.is_normal_phase() {
322            // Do not start election before PoW enters normal phase so we will
323            // not be force retired unexpectedly because we are
324            // elected but cannot vote.
325            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            // TODO(lpl): Check if we want to wait here.
362            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        // Proposal generator will ensure that at most one proposal is generated
382        // per round
383        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        // return proposal
403        Ok(ProposalMsg::new(
404            signed_proposal,
405            self.block_store.sync_info(),
406        ))
407    }
408
409    /// Process the proposal message:
410    /// 1. ensure after processing sync info, we're at the same round as the
411    /// proposal 2. execute and decide whether to vode for the proposal
412    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                // If a proposal has been received and voted, it will return
440                // error or false.
441                //
442                // 1. For old leader elections where there is only one leader
443                // and we vote after receiving the first
444                // proposal, the error is returned in
445                // `execute_and_vote` because `vote_sent.
446                // is_none()` is false. 2. For VRF leader election, we
447                // return Ok(false) when we insert a proposal from the same
448                // author to proposal_candidates.
449                //
450                // This ensures that there is no broadcast storm
451                // because we only broadcast a proposal when we receive it for
452                // the first time.
453                // TODO(lpl): Do not send to the sender and the original author.
454                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    /// Sync to the sync info sending from peer if it has newer certificates, if
473    /// we have newer certificates and help_remote is set, send it back the
474    /// local sync info.
475    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            // Some information in SyncInfo is ahead of what we have locally.
496            // First verify the SyncInfo (didn't verify it in the yet).
497            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            /*
509            let result = self
510                .block_store
511                .add_certs(&sync_info, self.create_block_retriever(author))
512                .await;
513             */
514            // TODO(lpl): Ensure this does not cause OOM.
515            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    /// This can only be used in `EpochManager.start_new_epoch`.
539    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    /// The function makes sure that it ensures the message_round equal to what
568    /// we have locally, brings the missing dependencies from the QC and
569    /// LedgerInfo of the given sync info and update the round_state with
570    /// the certificates if succeed. Returns Ok(true) if the sync succeeds
571    /// and the round matches so we can process further. Returns Ok(false)
572    /// if the message is stale. Returns Error in case sync mgr failed to
573    /// bring the missing dependencies. We'll try to help the remote if the
574    /// SyncInfo lags behind and the flag is set.
575    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    /// Process the SyncInfo sent by peers to catch up to latest state.
593    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        // To avoid a ping-pong cycle between two peers that move forward
605        // together.
606        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                // No certificate formed. This should not happen.
652                anyhow::bail!(
653                    "New round timeout without new certificate! round={}",
654                    round
655                );
656            }
657        }
658        Ok(())
659    }
660
661    /// The replica broadcasts a "timeout vote message", which includes the
662    /// round signature, which can be aggregated to a TimeoutCertificate.
663    /// The timeout vote message can be one of the following three options:
664    /// 1) In case a validator has previously voted in this round, it repeats
665    /// the same vote and sign a timeout.
666    /// 2) Otherwise vote for a NIL block and sign a timeout.
667    /// Note this function returns Err even if messages are broadcasted
668    /// successfully because timeout is considered as error. It only returns
669    /// Ok(()) when the timeout is stale.
670    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                // Certificate formed, so do not send timeout vote.
699                return Ok(());
700            }
701            _ => {
702                // No certificate formed, so enter normal timeout processing.
703            }
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                    // Didn't vote in this round yet, generate a backup vote
717                    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                // Vote for proposal
771                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                // Not a validator, just execute the block and wait for votes.
788                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            // No proposal to vote. Send Timeout earlier.
798            self.process_local_timeout(epoch_round).await
799        }
800    }
801
802    /// This function is called only after all the dependencies of the given QC
803    /// have been retrieved.
804    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    /// This function processes a proposal for the current round:
815    /// 1. Filter if it's proposed by valid proposer.
816    /// 2. Execute and add it to a block store.
817    /// 3. Try to vote for it following the safety rules.
818    /// 4. In case a validator chooses to vote, send the vote to the
819    /// representatives at the next round.
820    ///
821    /// Return `Ok(true)` if the block should be relayed.
822    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                // This proposal will not be chosen to vote, so we do not need
865                // to relay. A proposal received for several
866                // times also enters this branch because
867                // the vrf_output is the same.
868                Ok(false)
869            }
870        } else {
871            bail!("unsupported election rules")
872        }
873    }
874
875    /// The function generates a VoteMsg for a given proposed_block:
876    /// * first execute the block and add it to the block store
877    /// * then verify the voting rules
878    /// * save the updated state to consensus DB
879    /// * return a VoteMsg with the LedgerInfo to be committed in case the vote
880    ///   gathers QC.
881    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        // Short circuit if already voted.
890        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    /// Upon new vote:
919    /// 1. Ensures we're processing the vote from the same round as local round
920    /// 2. Filter out votes for rounds that should not be processed by this
921    /// validator (to avoid potential attacks).
922    /// 2. Add the vote to the pending votes and check whether it finishes a QC.
923    /// 3. Once the QC/TC successfully formed, notify the RoundState.
924    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    /// Add a vote to the pending votes.
959    /// If a new QC / TC is formed then
960    /// 1) fetch missing dependencies if required, and then
961    /// 2) call process_certificates(), which will start a new round in return.
962    ///
963    /// Return `Ok(true)` if the vote should be relayed.
964    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        // Add the vote and check whether it completes a new QC or a TC
976        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                // Wait for extra time to gather more votes before entering the
984                // next round.
985                self.round_state
986                    .setup_new_round_timeout(self.epoch_state.epoch);
987            }
988            VoteReceptionResult::VoteAdded(_) => {}
989            VoteReceptionResult::DuplicateVote => {
990                // Do not relay duplicate votes as we should have relayed it
991                // before.
992                relay = false;
993            }
994            VoteReceptionResult::EquivocateVote((vote1, vote2)) => {
995                // Attack detected!
996                // Construct a transaction to dispute this signer.
997                // TODO(lpl): Allow non-committee member to dispute?
998                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                        // TODO(lpl): Track disputed nodes to avoid sending
1036                        // multiple dispute, and retry if needed?
1037                        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            // Return error so that duplicate or invalid votes will not be
1046            // broadcast to others.
1047            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    /// Retrieve a n chained blocks from the block store starting from
1079    /// an initial parent id, returning with <n (as many as possible) if
1080    /// id or its ancestors can not be found.
1081    ///
1082    /// The current version of the function is not really async, but keeping it
1083    /// this way for future possible changes.
1084    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                // TODO(lpl): This error may be needed in the future.
1110                // status = BlockRetrievalStatus::NotEnoughBlocks;
1111                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    /// To jump start new round with the current certificates we have.
1130    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    /// Return true for blocks that we need to process
1162    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    /// Return true for votes that we need to process
1169    pub fn filter_vote(&self, v: &VoteMsg) -> bool {
1170        !self.round_state.vote_received(v.vote())
1171    }
1172}
1173
1174/// The functions used in tests to construct attack cases
1175impl RoundManager {
1176    /// Force the node to vote for a proposal without changing its consensus
1177    /// state. The node will still vote for the correct proposal
1178    /// independently if that's not disabled.
1179    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    /// Force the node to propose a block without changing its consensus
1209    /// state. The node will still propose a valid block independently if that's
1210    /// not disabled.
1211    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        // TODO: This vrf_output is incorrect if we want to propose a block in
1233        // another epoch.
1234        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        // TODO: The sync_info here may not be consistent with
1242        // `signed_proposal`.
1243        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        // It's allowed for a node to sign conflict pivot decision,
1263        // so we do not need to persist this signing event.
1264        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        // TODO(lpl): Check if we want to wait here.
1274        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        // This takes out the candidate, so we need to insert it back if it's
1281        // Some.
1282        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}