@ -28,7 +28,7 @@ func (consensus *Consensus) handleMessageUpdate(payload []byte) {
msg := & msg_pb . Message { }
msg := & msg_pb . Message { }
err := protobuf . Unmarshal ( payload , msg )
err := protobuf . Unmarshal ( payload , msg )
if err != nil {
if err != nil {
utils . Get Logger( ) . Error ( "Failed to unmarshal message payload." , "err" , err , "consensus" , consensus )
utils . Logger ( ) . Error ( ) . Err ( err ) . Interface ( "consensus" , consensus ) . Msg ( "Failed to unmarshal message payload." )
return
return
}
}
@ -40,14 +40,18 @@ func (consensus *Consensus) handleMessageUpdate(payload []byte) {
if msg . Type == msg_pb . MessageType_VIEWCHANGE || msg . Type == msg_pb . MessageType_NEWVIEW {
if msg . Type == msg_pb . MessageType_VIEWCHANGE || msg . Type == msg_pb . MessageType_NEWVIEW {
if msg . GetViewchange ( ) != nil && msg . GetViewchange ( ) . ShardId != consensus . ShardID {
if msg . GetViewchange ( ) != nil && msg . GetViewchange ( ) . ShardId != consensus . ShardID {
consensus . getLogger ( ) . Warn ( "Received view change message from different shard" ,
consensus . getLogger ( ) . Warn ( ) .
"myShardId" , consensus . ShardID , "receivedShardId" , msg . GetViewchange ( ) . ShardId )
Uint32 ( "myShardId" , consensus . ShardID ) .
Uint32 ( "receivedShardId" , msg . GetViewchange ( ) . ShardId ) .
Msg ( "Received view change message from different shard" )
return
return
}
}
} else {
} else {
if msg . GetConsensus ( ) != nil && msg . GetConsensus ( ) . ShardId != consensus . ShardID {
if msg . GetConsensus ( ) != nil && msg . GetConsensus ( ) . ShardId != consensus . ShardID {
consensus . getLogger ( ) . Warn ( "Received consensus message from different shard" ,
consensus . getLogger ( ) . Warn ( ) .
"myShardId" , consensus . ShardID , "receivedShardId" , msg . GetConsensus ( ) . ShardId )
Uint32 ( "myShardId" , consensus . ShardID ) .
Uint32 ( "receivedShardId" , msg . GetConsensus ( ) . ShardId ) .
Msg ( "Received consensus message from different shard" )
return
return
}
}
}
}
@ -68,7 +72,6 @@ func (consensus *Consensus) handleMessageUpdate(payload []byte) {
case msg_pb . MessageType_NEWVIEW :
case msg_pb . MessageType_NEWVIEW :
consensus . onNewView ( msg )
consensus . onNewView ( msg )
}
}
}
}
// TODO: move to consensus_leader.go later
// TODO: move to consensus_leader.go later
@ -79,12 +82,12 @@ func (consensus *Consensus) announce(block *types.Block) {
// prepare message and broadcast to validators
// prepare message and broadcast to validators
encodedBlock , err := rlp . EncodeToBytes ( block )
encodedBlock , err := rlp . EncodeToBytes ( block )
if err != nil {
if err != nil {
consensus . getLogger ( ) . Debug ( "[Announce] Failed encoding block" )
consensus . getLogger ( ) . Debug ( ) . Msg ( "[Announce] Failed encoding block" )
return
return
}
}
encodedBlockHeader , err := rlp . EncodeToBytes ( block . Header ( ) )
encodedBlockHeader , err := rlp . EncodeToBytes ( block . Header ( ) )
if err != nil {
if err != nil {
consensus . getLogger ( ) . Debug ( "[Announce] Failed encoding block header" )
consensus . getLogger ( ) . Debug ( ) . Msg ( "[Announce] Failed encoding block header" )
return
return
}
}
@ -98,55 +101,73 @@ func (consensus *Consensus) announce(block *types.Block) {
_ = protobuf . Unmarshal ( msgPayload , msg )
_ = protobuf . Unmarshal ( msgPayload , msg )
pbftMsg , err := ParsePbftMessage ( msg )
pbftMsg , err := ParsePbftMessage ( msg )
if err != nil {
if err != nil {
consensus . getLogger ( ) . Warn ( "[Announce] Unable to parse pbft message" , "error" , err )
consensus . getLogger ( ) . Warn ( ) . Err ( err ) . Msg ( "[Announce] Unable to parse pbft message" )
return
return
}
}
consensus . PbftLog . AddMessage ( pbftMsg )
consensus . PbftLog . AddMessage ( pbftMsg )
consensus . getLogger ( ) . Debug ( "[Announce] Added Announce message in pbftLog" , "MsgblockHash" , pbftMsg . BlockHash , "MsgViewID" , pbftMsg . ViewID , "MsgBlockNum" , pbftMsg . BlockNum )
consensus . getLogger ( ) . Debug ( ) .
Bytes ( "MsgblockHash" , pbftMsg . BlockHash [ : ] ) .
Uint64 ( "MsgViewID" , pbftMsg . ViewID ) .
Uint64 ( "MsgBlockNum" , pbftMsg . BlockNum ) .
Msg ( "[Announce] Added Announce message in pbftLog" )
consensus . PbftLog . AddBlock ( block )
consensus . PbftLog . AddBlock ( block )
// Leader sign the block hash itself
// Leader sign the block hash itself
consensus . prepareSigs [ consensus . PubKey . SerializeToHexStr ( ) ] = consensus . priKey . SignHash ( consensus . blockHash [ : ] )
consensus . prepareSigs [ consensus . PubKey . SerializeToHexStr ( ) ] = consensus . priKey . SignHash ( consensus . blockHash [ : ] )
if err := consensus . prepareBitmap . SetKey ( consensus . PubKey , true ) ; err != nil {
if err := consensus . prepareBitmap . SetKey ( consensus . PubKey , true ) ; err != nil {
consensus . getLogger ( ) . Warn ( "[Announce] Leader prepareBitmap SetKey failed" , "error" , err )
consensus . getLogger ( ) . Warn ( ) . Err ( err ) . Msg ( "[Announce] Leader prepareBitmap SetKey failed" )
return
return
}
}
// Construct broadcast p2p message
// Construct broadcast p2p message
if err := consensus . host . SendMessageToGroups ( [ ] p2p . GroupID { p2p . NewGroupIDByShardID ( p2p . ShardID ( consensus . ShardID ) ) } , host . ConstructP2pMessage ( byte ( 17 ) , msgToSend ) ) ; err != nil {
if err := consensus . host . SendMessageToGroups ( [ ] p2p . GroupID { p2p . NewGroupIDByShardID ( p2p . ShardID ( consensus . ShardID ) ) } , host . ConstructP2pMessage ( byte ( 17 ) , msgToSend ) ) ; err != nil {
consensus . getLogger ( ) . Warn ( "[Announce] Cannot send announce message" , "groupID" , p2p . NewGroupIDByShardID ( p2p . ShardID ( consensus . ShardID ) ) )
consensus . getLogger ( ) . Warn ( ) .
Str ( "groupID" , string ( p2p . NewGroupIDByShardID ( p2p . ShardID ( consensus . ShardID ) ) ) ) .
Msg ( "[Announce] Cannot send announce message" )
} else {
} else {
consensus . getLogger ( ) . Info ( "[Announce] Sent Announce Message!!" , "BlockHash" , block . Hash ( ) , "BlockNum" , block . NumberU64 ( ) )
consensus . getLogger ( ) . Info ( ) .
Bytes ( "BlockHash" , block . Hash ( ) . Bytes ( ) ) .
Uint64 ( "BlockNum" , block . NumberU64 ( ) ) .
Msg ( "[Announce] Sent Announce Message!!" )
}
}
consensus . getLogger ( ) . Debug ( "[Announce] Switching phase" , "From" , consensus . phase , "To" , Prepare )
consensus . getLogger ( ) . Debug ( ) .
Str ( "From" , consensus . phase . String ( ) ) .
Str ( "To" , Prepare . String ( ) ) .
Msg ( "[Announce] Switching phase" )
consensus . switchPhase ( Prepare , true )
consensus . switchPhase ( Prepare , true )
}
}
func ( consensus * Consensus ) onAnnounce ( msg * msg_pb . Message ) {
func ( consensus * Consensus ) onAnnounce ( msg * msg_pb . Message ) {
consensus . getLogger ( ) . Debug ( "[OnAnnounce] Receive announce message" )
consensus . getLogger ( ) . Debug ( ) . Msg ( "[OnAnnounce] Receive announce message" )
if consensus . PubKey . IsEqual ( consensus . LeaderPubKey ) && consensus . mode . Mode ( ) == Normal {
if consensus . PubKey . IsEqual ( consensus . LeaderPubKey ) && consensus . mode . Mode ( ) == Normal {
return
return
}
}
senderKey , err := consensus . verifySenderKey ( msg )
senderKey , err := consensus . verifySenderKey ( msg )
if err != nil {
if err != nil {
consensus . getLogger ( ) . Debug ( "[OnAnnounce] VerifySenderKey failed" , "error" , err )
consensus . getLogger ( ) . Error ( ) . Err ( err ) . Msg ( "[OnAnnounce] VerifySenderKey failed" )
return
return
}
}
if ! senderKey . IsEqual ( consensus . LeaderPubKey ) && consensus . mode . Mode ( ) == Normal && ! consensus . ignoreViewIDCheck {
if ! senderKey . IsEqual ( consensus . LeaderPubKey ) && consensus . mode . Mode ( ) == Normal && ! consensus . ignoreViewIDCheck {
consensus . getLogger ( ) . Warn ( "[OnAnnounce] SenderKey not match leader PubKey" , "senderKey" , senderKey . SerializeToHexStr ( ) , "leaderKey" , consensus . LeaderPubKey . SerializeToHexStr ( ) )
consensus . getLogger ( ) . Warn ( ) .
Str ( "senderKey" , senderKey . SerializeToHexStr ( ) ) .
Str ( "leaderKey" , consensus . LeaderPubKey . SerializeToHexStr ( ) ) .
Msg ( "[OnAnnounce] SenderKey not match leader PubKey" )
return
return
}
}
if err = verifyMessageSig ( senderKey , msg ) ; err != nil {
if err = verifyMessageSig ( senderKey , msg ) ; err != nil {
consensus . getLogger ( ) . Debug ( "[OnAnnounce] Failed to verify leader signature" , "error" , err )
consensus . getLogger ( ) . Error ( ) . Err ( err ) . Msg ( "[OnAnnounce] Failed to verify leader signature" )
return
return
}
}
recvMsg , err := ParsePbftMessage ( msg )
recvMsg , err := ParsePbftMessage ( msg )
if err != nil {
if err != nil {
consensus . getLogger ( ) . Debug ( "[OnAnnounce] Unparseable leader message" , "error" , err , "MsgBlockNum" , recvMsg . BlockNum )
consensus . getLogger ( ) . Error ( ) .
Err ( err ) .
Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) .
Msg ( "[OnAnnounce] Unparseable leader message" )
return
return
}
}
@ -155,17 +176,27 @@ func (consensus *Consensus) onAnnounce(msg *msg_pb.Message) {
var headerObj types . Header
var headerObj types . Header
err = rlp . DecodeBytes ( blockHeader , & headerObj )
err = rlp . DecodeBytes ( blockHeader , & headerObj )
if err != nil {
if err != nil {
consensus . getLogger ( ) . Warn ( "[OnAnnounce] Unparseable block header data" , "error" , err , "MsgBlockNum" , recvMsg . BlockNum )
consensus . getLogger ( ) . Warn ( ) .
Err ( err ) .
Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) .
Msg ( "[OnAnnounce] Unparseable block header data" )
return
return
}
}
if recvMsg . BlockNum < consensus . blockNum || recvMsg . BlockNum != headerObj . Number . Uint64 ( ) {
if recvMsg . BlockNum < consensus . blockNum || recvMsg . BlockNum != headerObj . Number . Uint64 ( ) {
consensus . getLogger ( ) . Debug ( "[OnAnnounce] BlockNum not match" , "MsgBlockNum" , recvMsg . BlockNum , "BlockNum" , headerObj . Number )
consensus . getLogger ( ) . Debug ( ) .
Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) .
Str ( "BlockNum" , headerObj . Number . String ( ) ) .
Msg ( "[OnAnnounce] BlockNum not match" )
return
return
}
}
if consensus . mode . Mode ( ) == Normal {
if consensus . mode . Mode ( ) == Normal {
if err = consensus . VerifyHeader ( consensus . ChainReader , & headerObj , true ) ; err != nil {
if err = consensus . VerifyHeader ( consensus . ChainReader , & headerObj , true ) ; err != nil {
consensus . getLogger ( ) . Warn ( "[OnAnnounce] Block content is not verified successfully" , "error" , err , "inChain" , consensus . ChainReader . CurrentHeader ( ) . Number , "MsgBlockNum" , headerObj . Number )
consensus . getLogger ( ) . Warn ( ) .
Err ( err ) .
Str ( "inChain" , consensus . ChainReader . CurrentHeader ( ) . Number . String ( ) ) .
Str ( "MsgBlockNum" , headerObj . Number . String ( ) ) .
Msg ( "[OnAnnounce] Block content is not verified successfully" )
return
return
}
}
}
}
@ -173,14 +204,21 @@ func (consensus *Consensus) onAnnounce(msg *msg_pb.Message) {
logMsgs := consensus . PbftLog . GetMessagesByTypeSeqView ( msg_pb . MessageType_ANNOUNCE , recvMsg . BlockNum , recvMsg . ViewID )
logMsgs := consensus . PbftLog . GetMessagesByTypeSeqView ( msg_pb . MessageType_ANNOUNCE , recvMsg . BlockNum , recvMsg . ViewID )
if len ( logMsgs ) > 0 {
if len ( logMsgs ) > 0 {
if logMsgs [ 0 ] . BlockHash != recvMsg . BlockHash {
if logMsgs [ 0 ] . BlockHash != recvMsg . BlockHash {
consensus . getLogger ( ) . Debug ( "[OnAnnounce] Leader is malicious" , "leaderKey" , consensus . LeaderPubKey . SerializeToHexStr ( ) )
consensus . getLogger ( ) . Debug ( ) .
Str ( "leaderKey" , consensus . LeaderPubKey . SerializeToHexStr ( ) ) .
Msg ( "[OnAnnounce] Leader is malicious" )
consensus . startViewChange ( consensus . viewID + 1 )
consensus . startViewChange ( consensus . viewID + 1 )
}
}
consensus . getLogger ( ) . Debug ( "[OnAnnounce] Announce message received again" , "leaderKey" , consensus . LeaderPubKey . SerializeToHexStr ( ) )
consensus . getLogger ( ) . Debug ( ) .
Str ( "leaderKey" , consensus . LeaderPubKey . SerializeToHexStr ( ) ) .
Msg ( "[OnAnnounce] Announce message received again" )
//return
//return
}
}
consensus . getLogger ( ) . Debug ( "[OnAnnounce] Announce message Added" , "MsgViewID" , recvMsg . ViewID , "MsgBlockNum" , recvMsg . BlockNum )
consensus . getLogger ( ) . Debug ( ) .
Uint64 ( "MsgViewID" , recvMsg . ViewID ) .
Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) .
Msg ( "[OnAnnounce] Announce message Added" )
consensus . PbftLog . AddMessage ( recvMsg )
consensus . PbftLog . AddMessage ( recvMsg )
consensus . mutex . Lock ( )
consensus . mutex . Lock ( )
@ -190,13 +228,16 @@ func (consensus *Consensus) onAnnounce(msg *msg_pb.Message) {
// we have already added message and block, skip check viewID and send prepare message if is in ViewChanging mode
// we have already added message and block, skip check viewID and send prepare message if is in ViewChanging mode
if consensus . mode . Mode ( ) == ViewChanging {
if consensus . mode . Mode ( ) == ViewChanging {
consensus . getLogger ( ) . Debug ( "[OnAnnounce] Still in ViewChanging Mode, Exiting !!" )
consensus . getLogger ( ) . Debug ( ) . Msg ( "[OnAnnounce] Still in ViewChanging Mode, Exiting !!" )
return
return
}
}
if consensus . checkViewID ( recvMsg ) != nil {
if consensus . checkViewID ( recvMsg ) != nil {
if consensus . mode . Mode ( ) == Normal {
if consensus . mode . Mode ( ) == Normal {
consensus . getLogger ( ) . Debug ( "[OnAnnounce] ViewID check failed" , "MsgViewID" , recvMsg . ViewID , "msgBlockNum" , recvMsg . BlockNum )
consensus . getLogger ( ) . Debug ( ) .
Uint64 ( "MsgViewID" , recvMsg . ViewID ) .
Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) .
Msg ( "[OnAnnounce] ViewID check failed" )
}
}
return
return
}
}
@ -211,11 +252,16 @@ func (consensus *Consensus) prepare() {
msgToSend := consensus . constructPrepareMessage ( )
msgToSend := consensus . constructPrepareMessage ( )
// TODO: this will not return immediatey, may block
// TODO: this will not return immediatey, may block
if err := consensus . host . SendMessageToGroups ( [ ] p2p . GroupID { p2p . NewGroupIDByShardID ( p2p . ShardID ( consensus . ShardID ) ) } , host . ConstructP2pMessage ( byte ( 17 ) , msgToSend ) ) ; err != nil {
if err := consensus . host . SendMessageToGroups ( [ ] p2p . GroupID { p2p . NewGroupIDByShardID ( p2p . ShardID ( consensus . ShardID ) ) } , host . ConstructP2pMessage ( byte ( 17 ) , msgToSend ) ) ; err != nil {
consensus . getLogger ( ) . Warn ( "[OnAnnounce] Cannot send prepare message" )
consensus . getLogger ( ) . Warn ( ) . Err ( err ) . Msg ( "[OnAnnounce] Cannot send prepare message" )
} else {
} else {
consensus . getLogger ( ) . Info ( "[OnAnnounce] Sent Prepare Message!!" , "BlockHash" , hex . EncodeToString ( consensus . blockHash [ : ] ) )
consensus . getLogger ( ) . Info ( ) .
}
Str ( "BlockHash" , hex . EncodeToString ( consensus . blockHash [ : ] ) ) .
consensus . getLogger ( ) . Debug ( "[Announce] Switching Phase" , "From" , consensus . phase , "To" , Prepare )
Msg ( "[OnAnnounce] Sent Prepare Message!!" )
}
consensus . getLogger ( ) . Debug ( ) .
Str ( "From" , consensus . phase . String ( ) ) .
Str ( "To" , Prepare . String ( ) ) .
Msg ( "[Announce] Switching Phase" )
consensus . switchPhase ( Prepare , true )
consensus . switchPhase ( Prepare , true )
}
}
@ -227,28 +273,33 @@ func (consensus *Consensus) onPrepare(msg *msg_pb.Message) {
senderKey , err := consensus . verifySenderKey ( msg )
senderKey , err := consensus . verifySenderKey ( msg )
if err != nil {
if err != nil {
consensus . getLogger ( ) . Debug ( "[OnPrepare] VerifySenderKey failed" , "error" , err )
consensus . getLogger ( ) . Error ( ) . Err ( err ) . Msg ( "[OnPrepare] VerifySenderKey failed" )
return
return
}
}
if err = verifyMessageSig ( senderKey , msg ) ; err != nil {
if err = verifyMessageSig ( senderKey , msg ) ; err != nil {
consensus . getLogger ( ) . Debug ( "[OnPrepare] Failed to verify sender's signature" , "error" , err )
consensus . getLogger ( ) . Error ( ) . Err ( err ) . Msg ( "[OnPrepare] Failed to verify sender's signature" )
return
return
}
}
recvMsg , err := ParsePbftMessage ( msg )
recvMsg , err := ParsePbftMessage ( msg )
if err != nil {
if err != nil {
consensus . getLogger ( ) . Debug ( "[OnPrepare] Unparseable validator message" , "error" , err )
consensus . getLogger ( ) . Error ( ) . Err ( err ) . Msg ( "[OnPrepare] Unparseable validator message" )
return
return
}
}
if recvMsg . ViewID != consensus . viewID || recvMsg . BlockNum != consensus . blockNum {
if recvMsg . ViewID != consensus . viewID || recvMsg . BlockNum != consensus . blockNum {
consensus . getLogger ( ) . Debug ( "[OnPrepare] Message ViewId or BlockNum not match" ,
consensus . getLogger ( ) . Debug ( ) .
"MsgViewID" , recvMsg . ViewID , "MsgBlockNum" , recvMsg . BlockNum )
Uint64 ( "MsgViewID" , recvMsg . ViewID ) .
Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) .
Msg ( "[OnPrepare] Message ViewId or BlockNum not match" )
return
return
}
}
if ! consensus . PbftLog . HasMatchingViewAnnounce ( consensus . blockNum , consensus . viewID , recvMsg . BlockHash ) {
if ! consensus . PbftLog . HasMatchingViewAnnounce ( consensus . blockNum , consensus . viewID , recvMsg . BlockHash ) {
consensus . getLogger ( ) . Debug ( "[OnPrepare] No Matching Announce message" , "MsgblockHash" , recvMsg . BlockHash , "MsgBlockNum" , recvMsg . BlockNum )
consensus . getLogger ( ) . Debug ( ) .
Uint64 ( "MsgViewID" , recvMsg . ViewID ) .
Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) .
Msg ( "[OnPrepare] No Matching Announce message" )
//return
//return
}
}
@ -260,15 +311,16 @@ func (consensus *Consensus) onPrepare(msg *msg_pb.Message) {
consensus . mutex . Lock ( )
consensus . mutex . Lock ( )
defer consensus . mutex . Unlock ( )
defer consensus . mutex . Unlock ( )
logger := consensus . getLogger ( ) . With ( ) . Str ( "validatorPubKey" , validatorPubKey ) . Logger ( )
if len ( prepareSigs ) >= consensus . Quorum ( ) {
if len ( prepareSigs ) >= consensus . Quorum ( ) {
// already have enough signatures
// already have enough signatures
consensus . getLogger ( ) . Debu g( "[OnPrepare] Received Additional Prepare Message" , "ValidatorPubKey" , validatorPubKey )
logger . Debug ( ) . Ms g( "[OnPrepare] Received Additional Prepare Message" )
return
return
}
}
// proceed only when the message is not received before
// proceed only when the message is not received before
_ , ok := prepareSigs [ validatorPubKey ]
_ , ok := prepareSigs [ validatorPubKey ]
if ok {
if ok {
consensus . getLogger ( ) . Debu g( "[OnPrepare] Already Received prepare message from the validator" , "ValidatorPubKey" , validatorPubKey )
logger . Debug ( ) . Ms g( "[OnPrepare] Already Received prepare message from the validator" )
return
return
}
}
@ -276,24 +328,25 @@ func (consensus *Consensus) onPrepare(msg *msg_pb.Message) {
var sign bls . Sign
var sign bls . Sign
err = sign . Deserialize ( prepareSig )
err = sign . Deserialize ( prepareSig )
if err != nil {
if err != nil {
consensus . getLogger ( ) . Error ( "[OnPrepare] Failed to deserialize bls signature" , "ValidatorPubKey" , validatorPubKey )
consensus . getLogger ( ) . Error ( ) . Err ( err ) . Msg ( "[OnPrepare] Failed to deserialize bls signature" )
return
return
}
}
if ! sign . VerifyHash ( recvMsg . SenderPubkey , consensus . blockHash [ : ] ) {
if ! sign . VerifyHash ( recvMsg . SenderPubkey , consensus . blockHash [ : ] ) {
consensus . getLogger ( ) . Error ( "[OnPrepare] Received invalid BLS signature" , "ValidatorPubKey" , validatorPubKey )
consensus . getLogger ( ) . Error ( ) . Msg ( "[OnPrepare] Received invalid BLS signature" )
return
return
}
}
consensus . getLogger ( ) . Info ( "[OnPrepare] Received New Prepare Signature" , "NumReceivedSoFar" , len ( prepareSigs ) , "validatorPubKey" , validatorPubKey , "PublicKeys" , len ( consensus . PublicKeys ) )
logger = logger . With ( ) . Int ( "NumReceivedSoFar" , len ( prepareSigs ) ) . Int ( "PublicKeys" , len ( consensus . PublicKeys ) ) . Logger ( )
logger . Info ( ) . Msg ( "[OnPrepare] Received New Prepare Signature" )
prepareSigs [ validatorPubKey ] = & sign
prepareSigs [ validatorPubKey ] = & sign
// Set the bitmap indicating that this validator signed.
// Set the bitmap indicating that this validator signed.
if err := prepareBitmap . SetKey ( recvMsg . SenderPubkey , true ) ; err != nil {
if err := prepareBitmap . SetKey ( recvMsg . SenderPubkey , true ) ; err != nil {
consensus . getLogger ( ) . Warn ( "[OnPrepare] prepareBitmap.SetKey failed" , "error" , err )
consensus . getLogger ( ) . Warn ( ) . Err ( err ) . Msg ( "[OnPrepare] prepareBitmap.SetKey failed" )
return
return
}
}
if len ( prepareSigs ) >= consensus . Quorum ( ) {
if len ( prepareSigs ) >= consensus . Quorum ( ) {
consensus . getLogger ( ) . Debu g( "[OnPrepare] Received Enough Prepare Signatures" , "NumReceivedSoFar" , len ( prepareSigs ) , "PublicKeys" , len ( consensus . PublicKeys ) )
logger . Debug ( ) . Ms g( "[OnPrepare] Received Enough Prepare Signatures" )
// Construct and broadcast prepared message
// Construct and broadcast prepared message
msgToSend , aggSig := consensus . constructPreparedMessage ( )
msgToSend , aggSig := consensus . constructPreparedMessage ( )
consensus . aggregatedPrepareSig = aggSig
consensus . aggregatedPrepareSig = aggSig
@ -304,7 +357,7 @@ func (consensus *Consensus) onPrepare(msg *msg_pb.Message) {
_ = protobuf . Unmarshal ( msgPayload , msg )
_ = protobuf . Unmarshal ( msgPayload , msg )
pbftMsg , err := ParsePbftMessage ( msg )
pbftMsg , err := ParsePbftMessage ( msg )
if err != nil {
if err != nil {
consensus . getLogger ( ) . Warn ( "[OnPrepare] Unable to parse pbft message" , "error" , err )
consensus . getLogger ( ) . Warn ( ) . Err ( err ) . Msg ( "[OnPrepare] Unable to parse pbft message" )
return
return
}
}
consensus . PbftLog . AddMessage ( pbftMsg )
consensus . PbftLog . AddMessage ( pbftMsg )
@ -315,51 +368,59 @@ func (consensus *Consensus) onPrepare(msg *msg_pb.Message) {
commitPayload := append ( blockNumHash , consensus . blockHash [ : ] ... )
commitPayload := append ( blockNumHash , consensus . blockHash [ : ] ... )
consensus . commitSigs [ consensus . PubKey . SerializeToHexStr ( ) ] = consensus . priKey . SignHash ( commitPayload )
consensus . commitSigs [ consensus . PubKey . SerializeToHexStr ( ) ] = consensus . priKey . SignHash ( commitPayload )
if err := consensus . commitBitmap . SetKey ( consensus . PubKey , true ) ; err != nil {
if err := consensus . commitBitmap . SetKey ( consensus . PubKey , true ) ; err != nil {
consensus . getLogger ( ) . Debug ( "[OnPrepare] Leader commit bitmap set failed" )
consensus . getLogger ( ) . Debug ( ) . Msg ( "[OnPrepare] Leader commit bitmap set failed" )
return
return
}
}
if err := consensus . host . SendMessageToGroups ( [ ] p2p . GroupID { p2p . NewGroupIDByShardID ( p2p . ShardID ( consensus . ShardID ) ) } , host . ConstructP2pMessage ( byte ( 17 ) , msgToSend ) ) ; err != nil {
if err := consensus . host . SendMessageToGroups ( [ ] p2p . GroupID { p2p . NewGroupIDByShardID ( p2p . ShardID ( consensus . ShardID ) ) } , host . ConstructP2pMessage ( byte ( 17 ) , msgToSend ) ) ; err != nil {
consensus . getLogger ( ) . Warn ( "[OnPrepare] Cannot send prepared message" )
consensus . getLogger ( ) . Warn ( ) . Msg ( "[OnPrepare] Cannot send prepared message" )
} else {
} else {
consensus . getLogger ( ) . Debug ( "[OnPrepare] Sent Prepared Message!!" , "BlockHash" , consensus . blockHash , "BlockNum" , consensus . blockNum )
consensus . getLogger ( ) . Debug ( ) .
Bytes ( "BlockHash" , consensus . blockHash [ : ] ) .
Uint64 ( "BlockNum" , consensus . blockNum ) .
Msg ( "[OnPrepare] Sent Prepared Message!!" )
}
}
consensus . getLogger ( ) . Debug ( "[OnPrepare] Switching phase" , "From" , consensus . phase , "To" , Commit )
consensus . getLogger ( ) . Debug ( ) .
Str ( "From" , consensus . phase . String ( ) ) .
Str ( "To" , Commit . String ( ) ) .
Msg ( "[OnPrepare] Switching phase" )
consensus . switchPhase ( Commit , true )
consensus . switchPhase ( Commit , true )
}
}
return
return
}
}
func ( consensus * Consensus ) onPrepared ( msg * msg_pb . Message ) {
func ( consensus * Consensus ) onPrepared ( msg * msg_pb . Message ) {
consensus . getLogger ( ) . Debug ( "[OnPrepared] Received Prepared message" )
consensus . getLogger ( ) . Debug ( ) . Msg ( "[OnPrepared] Received Prepared message" )
if consensus . PubKey . IsEqual ( consensus . LeaderPubKey ) && consensus . mode . Mode ( ) == Normal {
if consensus . PubKey . IsEqual ( consensus . LeaderPubKey ) && consensus . mode . Mode ( ) == Normal {
return
return
}
}
senderKey , err := consensus . verifySenderKey ( msg )
senderKey , err := consensus . verifySenderKey ( msg )
if err != nil {
if err != nil {
consensus . getLogger ( ) . Debug ( "[OnPrepared] VerifySenderKey failed" , "error" , err )
consensus . getLogger ( ) . Debug ( ) . Err ( err ) . Msg ( "[OnPrepared] VerifySenderKey failed" )
return
return
}
}
if ! senderKey . IsEqual ( consensus . LeaderPubKey ) && consensus . mode . Mode ( ) == Normal && ! consensus . ignoreViewIDCheck {
if ! senderKey . IsEqual ( consensus . LeaderPubKey ) && consensus . mode . Mode ( ) == Normal && ! consensus . ignoreViewIDCheck {
consensus . getLogger ( ) . Warn ( "[OnPrepared] SenderKey not match leader PubKey" )
consensus . getLogger ( ) . Warn ( ) . Msg ( "[OnPrepared] SenderKey not match leader PubKey" )
return
return
}
}
if err := verifyMessageSig ( senderKey , msg ) ; err != nil {
if err := verifyMessageSig ( senderKey , msg ) ; err != nil {
consensus . getLogger ( ) . Debug ( "[OnPrepared] Failed to verify sender's signature" , "error" , err )
consensus . getLogger ( ) . Debug ( ) . Err ( err ) . Msg ( "[OnPrepared] Failed to verify sender's signature" )
return
return
}
}
recvMsg , err := ParsePbftMessage ( msg )
recvMsg , err := ParsePbftMessage ( msg )
if err != nil {
if err != nil {
consensus . getLogger ( ) . Debug ( "[OnPrepared] Unparseable validator message" , "error" , err )
consensus . getLogger ( ) . Debug ( ) . Err ( err ) . Msg ( "[OnPrepared] Unparseable validator message" )
return
return
}
}
consensus . getLogger ( ) . Info ( "[OnPrepared] Received prepared message" , "MsgBlockNum" , recvMsg . BlockNum , "MsgViewID" , recvMsg . ViewID )
consensus . getLogger ( ) . Info ( ) .
Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) .
Uint64 ( "MsgViewID" , recvMsg . ViewID ) .
Msg ( "[OnPrepared] Received prepared message" )
if recvMsg . BlockNum < consensus . blockNum {
if recvMsg . BlockNum < consensus . blockNum {
consensus . getLogger ( ) . Debug ( "Old Block Received, ignoring!!" ,
consensus . getLogger ( ) . Debug ( ) . Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) . Msg ( "Old Block Received, ignoring!!" )
"MsgBlockNum" , recvMsg . BlockNum )
return
return
}
}
@ -367,17 +428,23 @@ func (consensus *Consensus) onPrepared(msg *msg_pb.Message) {
blockHash := recvMsg . BlockHash
blockHash := recvMsg . BlockHash
aggSig , mask , err := consensus . ReadSignatureBitmapPayload ( recvMsg . Payload , 0 )
aggSig , mask , err := consensus . ReadSignatureBitmapPayload ( recvMsg . Payload , 0 )
if err != nil {
if err != nil {
consensus . getLogger ( ) . Error ( "ReadSignatureBitmapPayload failed!!" , "error" , err )
consensus . getLogger ( ) . Error ( ) . Err ( err ) . Msg ( "ReadSignatureBitmapPayload failed!!" )
return
return
}
}
if count := utils . CountOneBits ( mask . Bitmap ) ; count < consensus . Quorum ( ) {
if count := utils . CountOneBits ( mask . Bitmap ) ; count < consensus . Quorum ( ) {
consensus . getLogger ( ) . Debug ( "Not enough signatures in the Prepared msg" , "Need" , consensus . Quorum ( ) , "Got" , count )
consensus . getLogger ( ) . Debug ( ) .
Int ( "Need" , consensus . Quorum ( ) ) .
Int ( "Got" , count ) .
Msg ( "Not enough signatures in the Prepared msg" )
return
return
}
}
if ! aggSig . VerifyHash ( mask . AggregatePublic , blockHash [ : ] ) {
if ! aggSig . VerifyHash ( mask . AggregatePublic , blockHash [ : ] ) {
myBlockHash := common . Hash { }
myBlockHash := common . Hash { }
myBlockHash . SetBytes ( consensus . blockHash [ : ] )
myBlockHash . SetBytes ( consensus . blockHash [ : ] )
consensus . getLogger ( ) . Warn ( "[OnPrepared] failed to verify multi signature for prepare phase" , "MsgBlockHash" , recvMsg . BlockHash , "myBlockHash" , myBlockHash )
consensus . getLogger ( ) . Warn ( ) .
Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) .
Uint64 ( "MsgViewID" , recvMsg . ViewID ) .
Msg ( "[OnPrepared] failed to verify multi signature for prepare phase" )
return
return
}
}
@ -386,26 +453,40 @@ func (consensus *Consensus) onPrepared(msg *msg_pb.Message) {
var blockObj types . Block
var blockObj types . Block
err = rlp . DecodeBytes ( block , & blockObj )
err = rlp . DecodeBytes ( block , & blockObj )
if err != nil {
if err != nil {
consensus . getLogger ( ) . Warn ( "[OnPrepared] Unparseable block header data" , "error" , err , "MsgBlockNum" , recvMsg . BlockNum )
consensus . getLogger ( ) . Warn ( ) .
Err ( err ) .
Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) .
Msg ( "[OnPrepared] Unparseable block header data" )
return
return
}
}
if blockObj . NumberU64 ( ) != recvMsg . BlockNum || recvMsg . BlockNum < consensus . blockNum {
if blockObj . NumberU64 ( ) != recvMsg . BlockNum || recvMsg . BlockNum < consensus . blockNum {
consensus . getLogger ( ) . Warn ( "[OnPrepared] BlockNum not match" , "MsgBlockNum" , recvMsg . BlockNum , "blockNum" , blockObj . NumberU64 ( ) )
consensus . getLogger ( ) . Warn ( ) .
Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) .
Uint64 ( "blockNum" , blockObj . NumberU64 ( ) ) .
Msg ( "[OnPrepared] BlockNum not match" )
return
return
}
}
if blockObj . Header ( ) . Hash ( ) != recvMsg . BlockHash {
if blockObj . Header ( ) . Hash ( ) != recvMsg . BlockHash {
consensus . getLogger ( ) . Warn ( "[OnPrepared] BlockHash not match" , "MsgBlockNum" , recvMsg . BlockNum , "MsgBlockHash" , recvMsg . BlockHash , "blockObjHash" , blockObj . Header ( ) . Hash ( ) )
consensus . getLogger ( ) . Warn ( ) .
Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) .
Bytes ( "MsgBlockHash" , recvMsg . BlockHash [ : ] ) .
Bytes ( "blockObjHash" , blockObj . Header ( ) . Hash ( ) . Bytes ( ) ) .
Msg ( "[OnPrepared] BlockHash not match" )
return
return
}
}
if consensus . mode . Mode ( ) == Normal {
if consensus . mode . Mode ( ) == Normal {
if err := consensus . VerifyHeader ( consensus . ChainReader , blockObj . Header ( ) , true ) ; err != nil {
if err := consensus . VerifyHeader ( consensus . ChainReader , blockObj . Header ( ) , true ) ; err != nil {
consensus . getLogger ( ) . Warn ( "[OnPrepared] Block header is not verified successfully" , "error" , err , "inChain" , consensus . ChainReader . CurrentHeader ( ) . Number , "MsgBlockNum" , blockObj . Header ( ) . Number )
consensus . getLogger ( ) . Warn ( ) .
Err ( err ) .
Str ( "inChain" , consensus . ChainReader . CurrentHeader ( ) . Number . String ( ) ) .
Str ( "MsgBlockNum" , blockObj . Header ( ) . Number . String ( ) ) .
Msg ( "[OnPrepared] Block header is not verified successfully" )
return
return
}
}
if consensus . BlockVerifier == nil {
if consensus . BlockVerifier == nil {
// do nothing
// do nothing
} else if err := consensus . BlockVerifier ( & blockObj ) ; err != nil {
} else if err := consensus . BlockVerifier ( & blockObj ) ; err != nil {
consensus . getLogger ( ) . Info ( "[OnPrepared] Block verification failed" , "error" , err )
consensus . getLogger ( ) . Error ( ) . Err ( err ) . Msg ( "[OnPrepared] Block verification failed" )
return
return
}
}
}
}
@ -413,26 +494,34 @@ func (consensus *Consensus) onPrepared(msg *msg_pb.Message) {
consensus . PbftLog . AddBlock ( & blockObj )
consensus . PbftLog . AddBlock ( & blockObj )
recvMsg . Block = [ ] byte { } // save memory space
recvMsg . Block = [ ] byte { } // save memory space
consensus . PbftLog . AddMessage ( recvMsg )
consensus . PbftLog . AddMessage ( recvMsg )
consensus . getLogger ( ) . Debug ( "[OnPrepared] Prepared message and block added" , "MsgViewID" , recvMsg . ViewID , "MsgBlockNum" , recvMsg . BlockNum , "blockHash" , recvMsg . BlockHash )
consensus . getLogger ( ) . Debug ( ) .
Uint64 ( "MsgViewID" , recvMsg . ViewID ) .
Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) .
Bytes ( "blockHash" , recvMsg . BlockHash [ : ] ) .
Msg ( "[OnPrepared] Prepared message and block added" )
consensus . mutex . Lock ( )
consensus . mutex . Lock ( )
defer consensus . mutex . Unlock ( )
defer consensus . mutex . Unlock ( )
consensus . tryCatchup ( )
consensus . tryCatchup ( )
if consensus . mode . Mode ( ) == ViewChanging {
if consensus . mode . Mode ( ) == ViewChanging {
consensus . getLogger ( ) . Debug ( "[OnPrepared] Still in ViewChanging mode, Exiting !!" )
consensus . getLogger ( ) . Debug ( ) . Msg ( "[OnPrepared] Still in ViewChanging mode, Exiting!!" )
return
return
}
}
if consensus . checkViewID ( recvMsg ) != nil {
if consensus . checkViewID ( recvMsg ) != nil {
if consensus . mode . Mode ( ) == Normal {
if consensus . mode . Mode ( ) == Normal {
consensus . getLogger ( ) . Debug ( "[OnPrepared] ViewID check failed" , "MsgViewID" , recvMsg . ViewID , "MsgBlockNum" , recvMsg . BlockNum )
consensus . getLogger ( ) . Debug ( ) .
Uint64 ( "MsgViewID" , recvMsg . ViewID ) .
Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) .
Msg ( "[OnPrepared] ViewID check failed" )
}
}
return
return
}
}
if recvMsg . BlockNum > consensus . blockNum {
if recvMsg . BlockNum > consensus . blockNum {
consensus . getLogger ( ) . Debug ( "[OnPrepared] Future Block Received, ignoring!!" ,
consensus . getLogger ( ) . Debug ( ) .
"MsgBlockNum" , recvMsg . BlockNum )
Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) .
Msg ( "[OnPrepared] Future Block Received, ignoring!!" )
return
return
}
}
@ -463,12 +552,18 @@ func (consensus *Consensus) onPrepared(msg *msg_pb.Message) {
}
}
if err := consensus . host . SendMessageToGroups ( [ ] p2p . GroupID { p2p . NewGroupIDByShardID ( p2p . ShardID ( consensus . ShardID ) ) } , host . ConstructP2pMessage ( byte ( 17 ) , msgToSend ) ) ; err != nil {
if err := consensus . host . SendMessageToGroups ( [ ] p2p . GroupID { p2p . NewGroupIDByShardID ( p2p . ShardID ( consensus . ShardID ) ) } , host . ConstructP2pMessage ( byte ( 17 ) , msgToSend ) ) ; err != nil {
consensus . getLogger ( ) . Warn ( "[OnPrepared] Cannot send commit message!!" )
consensus . getLogger ( ) . Warn ( ) . Msg ( "[OnPrepared] Cannot send commit message!!" )
} else {
} else {
consensus . getLogger ( ) . Info ( "[OnPrepared] Sent Commit Message!!" , "BlockHash" , consensus . blockHash , "BlockNum" , consensus . blockNum )
consensus . getLogger ( ) . Info ( ) .
Uint64 ( "BlockNum" , consensus . blockNum ) .
Bytes ( "BlockHash" , consensus . blockHash [ : ] ) .
Msg ( "[OnPrepared] Sent Commit Message!!" )
}
}
consensus . getLogger ( ) . Debug ( "[OnPrepared] Switching phase" , "From" , consensus . phase , "To" , Commit )
consensus . getLogger ( ) . Debug ( ) .
Str ( "From" , consensus . phase . String ( ) ) .
Str ( "To" , Commit . String ( ) ) .
Msg ( "[OnPrepared] Switching phase" )
consensus . switchPhase ( Commit , true )
consensus . switchPhase ( Commit , true )
return
return
@ -482,32 +577,41 @@ func (consensus *Consensus) onCommit(msg *msg_pb.Message) {
senderKey , err := consensus . verifySenderKey ( msg )
senderKey , err := consensus . verifySenderKey ( msg )
if err != nil {
if err != nil {
consensus . getLogger ( ) . Debug ( "[OnCommit] VerifySenderKey Failed" , "error" , err )
consensus . getLogger ( ) . Debug ( ) . Err ( err ) . Msg ( "[OnCommit] VerifySenderKey Failed" )
return
return
}
}
if err = verifyMessageSig ( senderKey , msg ) ; err != nil {
if err = verifyMessageSig ( senderKey , msg ) ; err != nil {
consensus . getLogger ( ) . Debug ( "[OnCommit] Failed to verify sender's signature" , "error" , err )
consensus . getLogger ( ) . Debug ( ) . Err ( err ) . Msg ( "[OnCommit] Failed to verify sender's signature" )
return
return
}
}
recvMsg , err := ParsePbftMessage ( msg )
recvMsg , err := ParsePbftMessage ( msg )
if err != nil {
if err != nil {
consensus . getLogger ( ) . Debug ( "[OnCommit] Parse pbft message failed" , "error" , err )
consensus . getLogger ( ) . Debug ( ) . Err ( err ) . Msg ( "[OnCommit] Parse pbft message failed" )
return
return
}
}
if recvMsg . ViewID != consensus . viewID || recvMsg . BlockNum != consensus . blockNum {
if recvMsg . ViewID != consensus . viewID || recvMsg . BlockNum != consensus . blockNum {
consensus . getLogger ( ) . Debug ( "[OnCommit] BlockNum/viewID not match" , "MsgViewID" , recvMsg . ViewID , "MsgBlockNum" , recvMsg . BlockNum , "ValidatorPubKey" , recvMsg . SenderPubkey . SerializeToHexStr ( ) )
consensus . getLogger ( ) . Debug ( ) .
Uint64 ( "MsgViewID" , recvMsg . ViewID ) .
Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) .
Str ( "ValidatorPubKey" , recvMsg . SenderPubkey . SerializeToHexStr ( ) ) .
Msg ( "[OnCommit] BlockNum/viewID not match" )
return
return
}
}
if ! consensus . PbftLog . HasMatchingAnnounce ( consensus . blockNum , recvMsg . BlockHash ) {
if ! consensus . PbftLog . HasMatchingAnnounce ( consensus . blockNum , recvMsg . BlockHash ) {
consensus . getLogger ( ) . Debug ( "[OnCommit] Cannot find matching blockhash" , "MsgBlockHash" , recvMsg . BlockHash , "MsgBlockNum" , recvMsg . BlockNum )
consensus . getLogger ( ) . Debug ( ) .
Bytes ( "MsgBlockHash" , recvMsg . BlockHash [ : ] ) .
Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) .
Msg ( "[OnCommit] Cannot find matching blockhash" )
return
return
}
}
if ! consensus . PbftLog . HasMatchingPrepared ( consensus . blockNum , recvMsg . BlockHash ) {
if ! consensus . PbftLog . HasMatchingPrepared ( consensus . blockNum , recvMsg . BlockHash ) {
consensus . getLogger ( ) . Debug ( "[OnCommit] Cannot find matching prepared message" , "blockHash" , recvMsg . BlockHash )
consensus . getLogger ( ) . Debug ( ) .
Bytes ( "blockHash" , recvMsg . BlockHash [ : ] ) .
Msg ( "[OnCommit] Cannot find matching prepared message" )
return
return
}
}
@ -518,8 +622,9 @@ func (consensus *Consensus) onCommit(msg *msg_pb.Message) {
consensus . mutex . Lock ( )
consensus . mutex . Lock ( )
defer consensus . mutex . Unlock ( )
defer consensus . mutex . Unlock ( )
logger := consensus . getLogger ( ) . With ( ) . Str ( "validatorPubKey" , validatorPubKey ) . Logger ( )
if ! consensus . IsValidatorInCommittee ( recvMsg . SenderPubkey ) {
if ! consensus . IsValidatorInCommittee ( recvMsg . SenderPubkey ) {
consensus . getLogger ( ) . Error ( "[OnCommit] Invalid validator" , "validatorPubKey" , validatorPubKey )
logger . Error ( ) . Msg ( "[OnCommit] Invalid validator" )
return
return
}
}
@ -529,7 +634,7 @@ func (consensus *Consensus) onCommit(msg *msg_pb.Message) {
// proceed only when the message is not received before
// proceed only when the message is not received before
_ , ok := commitSigs [ validatorPubKey ]
_ , ok := commitSigs [ validatorPubKey ]
if ok {
if ok {
consensus . getLogger ( ) . Debu g( "[OnCommit] Already received commit message from the validator" , "validatorPubKey" , validatorPubKey )
logger . Debug ( ) . Ms g( "[OnCommit] Already received commit message from the validator" )
return
return
}
}
@ -539,22 +644,24 @@ func (consensus *Consensus) onCommit(msg *msg_pb.Message) {
var sign bls . Sign
var sign bls . Sign
err = sign . Deserialize ( commitSig )
err = sign . Deserialize ( commitSig )
if err != nil {
if err != nil {
consensus . getLogger ( ) . Debu g( "[OnCommit] Failed to deserialize bls signature" , "validatorPubKey" , validatorPubKey )
logger . Debug ( ) . Ms g( "[OnCommit] Failed to deserialize bls signature" )
return
return
}
}
blockNumHash := make ( [ ] byte , 8 )
blockNumHash := make ( [ ] byte , 8 )
binary . LittleEndian . PutUint64 ( blockNumHash , recvMsg . BlockNum )
binary . LittleEndian . PutUint64 ( blockNumHash , recvMsg . BlockNum )
commitPayload := append ( blockNumHash , recvMsg . BlockHash [ : ] ... )
commitPayload := append ( blockNumHash , recvMsg . BlockHash [ : ] ... )
logger = logger . With ( ) . Uint64 ( "MsgViewID" , recvMsg . ViewID ) . Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) . Logger ( )
if ! sign . VerifyHash ( recvMsg . SenderPubkey , commitPayload ) {
if ! sign . VerifyHash ( recvMsg . SenderPubkey , commitPayload ) {
consensus . getLogger ( ) . Error ( "[OnCommit] Cannot verify commit message" , "MsgViewID" , recvMsg . ViewID , "MsgBlockNum" , recvMsg . BlockNum )
logger . Error ( ) . Msg ( "[OnCommit] Cannot verify commit message" )
return
return
}
}
consensus . getLogger ( ) . Info ( "[OnCommit] Received new commit message" , "numReceivedSoFar" , len ( commitSigs ) , "MsgViewID" , recvMsg . ViewID , "MsgBlockNum" , recvMsg . BlockNum , "validatorPubKey" , validatorPubKey )
logger = logger . With ( ) . Int ( "numReceivedSoFar" , len ( commitSigs ) ) . Logger ( )
logger . Info ( ) . Msg ( "[OnCommit] Received new commit message" )
commitSigs [ validatorPubKey ] = & sign
commitSigs [ validatorPubKey ] = & sign
// Set the bitmap indicating that this validator signed.
// Set the bitmap indicating that this validator signed.
if err := commitBitmap . SetKey ( recvMsg . SenderPubkey , true ) ; err != nil {
if err := commitBitmap . SetKey ( recvMsg . SenderPubkey , true ) ; err != nil {
consensus . getLogger ( ) . Warn ( "[OnCommit] commitBitmap.SetKey failed" , "error" , err )
consensus . getLogger ( ) . Warn ( ) . Err ( err ) . Msg ( "[OnCommit] commitBitmap.SetKey failed" )
return
return
}
}
@ -562,10 +669,10 @@ func (consensus *Consensus) onCommit(msg *msg_pb.Message) {
rewardThresholdIsMet := len ( commitSigs ) >= consensus . RewardThreshold ( )
rewardThresholdIsMet := len ( commitSigs ) >= consensus . RewardThreshold ( )
if ! quorumWasMet && quorumIsMet {
if ! quorumWasMet && quorumIsMet {
consensus . getLogger ( ) . Info ( "[OnCommit] 2/3 Enough commits received" , "NumCommits" , len ( commitSigs ) )
logger . Info ( ) . Msg ( "[OnCommit] 2/3 Enough commits received" )
go func ( viewID uint64 ) {
go func ( viewID uint64 ) {
time . Sleep ( 2 * time . Second )
time . Sleep ( 2 * time . Second )
consensus . getLogger ( ) . Debu g( "[OnCommit] Commit Grace Period Ended" , "NumCommits" , len ( commitSigs ) )
logger . Debug ( ) . Ms g( "[OnCommit] Commit Grace Period Ended" )
consensus . commitFinishChan <- viewID
consensus . commitFinishChan <- viewID
} ( consensus . viewID )
} ( consensus . viewID )
}
}
@ -573,13 +680,13 @@ func (consensus *Consensus) onCommit(msg *msg_pb.Message) {
if rewardThresholdIsMet {
if rewardThresholdIsMet {
go func ( viewID uint64 ) {
go func ( viewID uint64 ) {
consensus . commitFinishChan <- viewID
consensus . commitFinishChan <- viewID
consensus . getLogger ( ) . Info ( "[OnCommit] 90% Enough commits received" , "NumCommits" , len ( commitSigs ) )
logger . Info ( ) . Msg ( "[OnCommit] 90% Enough commits received" )
} ( consensus . viewID )
} ( consensus . viewID )
}
}
}
}
func ( consensus * Consensus ) finalizeCommits ( ) {
func ( consensus * Consensus ) finalizeCommits ( ) {
consensus . getLogger ( ) . Info ( "[Finalizing] Finalizing Block" , "NumCommits" , len ( consensus . commitSigs ) )
consensus . getLogger ( ) . Info ( ) . Int ( "NumCommits" , len ( consensus . commitSigs ) ) . Msg ( "[Finalizing] Finalizing Block" )
beforeCatchupNum := consensus . blockNum
beforeCatchupNum := consensus . blockNum
beforeCatchupViewID := consensus . viewID
beforeCatchupViewID := consensus . viewID
@ -594,7 +701,7 @@ func (consensus *Consensus) finalizeCommits() {
_ = protobuf . Unmarshal ( msgPayload , msg )
_ = protobuf . Unmarshal ( msgPayload , msg )
pbftMsg , err := ParsePbftMessage ( msg )
pbftMsg , err := ParsePbftMessage ( msg )
if err != nil {
if err != nil {
consensus . getLogger ( ) . Warn ( "[FinalizeCommits] Unable to parse pbft message" , "error" , err )
consensus . getLogger ( ) . Warn ( ) . Err ( err ) . Msg ( "[FinalizeCommits] Unable to parse pbft message" )
return
return
}
}
consensus . PbftLog . AddMessage ( pbftMsg )
consensus . PbftLog . AddMessage ( pbftMsg )
@ -602,19 +709,26 @@ func (consensus *Consensus) finalizeCommits() {
// find correct block content
// find correct block content
block := consensus . PbftLog . GetBlockByHash ( consensus . blockHash )
block := consensus . PbftLog . GetBlockByHash ( consensus . blockHash )
if block == nil {
if block == nil {
consensus . getLogger ( ) . Warn ( "[FinalizeCommits] Cannot find block by hash" , "blockHash" , hex . EncodeToString ( consensus . blockHash [ : ] ) )
consensus . getLogger ( ) . Warn ( ) .
Str ( "blockHash" , hex . EncodeToString ( consensus . blockHash [ : ] ) ) .
Msg ( "[FinalizeCommits] Cannot find block by hash" )
return
return
}
}
consensus . tryCatchup ( )
consensus . tryCatchup ( )
if consensus . blockNum - beforeCatchupNum != 1 {
if consensus . blockNum - beforeCatchupNum != 1 {
consensus . getLogger ( ) . Warn ( "[FinalizeCommits] Leader cannot provide the correct block for committed message" , "beforeCatchupBlockNum" , beforeCatchupNum )
consensus . getLogger ( ) . Warn ( ) .
Uint64 ( "beforeCatchupBlockNum" , beforeCatchupNum ) .
Msg ( "[FinalizeCommits] Leader cannot provide the correct block for committed message" )
return
return
}
}
// if leader success finalize the block, send committed message to validators
// if leader success finalize the block, send committed message to validators
if err := consensus . host . SendMessageToGroups ( [ ] p2p . GroupID { p2p . NewGroupIDByShardID ( p2p . ShardID ( consensus . ShardID ) ) } , host . ConstructP2pMessage ( byte ( 17 ) , msgToSend ) ) ; err != nil {
if err := consensus . host . SendMessageToGroups ( [ ] p2p . GroupID { p2p . NewGroupIDByShardID ( p2p . ShardID ( consensus . ShardID ) ) } , host . ConstructP2pMessage ( byte ( 17 ) , msgToSend ) ) ; err != nil {
consensus . getLogger ( ) . Warn ( "[Finalizing] Cannot send committed message" , "error" , err )
consensus . getLogger ( ) . Warn ( ) . Err ( err ) . Msg ( "[Finalizing] Cannot send committed message" )
} else {
} else {
consensus . getLogger ( ) . Info ( "[Finalizing] Sent Committed Message" , "BlockHash" , consensus . blockHash , "BlockNum" , consensus . blockNum )
consensus . getLogger ( ) . Info ( ) .
Bytes ( "BlockHash" , consensus . blockHash [ : ] ) .
Uint64 ( "BlockNum" , consensus . blockNum ) .
Msg ( "[Finalizing] Sent Committed Message" )
}
}
consensus . reportMetrics ( * block )
consensus . reportMetrics ( * block )
@ -627,20 +741,25 @@ func (consensus *Consensus) finalizeCommits() {
if consensus . consensusTimeout [ timeoutBootstrap ] . IsActive ( ) {
if consensus . consensusTimeout [ timeoutBootstrap ] . IsActive ( ) {
consensus . consensusTimeout [ timeoutBootstrap ] . Stop ( )
consensus . consensusTimeout [ timeoutBootstrap ] . Stop ( )
consensus . getLogger ( ) . Debug ( "[Finalizing] Start consensus timer; stop bootstrap timer only once" )
consensus . getLogger ( ) . Debug ( ) . Msg ( "[Finalizing] Start consensus timer; stop bootstrap timer only once" )
} else {
} else {
consensus . getLogger ( ) . Debug ( "[Finalizing] Start consensus timer" )
consensus . getLogger ( ) . Debug ( ) . Msg ( "[Finalizing] Start consensus timer" )
}
}
consensus . consensusTimeout [ timeoutConsensus ] . Start ( )
consensus . consensusTimeout [ timeoutConsensus ] . Start ( )
consensus . getLogger ( ) . Info ( "HOORAY!!!!!!! CONSENSUS REACHED!!!!!!!" , "BlockNum" , beforeCatchupNum , "ViewId" , beforeCatchupViewID , "BlockHash" , block . Hash ( ) , "index" , consensus . getIndexOfPubKey ( consensus . PubKey ) )
consensus . getLogger ( ) . Info ( ) .
Uint64 ( "BlockNum" , beforeCatchupNum ) .
Uint64 ( "ViewId" , beforeCatchupViewID ) .
Str ( "BlockHash" , block . Hash ( ) . String ( ) ) .
Int ( "index" , consensus . getIndexOfPubKey ( consensus . PubKey ) ) .
Msg ( "HOORAY!!!!!!! CONSENSUS REACHED!!!!!!!" )
// Send signal to Node so the new block can be added and new round of consensus can be triggered
// Send signal to Node so the new block can be added and new round of consensus can be triggered
consensus . ReadySignal <- struct { } { }
consensus . ReadySignal <- struct { } { }
}
}
func ( consensus * Consensus ) onCommitted ( msg * msg_pb . Message ) {
func ( consensus * Consensus ) onCommitted ( msg * msg_pb . Message ) {
consensus . getLogger ( ) . Debug ( "[OnCommitted] Receive committed message" )
consensus . getLogger ( ) . Debug ( ) . Msg ( "[OnCommitted] Receive committed message" )
if consensus . PubKey . IsEqual ( consensus . LeaderPubKey ) && consensus . mode . Mode ( ) == Normal {
if consensus . PubKey . IsEqual ( consensus . LeaderPubKey ) && consensus . mode . Mode ( ) == Normal {
return
return
@ -648,37 +767,42 @@ func (consensus *Consensus) onCommitted(msg *msg_pb.Message) {
senderKey , err := consensus . verifySenderKey ( msg )
senderKey , err := consensus . verifySenderKey ( msg )
if err != nil {
if err != nil {
consensus . getLogger ( ) . Warn ( "[OnCommitted] verifySenderKey failed" , "error" , err )
consensus . getLogger ( ) . Warn ( ) . Err ( err ) . Msg ( "[OnCommitted] verifySenderKey failed" )
return
return
}
}
if ! senderKey . IsEqual ( consensus . LeaderPubKey ) && consensus . mode . Mode ( ) == Normal && ! consensus . ignoreViewIDCheck {
if ! senderKey . IsEqual ( consensus . LeaderPubKey ) && consensus . mode . Mode ( ) == Normal && ! consensus . ignoreViewIDCheck {
consensus . getLogger ( ) . Warn ( "[OnCommitted] senderKey not match leader PubKey" )
consensus . getLogger ( ) . Warn ( ) . Msg ( "[OnCommitted] senderKey not match leader PubKey" )
return
return
}
}
if err = verifyMessageSig ( senderKey , msg ) ; err != nil {
if err = verifyMessageSig ( senderKey , msg ) ; err != nil {
consensus . getLogger ( ) . Warn ( "[OnCommitted] Failed to verify sender's signature" , "error" , err )
consensus . getLogger ( ) . Warn ( ) . Err ( err ) . Msg ( "[OnCommitted] Failed to verify sender's signature" )
return
return
}
}
recvMsg , err := ParsePbftMessage ( msg )
recvMsg , err := ParsePbftMessage ( msg )
if err != nil {
if err != nil {
consensus . getLogger ( ) . Warn ( "[OnCommitted] unable to parse msg" , "error" , err )
consensus . getLogger ( ) . Warn ( ) . Msg ( "[OnCommitted] unable to parse msg" )
return
return
}
}
if recvMsg . BlockNum < consensus . blockNum {
if recvMsg . BlockNum < consensus . blockNum {
consensus . getLogger ( ) . Info ( "[OnCommitted] Received Old Blocks!!" , "MsgBlockNum" , recvMsg . BlockNum )
consensus . getLogger ( ) . Info ( ) .
Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) .
Msg ( "[OnCommitted] Received Old Blocks!!" )
return
return
}
}
aggSig , mask , err := consensus . ReadSignatureBitmapPayload ( recvMsg . Payload , 0 )
aggSig , mask , err := consensus . ReadSignatureBitmapPayload ( recvMsg . Payload , 0 )
if err != nil {
if err != nil {
consensus . getLogger ( ) . Error ( "[OnCommitted] readSignatureBitmapPayload failed" , "error" , err )
consensus . getLogger ( ) . Error ( ) . Err ( err ) . Msg ( "[OnCommitted] readSignatureBitmapPayload failed" )
return
return
}
}
// check has 2f+1 signatures
// check has 2f+1 signatures
if count := utils . CountOneBits ( mask . Bitmap ) ; count < consensus . Quorum ( ) {
if count := utils . CountOneBits ( mask . Bitmap ) ; count < consensus . Quorum ( ) {
consensus . getLogger ( ) . Warn ( "[OnCommitted] Not enough signature in committed msg" , "need" , consensus . Quorum ( ) , "got" , count )
consensus . getLogger ( ) . Warn ( ) .
Int ( "need" , consensus . Quorum ( ) ) .
Int ( "got" , count ) .
Msg ( "[OnCommitted] Not enough signature in committed msg" )
return
return
}
}
@ -686,12 +810,17 @@ func (consensus *Consensus) onCommitted(msg *msg_pb.Message) {
binary . LittleEndian . PutUint64 ( blockNumBytes , recvMsg . BlockNum )
binary . LittleEndian . PutUint64 ( blockNumBytes , recvMsg . BlockNum )
commitPayload := append ( blockNumBytes , recvMsg . BlockHash [ : ] ... )
commitPayload := append ( blockNumBytes , recvMsg . BlockHash [ : ] ... )
if ! aggSig . VerifyHash ( mask . AggregatePublic , commitPayload ) {
if ! aggSig . VerifyHash ( mask . AggregatePublic , commitPayload ) {
consensus . getLogger ( ) . Error ( "[OnCommitted] Failed to verify the multi signature for commit phase" , "MsgBlockNum" , recvMsg . BlockNum )
consensus . getLogger ( ) . Error ( ) .
Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) .
Msg ( "[OnCommitted] Failed to verify the multi signature for commit phase" )
return
return
}
}
consensus . PbftLog . AddMessage ( recvMsg )
consensus . PbftLog . AddMessage ( recvMsg )
consensus . getLogger ( ) . Debug ( "[OnCommitted] Committed message added" , "MsgViewID" , recvMsg . ViewID , "MsgBlockNum" , recvMsg . BlockNum )
consensus . getLogger ( ) . Debug ( ) .
Uint64 ( "MsgViewID" , recvMsg . ViewID ) .
Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) .
Msg ( "[OnCommitted] Committed message added" )
consensus . mutex . Lock ( )
consensus . mutex . Lock ( )
defer consensus . mutex . Unlock ( )
defer consensus . mutex . Unlock ( )
@ -700,7 +829,7 @@ func (consensus *Consensus) onCommitted(msg *msg_pb.Message) {
consensus . commitBitmap = mask
consensus . commitBitmap = mask
if recvMsg . BlockNum - consensus . blockNum > consensusBlockNumBuffer {
if recvMsg . BlockNum - consensus . blockNum > consensusBlockNumBuffer {
consensus . getLogger ( ) . Debug ( "[OnCommitted] out of sync" , "MsgBlockNum" , recvMsg . BlockNum )
consensus . getLogger ( ) . Debug ( ) . Uint64 ( "MsgBlockNum" , recvMsg . BlockNum ) . Msg ( "[OnCommitted] out of sync" )
go func ( ) {
go func ( ) {
select {
select {
case consensus . blockNumLowChan <- struct { } { } :
case consensus . blockNumLowChan <- struct { } { } :
@ -721,15 +850,15 @@ func (consensus *Consensus) onCommitted(msg *msg_pb.Message) {
consensus . tryCatchup ( )
consensus . tryCatchup ( )
if consensus . mode . Mode ( ) == ViewChanging {
if consensus . mode . Mode ( ) == ViewChanging {
consensus . getLogger ( ) . Debug ( "[OnCommitted] Still in ViewChanging mode, Exiting !!" )
consensus . getLogger ( ) . Debug ( ) . Msg ( "[OnCommitted] Still in ViewChanging mode, Exiting!!" )
return
return
}
}
if consensus . consensusTimeout [ timeoutBootstrap ] . IsActive ( ) {
if consensus . consensusTimeout [ timeoutBootstrap ] . IsActive ( ) {
consensus . consensusTimeout [ timeoutBootstrap ] . Stop ( )
consensus . consensusTimeout [ timeoutBootstrap ] . Stop ( )
consensus . getLogger ( ) . Debug ( "[OnCommitted] Start consensus timer; stop bootstrap timer only once" )
consensus . getLogger ( ) . Debug ( ) . Msg ( "[OnCommitted] Start consensus timer; stop bootstrap timer only once" )
} else {
} else {
consensus . getLogger ( ) . Debug ( "[OnCommitted] Start consensus timer" )
consensus . getLogger ( ) . Debug ( ) . Msg ( "[OnCommitted] Start consensus timer" )
}
}
consensus . consensusTimeout [ timeoutConsensus ] . Start ( )
consensus . consensusTimeout [ timeoutConsensus ] . Start ( )
return
return
@ -757,7 +886,7 @@ func (consensus *Consensus) LastCommitSig() ([]byte, []byte, error) {
// try to catch up if fall behind
// try to catch up if fall behind
func ( consensus * Consensus ) tryCatchup ( ) {
func ( consensus * Consensus ) tryCatchup ( ) {
consensus . getLogger ( ) . Info ( "[TryCatchup] commit new blocks" )
consensus . getLogger ( ) . Info ( ) . Msg ( "[TryCatchup] commit new blocks" )
// if consensus.phase != Commit && consensus.mode.Mode() == Normal {
// if consensus.phase != Commit && consensus.mode.Mode() == Normal {
// return
// return
// }
// }
@ -768,9 +897,11 @@ func (consensus *Consensus) tryCatchup() {
break
break
}
}
if len ( msgs ) > 1 {
if len ( msgs ) > 1 {
consensus . getLogger ( ) . Error ( "[TryCatchup] DANGER!!! we should only get one committed message for a given blockNum" , "numMsgs" , len ( msgs ) )
consensus . getLogger ( ) . Error ( ) .
Int ( "numMsgs" , len ( msgs ) ) .
Msg ( "[TryCatchup] DANGER!!! we should only get one committed message for a given blockNum" )
}
}
consensus . getLogger ( ) . Info ( "[TryCatchup] committed message found" )
consensus . getLogger ( ) . Info ( ) . Msg ( "[TryCatchup] committed message found" )
block := consensus . PbftLog . GetBlockByHash ( msgs [ 0 ] . BlockHash )
block := consensus . PbftLog . GetBlockByHash ( msgs [ 0 ] . BlockHash )
if block == nil {
if block == nil {
@ -780,43 +911,48 @@ func (consensus *Consensus) tryCatchup() {
if consensus . BlockVerifier == nil {
if consensus . BlockVerifier == nil {
// do nothing
// do nothing
} else if err := consensus . BlockVerifier ( block ) ; err != nil {
} else if err := consensus . BlockVerifier ( block ) ; err != nil {
consensus . getLogger ( ) . Info ( "[TryCatchup]block verification faied" )
consensus . getLogger ( ) . Info ( ) . Msg ( "[TryCatchup]block verification faied" )
return
return
}
}
if block . ParentHash ( ) != consensus . ChainReader . CurrentHeader ( ) . Hash ( ) {
if block . ParentHash ( ) != consensus . ChainReader . CurrentHeader ( ) . Hash ( ) {
consensus . getLogger ( ) . Debug ( "[TryCatchup] parent block hash not match" )
consensus . getLogger ( ) . Debug ( ) . Msg ( "[TryCatchup] parent block hash not match" )
break
break
}
}
consensus . getLogger ( ) . Info ( "[TryCatchup] block found to commit" )
consensus . getLogger ( ) . Info ( ) . Msg ( "[TryCatchup] block found to commit" )
preparedMsgs := consensus . PbftLog . GetMessagesByTypeSeqHash ( msg_pb . MessageType_PREPARED , msgs [ 0 ] . BlockNum , msgs [ 0 ] . BlockHash )
preparedMsgs := consensus . PbftLog . GetMessagesByTypeSeqHash ( msg_pb . MessageType_PREPARED , msgs [ 0 ] . BlockNum , msgs [ 0 ] . BlockHash )
msg := consensus . PbftLog . FindMessageByMaxViewID ( preparedMsgs )
msg := consensus . PbftLog . FindMessageByMaxViewID ( preparedMsgs )
if msg == nil {
if msg == nil {
break
break
}
}
consensus . getLogger ( ) . Info ( "[TryCatchup] prepared message found to commit" )
consensus . getLogger ( ) . Info ( ) . Msg ( "[TryCatchup] prepared message found to commit" )
consensus . blockHash = [ 32 ] byte { }
consensus . blockHash = [ 32 ] byte { }
consensus . blockNum = consensus . blockNum + 1
consensus . blockNum = consensus . blockNum + 1
consensus . viewID = msgs [ 0 ] . ViewID + 1
consensus . viewID = msgs [ 0 ] . ViewID + 1
consensus . LeaderPubKey = msgs [ 0 ] . SenderPubkey
consensus . LeaderPubKey = msgs [ 0 ] . SenderPubkey
consensus . getLogger ( ) . Info ( "[TryCatchup] Adding block to chain" )
consensus . getLogger ( ) . Info ( ) . Msg ( "[TryCatchup] Adding block to chain" )
consensus . OnConsensusDone ( block )
consensus . OnConsensusDone ( block )
consensus . ResetState ( )
consensus . ResetState ( )
select {
select {
case consensus . VerifiedNewBlock <- block :
case consensus . VerifiedNewBlock <- block :
default :
default :
consensus . getLogger ( ) . Info ( "[TryCatchup] consensus verified block send to chan failed" , "blockHash" , block . Hash ( ) )
consensus . getLogger ( ) . Info ( ) .
Str ( "blockHash" , block . Hash ( ) . String ( ) ) .
Msg ( "[TryCatchup] consensus verified block send to chan failed" )
continue
continue
}
}
break
break
}
}
if currentBlockNum < consensus . blockNum {
if currentBlockNum < consensus . blockNum {
consensus . getLogger ( ) . Info ( "[TryCatchup] Catched up!" , "From" , currentBlockNum , "To" , consensus . blockNum )
consensus . getLogger ( ) . Info ( ) .
Uint64 ( "From" , currentBlockNum ) .
Uint64 ( "To" , consensus . blockNum ) .
Msg ( "[TryCatchup] Catched up!" )
consensus . switchPhase ( Announce , true )
consensus . switchPhase ( Announce , true )
}
}
// catup up and skip from view change trap
// catup up and skip from view change trap
@ -833,7 +969,7 @@ func (consensus *Consensus) tryCatchup() {
func ( consensus * Consensus ) Start ( blockChannel chan * types . Block , stopChan chan struct { } , stoppedChan chan struct { } , startChannel chan struct { } ) {
func ( consensus * Consensus ) Start ( blockChannel chan * types . Block , stopChan chan struct { } , stoppedChan chan struct { } , startChannel chan struct { } ) {
go func ( ) {
go func ( ) {
if nodeconfig . GetDefaultConfig ( ) . IsLeader ( ) {
if nodeconfig . GetDefaultConfig ( ) . IsLeader ( ) {
consensus . getLogger ( ) . Info ( "[ConsensusMainLoop] Waiting for consensus start" , "time" , time . Now ( ) )
consensus . getLogger ( ) . Info ( ) . Time ( "time" , time . Now ( ) ) . Msg ( "[ConsensusMainLoop] Waiting for consensus start" )
<- startChannel
<- startChannel
if nodeconfig . GetDefaultConfig ( ) . IsLeader ( ) {
if nodeconfig . GetDefaultConfig ( ) . IsLeader ( ) {
@ -844,11 +980,14 @@ func (consensus *Consensus) Start(blockChannel chan *types.Block, stopChan chan
} ( )
} ( )
}
}
}
}
consensus . getLogger ( ) . Info ( "[ConsensusMainLoop] Consensus started" , "time" , time . Now ( ) )
consensus . getLogger ( ) . Info ( ) . Time ( "time" , time . Now ( ) ) . Msg ( "[ConsensusMainLoop] Consensus started" )
defer close ( stoppedChan )
defer close ( stoppedChan )
ticker := time . NewTicker ( 3 * time . Second )
ticker := time . NewTicker ( 3 * time . Second )
consensus . consensusTimeout [ timeoutBootstrap ] . Start ( )
consensus . consensusTimeout [ timeoutBootstrap ] . Start ( )
consensus . getLogger ( ) . Debug ( "[ConsensusMainLoop] Start bootstrap timeout (only once)" , "viewID" , consensus . viewID , "block" , consensus . blockNum )
consensus . getLogger ( ) . Debug ( ) .
Uint64 ( "viewID" , consensus . viewID ) .
Uint64 ( "block" , consensus . blockNum ) .
Msg ( "[ConsensusMainLoop] Start bootstrap timeout (only once)" )
for {
for {
select {
select {
case <- ticker . C :
case <- ticker . C :
@ -860,11 +999,11 @@ func (consensus *Consensus) Start(blockChannel chan *types.Block, stopChan chan
continue
continue
}
}
if k != timeoutViewChange {
if k != timeoutViewChange {
consensus . getLogger ( ) . Debug ( "[ConsensusMainLoop] Ops Consensus Timeout!!!" )
consensus . getLogger ( ) . Debug ( ) . Msg ( "[ConsensusMainLoop] Ops Consensus Timeout!!!" )
consensus . startViewChange ( consensus . viewID + 1 )
consensus . startViewChange ( consensus . viewID + 1 )
break
break
} else {
} else {
consensus . getLogger ( ) . Debug ( "[ConsensusMainLoop] Ops View Change Timeout!!!" )
consensus . getLogger ( ) . Debug ( ) . Msg ( "[ConsensusMainLoop] Ops View Change Timeout!!!" )
viewID := consensus . mode . ViewID ( )
viewID := consensus . mode . ViewID ( )
consensus . startViewChange ( viewID + 1 )
consensus . startViewChange ( viewID + 1 )
break
break
@ -872,15 +1011,17 @@ func (consensus *Consensus) Start(blockChannel chan *types.Block, stopChan chan
}
}
case <- consensus . syncReadyChan :
case <- consensus . syncReadyChan :
consensus . updateConsensusInformation ( )
consensus . updateConsensusInformation ( )
consensus . getLogger ( ) . Info ( "Node is in sync" )
consensus . getLogger ( ) . Info ( ) . Msg ( "Node is in sync" )
case <- consensus . syncNotReadyChan :
case <- consensus . syncNotReadyChan :
consensus . SetBlockNum ( consensus . ChainReader . CurrentHeader ( ) . Number . Uint64 ( ) + 1 )
consensus . SetBlockNum ( consensus . ChainReader . CurrentHeader ( ) . Number . Uint64 ( ) + 1 )
consensus . mode . SetMode ( Syncing )
consensus . mode . SetMode ( Syncing )
consensus . getLogger ( ) . Info ( "Node is out of sync" )
consensus . getLogger ( ) . Info ( ) . Msg ( "Node is out of sync" )
case newBlock := <- blockChannel :
case newBlock := <- blockChannel :
consensus . getLogger ( ) . Info ( "[ConsensusMainLoop] Received Proposed New Block!" , "MsgBlockNum" , newBlock . NumberU64 ( ) )
consensus . getLogger ( ) . Info ( ) .
Uint64 ( "MsgBlockNum" , newBlock . NumberU64 ( ) ) .
Msg ( "[ConsensusMainLoop] Received Proposed New Block!" )
if consensus . ShardID == 0 {
if consensus . ShardID == 0 {
// TODO ek/rj - re-enable this after fixing DRand
// TODO ek/rj - re-enable this after fixing DRand
//if core.IsEpochBlock(newBlock) { // Only beacon chain do randomness generation
//if core.IsEpochBlock(newBlock) { // Only beacon chain do randomness generation
@ -902,7 +1043,9 @@ func (consensus *Consensus) Start(blockChannel chan *types.Block, stopChan chan
if err == nil {
if err == nil {
// Verify the randomness
// Verify the randomness
_ = blockHash
_ = blockHash
consensus . getLogger ( ) . Info ( "[ConsensusMainLoop] Adding randomness into new block" , "rnd" , rnd )
consensus . getLogger ( ) . Info ( ) .
Bytes ( "rnd" , rnd [ : ] ) .
Msg ( "[ConsensusMainLoop] Adding randomness into new block" )
// newBlock.AddVdf([258]byte{}) // TODO(HB): add real vdf
// newBlock.AddVdf([258]byte{}) // TODO(HB): add real vdf
} else {
} else {
//consensus.getLogger().Info("Failed to get randomness", "error", err)
//consensus.getLogger().Info("Failed to get randomness", "error", err)
@ -910,7 +1053,12 @@ func (consensus *Consensus) Start(blockChannel chan *types.Block, stopChan chan
}
}
startTime = time . Now ( )
startTime = time . Now ( )
consensus . getLogger ( ) . Debug ( "[ConsensusMainLoop] STARTING CONSENSUS" , "numTxs" , len ( newBlock . Transactions ( ) ) , "consensus" , consensus , "startTime" , startTime , "publicKeys" , len ( consensus . PublicKeys ) )
consensus . getLogger ( ) . Debug ( ) .
Int ( "numTxs" , len ( newBlock . Transactions ( ) ) ) .
Interface ( "consensus" , consensus ) .
Time ( "startTime" , startTime ) .
Int ( "publicKeys" , len ( consensus . PublicKeys ) ) .
Msg ( "[ConsensusMainLoop] STARTING CONSENSUS" )
consensus . announce ( newBlock )
consensus . announce ( newBlock )
case msg := <- consensus . MsgChan :
case msg := <- consensus . MsgChan :