@ -1,25 +1,21 @@
package consensus
package consensus
import (
import (
"encoding/binary"
"fmt"
"fmt"
"sync"
"github.com/harmony-one/harmony/crypto/bls"
mapset "github.com/deckarep/golang-set"
"github.com/ethereum/go-ethereum/common"
"github.com/ethereum/go-ethereum/common"
bls_core "github.com/harmony-one/bls/ffi/go/bls"
bls_core "github.com/harmony-one/bls/ffi/go/bls"
"github.com/pkg/errors"
"github.com/rs/zerolog"
msg_pb "github.com/harmony-one/harmony/api/proto/message"
msg_pb "github.com/harmony-one/harmony/api/proto/message"
"github.com/harmony-one/harmony/core/types"
"github.com/harmony-one/harmony/core/types"
"github.com/harmony-one/harmony/crypto/bls"
bls_cosi "github.com/harmony-one/harmony/crypto/bls"
bls_cosi "github.com/harmony-one/harmony/crypto/bls"
)
)
// FBFTLog represents the log stored by a node during FBFT process
type FBFTLog struct {
blocks mapset . Set //store blocks received in FBFT
messages mapset . Set // store messages received in FBFT
maxLogSize uint32
}
// FBFTMessage is the record of pbft messages received by a node during FBFT process
// FBFTMessage is the record of pbft messages received by a node during FBFT process
type FBFTMessage struct {
type FBFTMessage struct {
MessageType msg_pb . MessageType
MessageType msg_pb . MessageType
@ -59,106 +55,153 @@ func (m *FBFTMessage) String() string {
)
)
}
}
const (
idTypeBytes = 4
idHashBytes = common . HashLength
idSenderBytes = bls . PublicKeySizeInBytes
idBytes = idTypeBytes + idHashBytes + idSenderBytes
)
type (
// fbftMsgID is the id that uniquely defines a fbft message.
fbftMsgID [ idBytes ] byte
)
// id return the ID of the FBFT message which uniquely identifies a FBFT message.
// The ID is a concatenation of MsgType, BlockHash, and sender key
func ( m * FBFTMessage ) id ( ) fbftMsgID {
var id fbftMsgID
binary . LittleEndian . PutUint32 ( id [ : ] , uint32 ( m . MessageType ) )
copy ( id [ idTypeBytes : ] , m . BlockHash [ : ] )
if m . SenderPubkey != nil {
copy ( id [ idTypeBytes + idHashBytes : ] , m . SenderPubkey . Bytes [ : ] )
}
return id
}
// FBFTLog represents the log stored by a node during FBFT process
type FBFTLog struct {
blocks map [ common . Hash ] * types . Block // store blocks received in FBFT
verifiedBlocks map [ common . Hash ] struct { } // store block hashes for blocks that has already been verified
blockLock sync . RWMutex
messages map [ fbftMsgID ] * FBFTMessage // store messages received in FBFT
msgLock sync . RWMutex
}
// NewFBFTLog returns new instance of FBFTLog
// NewFBFTLog returns new instance of FBFTLog
func NewFBFTLog ( ) * FBFTLog {
func NewFBFTLog ( ) * FBFTLog {
blocks := mapset . NewSet ( )
pbftLog := FBFTLog {
messages := mapset . NewSet ( )
blocks : make ( map [ common . Hash ] * types . Block ) ,
logSize := maxLogSize
messages : make ( map [ fbftMsgID ] * FBFTMessage ) ,
pbftLog := FBFTLog { blocks : blocks , messages : messages , maxLogSize : logSize }
verifiedBlocks : make ( map [ common . Hash ] struct { } ) ,
}
return & pbftLog
return & pbftLog
}
}
// Blocks return the blocks stored in the log
// AddBlock add a new block into the log
func ( log * FBFTLog ) Blocks ( ) mapset . Set {
func ( log * FBFTLog ) AddBlock ( block * types . Block ) {
return log . blocks
log . blockLock . Lock ( )
defer log . blockLock . Unlock ( )
log . blocks [ block . Hash ( ) ] = block
}
}
// Messages return the messages stored in the log
// MarkBlockVerified marks the block as verified
func ( log * FBFTLog ) Messages ( ) mapset . Set {
func ( log * FBFTLog ) MarkBlockVerified ( block * types . Block ) {
return log . messages
log . blockLock . Lock ( )
defer log . blockLock . Unlock ( )
log . verifiedBlocks [ block . Hash ( ) ] = struct { } { }
}
}
// AddBlock add a new block into the log
// IsBlockVerified checks whether the block is verified
func ( log * FBFTLog ) AddBlock ( block * types . Block ) {
func ( log * FBFTLog ) IsBlockVerified ( block * types . Block ) bool {
log . blocks . Add ( block )
log . blockLock . RLock ( )
defer log . blockLock . RUnlock ( )
_ , exist := log . verifiedBlocks [ block . Hash ( ) ]
return exist
}
}
// GetBlockByHash returns the block matches the given block hash
// GetBlockByHash returns the block matches the given block hash
func ( log * FBFTLog ) GetBlockByHash ( hash common . Hash ) * types . Block {
func ( log * FBFTLog ) GetBlockByHash ( hash common . Hash ) * types . Block {
var found * types . Block
log . blockLock . RLock ( )
it := log . Blocks ( ) . Iterator ( )
defer log . blockLock . RUnlock ( )
for block := range it . C {
if block . ( * types . Block ) . Header ( ) . Hash ( ) == hash {
return log . blocks [ hash ]
found = block . ( * types . Block )
it . Stop ( )
}
}
return found
}
}
// GetBlocksByNumber returns the blocks match the given block number
// GetBlocksByNumber returns the blocks match the given block number
func ( log * FBFTLog ) GetBlocksByNumber ( number uint64 ) [ ] * types . Block {
func ( log * FBFTLog ) GetBlocksByNumber ( number uint64 ) [ ] * types . Block {
found := [ ] * types . Block { }
log . blockLock . RLock ( )
it := log . Blocks ( ) . Iterator ( )
defer log . blockLock . RUnlock ( )
for block := range it . C {
if block . ( * types . Block ) . NumberU64 ( ) == number {
var blocks [ ] * types . Block
found = append ( found , block . ( * types . Block ) )
for _ , block := range log . blocks {
if block . NumberU64 ( ) == number {
blocks = append ( blocks , block )
}
}
}
}
return found
return blocks
}
}
// DeleteBlocksLessThan deletes blocks less than given block number
// DeleteBlocksLessThan deletes blocks less than given block number
func ( log * FBFTLog ) DeleteBlocksLessThan ( number uint64 ) {
func ( log * FBFTLog ) DeleteBlocksLessThan ( number uint64 ) {
found := mapset . NewSet ( )
log . blockLock . Lock ( )
it := log . Blocks ( ) . Iterator ( )
defer log . blockLock . Unlock ( )
for block := range it . C {
if block . ( * types . Block ) . NumberU64 ( ) < number {
for h , block := range log . blocks {
found . Add ( block )
if block . NumberU64 ( ) < number {
delete ( log . blocks , h )
delete ( log . verifiedBlocks , h )
}
}
}
}
log . blocks = log . blocks . Difference ( found )
}
}
// DeleteBlockByNumber deletes block of specific number
// DeleteBlockByNumber deletes block of specific number
func ( log * FBFTLog ) DeleteBlockByNumber ( number uint64 ) {
func ( log * FBFTLog ) DeleteBlockByNumber ( number uint64 ) {
found := mapset . NewSet ( )
log . blockLock . Lock ( )
it := log . Blocks ( ) . Iterator ( )
defer log . blockLock . Unlock ( )
for block := range it . C {
if block . ( * types . Block ) . NumberU64 ( ) == number {
for h , block := range log . blocks {
found . Add ( block )
if block . NumberU64 ( ) == number {
delete ( log . blocks , h )
delete ( log . verifiedBlocks , h )
}
}
}
}
log . blocks = log . blocks . Difference ( found )
}
}
// DeleteMessagesLessThan deletes messages less than given block number
// DeleteMessagesLessThan deletes messages less than given block number
func ( log * FBFTLog ) DeleteMessagesLessThan ( number uint64 ) {
func ( log * FBFTLog ) DeleteMessagesLessThan ( number uint64 ) {
found := mapset . NewSet ( )
log . msgLock . Lock ( )
it := log . Messages ( ) . Iterator ( )
defer log . msgLock . Unlock ( )
for msg := range it . C {
if msg . ( * FBFTMessage ) . BlockNum < number {
for h , msg := range log . messages {
found . Add ( msg )
if msg . BlockNum < number {
delete ( log . messages , h )
}
}
}
}
log . messages = log . messages . Difference ( found )
}
}
// AddMessage adds a pbft message into the log
// AddMessage adds a pbft message into the log
func ( log * FBFTLog ) AddMessage ( msg * FBFTMessage ) {
func ( log * FBFTLog ) AddMessage ( msg * FBFTMessage ) {
log . messages . Add ( msg )
log . msgLock . Lock ( )
defer log . msgLock . Unlock ( )
log . messages [ msg . id ( ) ] = msg
}
}
// GetMessagesByTypeSeqViewHash returns pbft messages with matching type, blockNum, viewID and blockHash
// GetMessagesByTypeSeqViewHash returns pbft messages with matching type, blockNum, viewID and blockHash
func ( log * FBFTLog ) GetMessagesByTypeSeqViewHash ( typ msg_pb . MessageType , blockNum uint64 , viewID uint64 , blockHash common . Hash ) [ ] * FBFTMessage {
func ( log * FBFTLog ) GetMessagesByTypeSeqViewHash ( typ msg_pb . MessageType , blockNum uint64 , viewID uint64 , blockHash common . Hash ) [ ] * FBFTMessage {
found := [ ] * FBFTMessage { }
log . msgLock . RLock ( )
it := log . Messages ( ) . Iterator ( )
defer log . msgLock . RUnlock ( )
for msg := range it . C {
if msg . ( * FBFTMessage ) . MessageType == typ &&
var found [ ] * FBFTMessage
msg . ( * FBFTMessage ) . BlockNum == blockNum &&
for _ , msg := range log . messages {
msg . ( * FBFTMessage ) . ViewID == viewID &&
if msg . MessageType == typ && msg . BlockNum == blockNum && msg . ViewID == viewID && msg . BlockHash == blockHash {
msg . ( * FBFTMessage ) . BlockHash == blockHash {
found = append ( found , msg )
found = append ( found , msg . ( * FBFTMessage ) )
}
}
}
}
return found
return found
@ -166,12 +209,13 @@ func (log *FBFTLog) GetMessagesByTypeSeqViewHash(typ msg_pb.MessageType, blockNu
// GetMessagesByTypeSeq returns pbft messages with matching type, blockNum
// GetMessagesByTypeSeq returns pbft messages with matching type, blockNum
func ( log * FBFTLog ) GetMessagesByTypeSeq ( typ msg_pb . MessageType , blockNum uint64 ) [ ] * FBFTMessage {
func ( log * FBFTLog ) GetMessagesByTypeSeq ( typ msg_pb . MessageType , blockNum uint64 ) [ ] * FBFTMessage {
found := [ ] * FBFTMessage { }
log . msgLock . RLock ( )
it := log . Messages ( ) . Iterator ( )
defer log . msgLock . RUnlock ( )
for msg := range it . C {
if msg . ( * FBFTMessage ) . MessageType == typ &&
var found [ ] * FBFTMessage
msg . ( * FBFTMessage ) . BlockNum == blockNum {
for _ , msg := range log . messages {
found = append ( found , msg . ( * FBFTMessage ) )
if msg . MessageType == typ && msg . BlockNum == blockNum {
found = append ( found , msg )
}
}
}
}
return found
return found
@ -179,13 +223,13 @@ func (log *FBFTLog) GetMessagesByTypeSeq(typ msg_pb.MessageType, blockNum uint64
// GetMessagesByTypeSeqHash returns pbft messages with matching type, blockNum
// GetMessagesByTypeSeqHash returns pbft messages with matching type, blockNum
func ( log * FBFTLog ) GetMessagesByTypeSeqHash ( typ msg_pb . MessageType , blockNum uint64 , blockHash common . Hash ) [ ] * FBFTMessage {
func ( log * FBFTLog ) GetMessagesByTypeSeqHash ( typ msg_pb . MessageType , blockNum uint64 , blockHash common . Hash ) [ ] * FBFTMessage {
found := [ ] * FBFTMessage { }
log . msgLock . RLock ( )
it := log . Messages ( ) . Iterator ( )
defer log . msgLock . RUnlock ( )
for msg := range it . C {
if msg . ( * FBFTMessage ) . MessageType == typ &&
var found [ ] * FBFTMessage
msg . ( * FBFTMessage ) . BlockNum == blockNum &&
for _ , msg := range log . messages {
msg . ( * FBFTMessage ) . BlockHash == blockHash {
if msg . MessageType == typ && msg . BlockNum == blockNum && msg . BlockHash == blockHash {
found = append ( found , msg . ( * FBFTMessage ) )
found = append ( found , msg )
}
}
}
}
return found
return found
@ -217,13 +261,15 @@ func (log *FBFTLog) HasMatchingViewPrepared(blockNum uint64, viewID uint64, bloc
// GetMessagesByTypeSeqView returns pbft messages with matching type, blockNum and viewID
// GetMessagesByTypeSeqView returns pbft messages with matching type, blockNum and viewID
func ( log * FBFTLog ) GetMessagesByTypeSeqView ( typ msg_pb . MessageType , blockNum uint64 , viewID uint64 ) [ ] * FBFTMessage {
func ( log * FBFTLog ) GetMessagesByTypeSeqView ( typ msg_pb . MessageType , blockNum uint64 , viewID uint64 ) [ ] * FBFTMessage {
found := [ ] * FBFTMessage { }
log . msgLock . RLock ( )
it := log . Messages ( ) . Iterator ( )
defer log . msgLock . RUnlock ( )
for msg := range it . C {
if msg . ( * FBFTMessage ) . MessageType != typ || msg . ( * FBFTMessage ) . BlockNum != blockNum || msg . ( * FBFTMessage ) . ViewID != viewID {
var found [ ] * FBFTMessage
for _ , msg := range log . messages {
if msg . MessageType != typ || msg . BlockNum != blockNum || msg . ViewID != viewID {
continue
continue
}
}
found = append ( found , msg . ( * FBFTMessage ) )
found = append ( found , msg )
}
}
return found
return found
}
}
@ -267,3 +313,37 @@ func ParseFBFTMessage(msg *msg_pb.Message) (*FBFTMessage, error) {
return & pbftMsg , nil
return & pbftMsg , nil
}
}
var errFBFTLogNotFound = errors . New ( "FBFT log not found" )
// GetCommittedBlockAndMsgsFromNumber get committed block and message starting from block number bn.
func ( log * FBFTLog ) GetCommittedBlockAndMsgsFromNumber ( bn uint64 , logger * zerolog . Logger ) ( * types . Block , * FBFTMessage , error ) {
msgs := log . GetMessagesByTypeSeq (
msg_pb . MessageType_COMMITTED , bn ,
)
if len ( msgs ) == 0 {
return nil , nil , errFBFTLogNotFound
}
if len ( msgs ) > 1 {
logger . Error ( ) . Int ( "numMsgs" , len ( msgs ) ) . Err ( errors . New ( "DANGER!!! multiple COMMITTED message in PBFT log observed" ) )
}
for i := range msgs {
block := log . GetBlockByHash ( msgs [ i ] . BlockHash )
if block == nil {
logger . Debug ( ) .
Uint64 ( "blockNum" , msgs [ i ] . BlockNum ) .
Uint64 ( "viewID" , msgs [ i ] . ViewID ) .
Str ( "blockHash" , msgs [ i ] . BlockHash . Hex ( ) ) .
Err ( errors . New ( "failed finding a matching block for committed message" ) )
continue
}
return block , msgs [ i ] , nil
}
return nil , nil , errFBFTLogNotFound
}
// PruneCacheBeforeBlock prune all blocks before bn
func ( log * FBFTLog ) PruneCacheBeforeBlock ( bn uint64 ) {
log . DeleteBlocksLessThan ( bn - 1 )
log . DeleteMessagesLessThan ( bn - 1 )
}