mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-18 09:53:48 +00:00
consensus/istanbul: Istanbul core implementation
This commit is contained in:
parent
f03a83a17b
commit
0abccdb81e
14 changed files with 1928 additions and 0 deletions
182
consensus/istanbul/core/backlog.go
Normal file
182
consensus/istanbul/core/backlog.go
Normal file
|
|
@ -0,0 +1,182 @@
|
||||||
|
// Copyright 2017 The go-ethereum Authors
|
||||||
|
// This file is part of the go-ethereum library.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is free software: you can redistribute it and/or modify
|
||||||
|
// it under the terms of the GNU Lesser General Public License as published by
|
||||||
|
// the Free Software Foundation, either version 3 of the License, or
|
||||||
|
// (at your option) any later version.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is distributed in the hope that it will be useful,
|
||||||
|
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||||
|
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||||
|
// GNU Lesser General Public License for more details.
|
||||||
|
//
|
||||||
|
// You should have received a copy of the GNU Lesser General Public License
|
||||||
|
// along with the go-ethereum library. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
|
package core
|
||||||
|
|
||||||
|
import (
|
||||||
|
"github.com/ethereum/go-ethereum/consensus/istanbul"
|
||||||
|
"gopkg.in/karalabe/cookiejar.v2/collections/prque"
|
||||||
|
)
|
||||||
|
|
||||||
|
var (
|
||||||
|
// msgPriority is defined for calculating processing priority to speedup consensus
|
||||||
|
// msgPreprepare > msgCommit > msgPrepare
|
||||||
|
msgPriority = map[uint64]int{
|
||||||
|
msgPreprepare: 1,
|
||||||
|
msgCommit: 2,
|
||||||
|
msgPrepare: 3,
|
||||||
|
}
|
||||||
|
)
|
||||||
|
|
||||||
|
// checkMessage checks the message state
|
||||||
|
// return errInvalidMessage if the message is invalid
|
||||||
|
// return errFutureMessage if the message view is larger than current view
|
||||||
|
// return errOldMessage if the message view is smaller than current view
|
||||||
|
func (c *core) checkMessage(msgCode uint64, view *istanbul.View) error {
|
||||||
|
if view == nil || view.Sequence == nil || view.Round == nil {
|
||||||
|
return errInvalidMessage
|
||||||
|
}
|
||||||
|
|
||||||
|
if msgCode == msgRoundChange {
|
||||||
|
if view.Sequence.Cmp(c.currentView().Sequence) > 0 {
|
||||||
|
return errFutureMessage
|
||||||
|
} else if view.Cmp(c.currentView()) < 0 {
|
||||||
|
return errOldMessage
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
if view.Cmp(c.currentView()) > 0 {
|
||||||
|
return errFutureMessage
|
||||||
|
}
|
||||||
|
|
||||||
|
if view.Cmp(c.currentView()) < 0 {
|
||||||
|
return errOldMessage
|
||||||
|
}
|
||||||
|
|
||||||
|
if c.waitingForRoundChange {
|
||||||
|
return errFutureMessage
|
||||||
|
}
|
||||||
|
|
||||||
|
// StateAcceptRequest only accepts msgPreprepare
|
||||||
|
// other messages are future messages
|
||||||
|
if c.state == StateAcceptRequest {
|
||||||
|
if msgCode > msgPreprepare {
|
||||||
|
return errFutureMessage
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// For states(StatePreprepared, StatePrepared, StateCommitted),
|
||||||
|
// can accept all message types if processing with same view
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) storeBacklog(msg *message, src istanbul.Validator) {
|
||||||
|
logger := c.logger.New("from", src, "state", c.state)
|
||||||
|
|
||||||
|
if src.Address() == c.Address() {
|
||||||
|
logger.Warn("Backlog from self")
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
logger.Trace("Store future message")
|
||||||
|
|
||||||
|
c.backlogsMu.Lock()
|
||||||
|
defer c.backlogsMu.Unlock()
|
||||||
|
|
||||||
|
backlog := c.backlogs[src]
|
||||||
|
if backlog == nil {
|
||||||
|
backlog = prque.New()
|
||||||
|
}
|
||||||
|
switch msg.Code {
|
||||||
|
case msgPreprepare:
|
||||||
|
var p *istanbul.Preprepare
|
||||||
|
err := msg.Decode(&p)
|
||||||
|
if err == nil {
|
||||||
|
backlog.Push(msg, toPriority(msg.Code, p.View))
|
||||||
|
}
|
||||||
|
// for msgRoundChange, msgPrepare and msgCommit cases
|
||||||
|
default:
|
||||||
|
var p *istanbul.Subject
|
||||||
|
err := msg.Decode(&p)
|
||||||
|
if err == nil {
|
||||||
|
backlog.Push(msg, toPriority(msg.Code, p.View))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
c.backlogs[src] = backlog
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) processBacklog() {
|
||||||
|
c.backlogsMu.Lock()
|
||||||
|
defer c.backlogsMu.Unlock()
|
||||||
|
|
||||||
|
for src, backlog := range c.backlogs {
|
||||||
|
if backlog == nil {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
logger := c.logger.New("from", src, "state", c.state)
|
||||||
|
isFuture := false
|
||||||
|
|
||||||
|
// We stop processing if
|
||||||
|
// 1. backlog is empty
|
||||||
|
// 2. The first message in queue is a future message
|
||||||
|
for !(backlog.Empty() || isFuture) {
|
||||||
|
m, prio := backlog.Pop()
|
||||||
|
msg := m.(*message)
|
||||||
|
var view *istanbul.View
|
||||||
|
switch msg.Code {
|
||||||
|
case msgPreprepare:
|
||||||
|
var m *istanbul.Preprepare
|
||||||
|
err := msg.Decode(&m)
|
||||||
|
if err == nil {
|
||||||
|
view = m.View
|
||||||
|
}
|
||||||
|
// for msgRoundChange, msgPrepare and msgCommit cases
|
||||||
|
default:
|
||||||
|
var sub *istanbul.Subject
|
||||||
|
err := msg.Decode(&sub)
|
||||||
|
if err == nil {
|
||||||
|
view = sub.View
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if view == nil {
|
||||||
|
logger.Debug("Nil view", "msg", msg)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
// Push back if it's a future message
|
||||||
|
err := c.checkMessage(msg.Code, view)
|
||||||
|
if err != nil {
|
||||||
|
if err == errFutureMessage {
|
||||||
|
logger.Trace("Stop processing backlog", "msg", msg)
|
||||||
|
backlog.Push(msg, prio)
|
||||||
|
isFuture = true
|
||||||
|
break
|
||||||
|
}
|
||||||
|
logger.Trace("Skip the backlog event", "msg", msg, "err", err)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
logger.Trace("Post backlog event", "msg", msg)
|
||||||
|
|
||||||
|
go c.sendEvent(backlogEvent{
|
||||||
|
src: src,
|
||||||
|
msg: msg,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func toPriority(msgCode uint64, view *istanbul.View) float32 {
|
||||||
|
if msgCode == msgRoundChange {
|
||||||
|
// For msgRoundChange, set the message priority based on its sequence
|
||||||
|
return -float32(view.Sequence.Uint64() * 1000)
|
||||||
|
}
|
||||||
|
// FIXME: round will be reset as 0 while new sequence
|
||||||
|
// 10 * Round limits the range of message code is from 0 to 9
|
||||||
|
// 1000 * Sequence limits the range of round is from 0 to 99
|
||||||
|
return -float32(view.Sequence.Uint64()*1000 + view.Round.Uint64()*10 + uint64(msgPriority[msgCode]))
|
||||||
|
}
|
||||||
107
consensus/istanbul/core/commit.go
Normal file
107
consensus/istanbul/core/commit.go
Normal file
|
|
@ -0,0 +1,107 @@
|
||||||
|
// Copyright 2017 The go-ethereum Authors
|
||||||
|
// This file is part of the go-ethereum library.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is free software: you can redistribute it and/or modify
|
||||||
|
// it under the terms of the GNU Lesser General Public License as published by
|
||||||
|
// the Free Software Foundation, either version 3 of the License, or
|
||||||
|
// (at your option) any later version.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is distributed in the hope that it will be useful,
|
||||||
|
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||||
|
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||||
|
// GNU Lesser General Public License for more details.
|
||||||
|
//
|
||||||
|
// You should have received a copy of the GNU Lesser General Public License
|
||||||
|
// along with the go-ethereum library. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
|
package core
|
||||||
|
|
||||||
|
import (
|
||||||
|
"reflect"
|
||||||
|
|
||||||
|
"github.com/ethereum/go-ethereum/common"
|
||||||
|
"github.com/ethereum/go-ethereum/consensus/istanbul"
|
||||||
|
)
|
||||||
|
|
||||||
|
func (c *core) sendCommit() {
|
||||||
|
sub := c.current.Subject()
|
||||||
|
c.broadcastCommit(sub)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) sendCommitForOldBlock(view *istanbul.View, digest common.Hash) {
|
||||||
|
sub := &istanbul.Subject{
|
||||||
|
View: view,
|
||||||
|
Digest: digest,
|
||||||
|
}
|
||||||
|
c.broadcastCommit(sub)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) broadcastCommit(sub *istanbul.Subject) {
|
||||||
|
logger := c.logger.New("state", c.state)
|
||||||
|
|
||||||
|
encodedSubject, err := Encode(sub)
|
||||||
|
if err != nil {
|
||||||
|
logger.Error("Failed to encode", "subject", sub)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
c.broadcast(&message{
|
||||||
|
Code: msgCommit,
|
||||||
|
Msg: encodedSubject,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) handleCommit(msg *message, src istanbul.Validator) error {
|
||||||
|
// Decode COMMIT message
|
||||||
|
var commit *istanbul.Subject
|
||||||
|
err := msg.Decode(&commit)
|
||||||
|
if err != nil {
|
||||||
|
return errFailedDecodeCommit
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := c.checkMessage(msgCommit, commit.View); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := c.verifyCommit(commit, src); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
c.acceptCommit(msg, src)
|
||||||
|
|
||||||
|
// Commit the proposal once we have enough COMMIT messages and we are not in the Committed state.
|
||||||
|
//
|
||||||
|
// If we already have a proposal, we may have chance to speed up the consensus process
|
||||||
|
// by committing the proposal without PREPARE messages.
|
||||||
|
if c.current.Commits.Size() > 2*c.valSet.F() && c.state.Cmp(StateCommitted) < 0 {
|
||||||
|
// Still need to call LockHash here since state can skip Prepared state and jump directly to the Committed state.
|
||||||
|
c.current.LockHash()
|
||||||
|
c.commit()
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// verifyCommit verifies if the received COMMIT message is equivalent to our subject
|
||||||
|
func (c *core) verifyCommit(commit *istanbul.Subject, src istanbul.Validator) error {
|
||||||
|
logger := c.logger.New("from", src, "state", c.state)
|
||||||
|
|
||||||
|
sub := c.current.Subject()
|
||||||
|
if !reflect.DeepEqual(commit, sub) {
|
||||||
|
logger.Warn("Inconsistent subjects between commit and proposal", "expected", sub, "got", commit)
|
||||||
|
return errInconsistentSubject
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) acceptCommit(msg *message, src istanbul.Validator) error {
|
||||||
|
logger := c.logger.New("from", src, "state", c.state)
|
||||||
|
|
||||||
|
// Add the COMMIT message to current round state
|
||||||
|
if err := c.current.Commits.Add(msg); err != nil {
|
||||||
|
logger.Error("Failed to record commit message", "msg", msg, "err", err)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
341
consensus/istanbul/core/core.go
Normal file
341
consensus/istanbul/core/core.go
Normal file
|
|
@ -0,0 +1,341 @@
|
||||||
|
// Copyright 2017 The go-ethereum Authors
|
||||||
|
// This file is part of the go-ethereum library.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is free software: you can redistribute it and/or modify
|
||||||
|
// it under the terms of the GNU Lesser General Public License as published by
|
||||||
|
// the Free Software Foundation, either version 3 of the License, or
|
||||||
|
// (at your option) any later version.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is distributed in the hope that it will be useful,
|
||||||
|
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||||
|
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||||
|
// GNU Lesser General Public License for more details.
|
||||||
|
//
|
||||||
|
// You should have received a copy of the GNU Lesser General Public License
|
||||||
|
// along with the go-ethereum library. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
|
package core
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"math"
|
||||||
|
"math/big"
|
||||||
|
"sync"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/ethereum/go-ethereum/common"
|
||||||
|
"github.com/ethereum/go-ethereum/consensus/istanbul"
|
||||||
|
"github.com/ethereum/go-ethereum/core/types"
|
||||||
|
"github.com/ethereum/go-ethereum/event"
|
||||||
|
"github.com/ethereum/go-ethereum/log"
|
||||||
|
"github.com/ethereum/go-ethereum/metrics"
|
||||||
|
"gopkg.in/karalabe/cookiejar.v2/collections/prque"
|
||||||
|
)
|
||||||
|
|
||||||
|
// New creates an Istanbul consensus core
|
||||||
|
func New(backend istanbul.Backend, config *istanbul.Config) Engine {
|
||||||
|
c := &core{
|
||||||
|
config: config,
|
||||||
|
address: backend.Address(),
|
||||||
|
state: StateAcceptRequest,
|
||||||
|
handlerWg: new(sync.WaitGroup),
|
||||||
|
logger: log.New("address", backend.Address()),
|
||||||
|
backend: backend,
|
||||||
|
backlogs: make(map[istanbul.Validator]*prque.Prque),
|
||||||
|
backlogsMu: new(sync.Mutex),
|
||||||
|
pendingRequests: prque.New(),
|
||||||
|
pendingRequestsMu: new(sync.Mutex),
|
||||||
|
consensusTimestamp: time.Time{},
|
||||||
|
roundMeter: metrics.NewRegisteredMeter("consensus/istanbul/core/round", nil),
|
||||||
|
sequenceMeter: metrics.NewRegisteredMeter("consensus/istanbul/core/sequence", nil),
|
||||||
|
consensusTimer: metrics.NewRegisteredTimer("consensus/istanbul/core/consensus", nil),
|
||||||
|
}
|
||||||
|
c.validateFn = c.checkValidatorSignature
|
||||||
|
return c
|
||||||
|
}
|
||||||
|
|
||||||
|
// ----------------------------------------------------------------------------
|
||||||
|
|
||||||
|
type core struct {
|
||||||
|
config *istanbul.Config
|
||||||
|
address common.Address
|
||||||
|
state State
|
||||||
|
logger log.Logger
|
||||||
|
|
||||||
|
backend istanbul.Backend
|
||||||
|
events *event.TypeMuxSubscription
|
||||||
|
finalCommittedSub *event.TypeMuxSubscription
|
||||||
|
timeoutSub *event.TypeMuxSubscription
|
||||||
|
futurePreprepareTimer *time.Timer
|
||||||
|
|
||||||
|
valSet istanbul.ValidatorSet
|
||||||
|
waitingForRoundChange bool
|
||||||
|
validateFn func([]byte, []byte) (common.Address, error)
|
||||||
|
|
||||||
|
backlogs map[istanbul.Validator]*prque.Prque
|
||||||
|
backlogsMu *sync.Mutex
|
||||||
|
|
||||||
|
current *roundState
|
||||||
|
handlerWg *sync.WaitGroup
|
||||||
|
|
||||||
|
roundChangeSet *roundChangeSet
|
||||||
|
roundChangeTimer *time.Timer
|
||||||
|
|
||||||
|
pendingRequests *prque.Prque
|
||||||
|
pendingRequestsMu *sync.Mutex
|
||||||
|
|
||||||
|
consensusTimestamp time.Time
|
||||||
|
// the meter to record the round change rate
|
||||||
|
roundMeter metrics.Meter
|
||||||
|
// the meter to record the sequence update rate
|
||||||
|
sequenceMeter metrics.Meter
|
||||||
|
// the timer to record consensus duration (from accepting a preprepare to final committed stage)
|
||||||
|
consensusTimer metrics.Timer
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) finalizeMessage(msg *message) ([]byte, error) {
|
||||||
|
var err error
|
||||||
|
// Add sender address
|
||||||
|
msg.Address = c.Address()
|
||||||
|
|
||||||
|
// Add proof of consensus
|
||||||
|
msg.CommittedSeal = []byte{}
|
||||||
|
// Assign the CommittedSeal if it's a COMMIT message and proposal is not nil
|
||||||
|
if msg.Code == msgCommit && c.current.Proposal() != nil {
|
||||||
|
seal := PrepareCommittedSeal(c.current.Proposal().Hash())
|
||||||
|
msg.CommittedSeal, err = c.backend.Sign(seal)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Sign message
|
||||||
|
data, err := msg.PayloadNoSig()
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
msg.Signature, err = c.backend.Sign(data)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
// Convert to payload
|
||||||
|
payload, err := msg.Payload()
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
return payload, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) broadcast(msg *message) {
|
||||||
|
logger := c.logger.New("state", c.state)
|
||||||
|
|
||||||
|
payload, err := c.finalizeMessage(msg)
|
||||||
|
if err != nil {
|
||||||
|
logger.Error("Failed to finalize message", "msg", msg, "err", err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
// Broadcast payload
|
||||||
|
if err = c.backend.Broadcast(c.valSet, payload); err != nil {
|
||||||
|
logger.Error("Failed to broadcast message", "msg", msg, "err", err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) currentView() *istanbul.View {
|
||||||
|
return &istanbul.View{
|
||||||
|
Sequence: new(big.Int).Set(c.current.Sequence()),
|
||||||
|
Round: new(big.Int).Set(c.current.Round()),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) isProposer() bool {
|
||||||
|
v := c.valSet
|
||||||
|
if v == nil {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
return v.IsProposer(c.backend.Address())
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) commit() {
|
||||||
|
c.setState(StateCommitted)
|
||||||
|
|
||||||
|
proposal := c.current.Proposal()
|
||||||
|
if proposal != nil {
|
||||||
|
committedSeals := make([][]byte, c.current.Commits.Size())
|
||||||
|
for i, v := range c.current.Commits.Values() {
|
||||||
|
committedSeals[i] = make([]byte, types.IstanbulExtraSeal)
|
||||||
|
copy(committedSeals[i][:], v.CommittedSeal[:])
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := c.backend.Commit(proposal, committedSeals); err != nil {
|
||||||
|
c.current.UnlockHash() //Unlock block when insertion fails
|
||||||
|
c.sendNextRoundChange()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// startNewRound starts a new round. if round equals to 0, it means to starts a new sequence
|
||||||
|
func (c *core) startNewRound(round *big.Int) {
|
||||||
|
var logger log.Logger
|
||||||
|
if c.current == nil {
|
||||||
|
logger = c.logger.New("old_round", -1, "old_seq", 0)
|
||||||
|
} else {
|
||||||
|
logger = c.logger.New("old_round", c.current.Round(), "old_seq", c.current.Sequence())
|
||||||
|
}
|
||||||
|
|
||||||
|
roundChange := false
|
||||||
|
// Try to get last proposal
|
||||||
|
lastProposal, lastProposer := c.backend.LastProposal()
|
||||||
|
if c.current == nil {
|
||||||
|
logger.Trace("Start to the initial round")
|
||||||
|
} else if lastProposal.Number().Cmp(c.current.Sequence()) >= 0 {
|
||||||
|
diff := new(big.Int).Sub(lastProposal.Number(), c.current.Sequence())
|
||||||
|
c.sequenceMeter.Mark(new(big.Int).Add(diff, common.Big1).Int64())
|
||||||
|
|
||||||
|
if !c.consensusTimestamp.IsZero() {
|
||||||
|
c.consensusTimer.UpdateSince(c.consensusTimestamp)
|
||||||
|
c.consensusTimestamp = time.Time{}
|
||||||
|
}
|
||||||
|
logger.Trace("Catch up latest proposal", "number", lastProposal.Number().Uint64(), "hash", lastProposal.Hash())
|
||||||
|
} else if lastProposal.Number().Cmp(big.NewInt(c.current.Sequence().Int64()-1)) == 0 {
|
||||||
|
if round.Cmp(common.Big0) == 0 {
|
||||||
|
// same seq and round, don't need to start new round
|
||||||
|
return
|
||||||
|
} else if round.Cmp(c.current.Round()) < 0 {
|
||||||
|
logger.Warn("New round should not be smaller than current round", "seq", lastProposal.Number().Int64(), "new_round", round, "old_round", c.current.Round())
|
||||||
|
return
|
||||||
|
}
|
||||||
|
roundChange = true
|
||||||
|
} else {
|
||||||
|
logger.Warn("New sequence should be larger than current sequence", "new_seq", lastProposal.Number().Int64())
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
var newView *istanbul.View
|
||||||
|
if roundChange {
|
||||||
|
newView = &istanbul.View{
|
||||||
|
Sequence: new(big.Int).Set(c.current.Sequence()),
|
||||||
|
Round: new(big.Int).Set(round),
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
newView = &istanbul.View{
|
||||||
|
Sequence: new(big.Int).Add(lastProposal.Number(), common.Big1),
|
||||||
|
Round: new(big.Int),
|
||||||
|
}
|
||||||
|
c.valSet = c.backend.Validators(lastProposal)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Update logger
|
||||||
|
logger = logger.New("old_proposer", c.valSet.GetProposer())
|
||||||
|
// Clear invalid ROUND CHANGE messages
|
||||||
|
c.roundChangeSet = newRoundChangeSet(c.valSet)
|
||||||
|
// New snapshot for new round
|
||||||
|
c.updateRoundState(newView, c.valSet, roundChange)
|
||||||
|
// Calculate new proposer
|
||||||
|
c.valSet.CalcProposer(lastProposer, newView.Round.Uint64())
|
||||||
|
c.waitingForRoundChange = false
|
||||||
|
c.setState(StateAcceptRequest)
|
||||||
|
if roundChange && c.isProposer() && c.current != nil {
|
||||||
|
// If it is locked, propose the old proposal
|
||||||
|
// If we have pending request, propose pending request
|
||||||
|
if c.current.IsHashLocked() {
|
||||||
|
r := &istanbul.Request{
|
||||||
|
Proposal: c.current.Proposal(), //c.current.Proposal would be the locked proposal by previous proposer, see updateRoundState
|
||||||
|
}
|
||||||
|
c.sendPreprepare(r)
|
||||||
|
} else if c.current.pendingRequest != nil {
|
||||||
|
c.sendPreprepare(c.current.pendingRequest)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
c.newRoundChangeTimer()
|
||||||
|
|
||||||
|
logger.Debug("New round", "new_round", newView.Round, "new_seq", newView.Sequence, "new_proposer", c.valSet.GetProposer(), "valSet", c.valSet.List(), "size", c.valSet.Size(), "isProposer", c.isProposer())
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) catchUpRound(view *istanbul.View) {
|
||||||
|
logger := c.logger.New("old_round", c.current.Round(), "old_seq", c.current.Sequence(), "old_proposer", c.valSet.GetProposer())
|
||||||
|
|
||||||
|
if view.Round.Cmp(c.current.Round()) > 0 {
|
||||||
|
c.roundMeter.Mark(new(big.Int).Sub(view.Round, c.current.Round()).Int64())
|
||||||
|
}
|
||||||
|
c.waitingForRoundChange = true
|
||||||
|
|
||||||
|
// Need to keep block locked for round catching up
|
||||||
|
c.updateRoundState(view, c.valSet, true)
|
||||||
|
c.roundChangeSet.Clear(view.Round)
|
||||||
|
c.newRoundChangeTimer()
|
||||||
|
|
||||||
|
logger.Trace("Catch up round", "new_round", view.Round, "new_seq", view.Sequence, "new_proposer", c.valSet)
|
||||||
|
}
|
||||||
|
|
||||||
|
// updateRoundState updates round state by checking if locking block is necessary
|
||||||
|
func (c *core) updateRoundState(view *istanbul.View, validatorSet istanbul.ValidatorSet, roundChange bool) {
|
||||||
|
// Lock only if both roundChange is true and it is locked
|
||||||
|
if roundChange && c.current != nil {
|
||||||
|
if c.current.IsHashLocked() {
|
||||||
|
c.current = newRoundState(view, validatorSet, c.current.GetLockedHash(), c.current.Preprepare, c.current.pendingRequest, c.backend.HasBadProposal)
|
||||||
|
} else {
|
||||||
|
c.current = newRoundState(view, validatorSet, common.Hash{}, nil, c.current.pendingRequest, c.backend.HasBadProposal)
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
c.current = newRoundState(view, validatorSet, common.Hash{}, nil, nil, c.backend.HasBadProposal)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) setState(state State) {
|
||||||
|
if c.state != state {
|
||||||
|
c.state = state
|
||||||
|
}
|
||||||
|
if state == StateAcceptRequest {
|
||||||
|
c.processPendingRequests()
|
||||||
|
}
|
||||||
|
c.processBacklog()
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) Address() common.Address {
|
||||||
|
return c.address
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) stopFuturePreprepareTimer() {
|
||||||
|
if c.futurePreprepareTimer != nil {
|
||||||
|
c.futurePreprepareTimer.Stop()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) stopTimer() {
|
||||||
|
c.stopFuturePreprepareTimer()
|
||||||
|
if c.roundChangeTimer != nil {
|
||||||
|
c.roundChangeTimer.Stop()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) newRoundChangeTimer() {
|
||||||
|
c.stopTimer()
|
||||||
|
|
||||||
|
// set timeout based on the round number
|
||||||
|
timeout := time.Duration(c.config.RequestTimeout) * time.Millisecond
|
||||||
|
round := c.current.Round().Uint64()
|
||||||
|
if round > 0 {
|
||||||
|
timeout += time.Duration(math.Pow(2, float64(round))) * time.Second
|
||||||
|
}
|
||||||
|
|
||||||
|
c.roundChangeTimer = time.AfterFunc(timeout, func() {
|
||||||
|
c.sendEvent(timeoutEvent{})
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) checkValidatorSignature(data []byte, sig []byte) (common.Address, error) {
|
||||||
|
return istanbul.CheckValidatorSignature(c.valSet, data, sig)
|
||||||
|
}
|
||||||
|
|
||||||
|
// PrepareCommittedSeal returns a committed seal for the given hash
|
||||||
|
func PrepareCommittedSeal(hash common.Hash) []byte {
|
||||||
|
var buf bytes.Buffer
|
||||||
|
buf.Write(hash.Bytes())
|
||||||
|
buf.Write([]byte{byte(msgCommit)})
|
||||||
|
return buf.Bytes()
|
||||||
|
}
|
||||||
46
consensus/istanbul/core/errors.go
Normal file
46
consensus/istanbul/core/errors.go
Normal file
|
|
@ -0,0 +1,46 @@
|
||||||
|
// Copyright 2017 The go-ethereum Authors
|
||||||
|
// This file is part of the go-ethereum library.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is free software: you can redistribute it and/or modify
|
||||||
|
// it under the terms of the GNU Lesser General Public License as published by
|
||||||
|
// the Free Software Foundation, either version 3 of the License, or
|
||||||
|
// (at your option) any later version.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is distributed in the hope that it will be useful,
|
||||||
|
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||||
|
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||||
|
// GNU Lesser General Public License for more details.
|
||||||
|
//
|
||||||
|
// You should have received a copy of the GNU Lesser General Public License
|
||||||
|
// along with the go-ethereum library. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
|
package core
|
||||||
|
|
||||||
|
import "errors"
|
||||||
|
|
||||||
|
var (
|
||||||
|
// errInconsistentSubject is returned when received subject is different from
|
||||||
|
// current subject.
|
||||||
|
errInconsistentSubject = errors.New("inconsistent subjects")
|
||||||
|
// errNotFromProposer is returned when received message is supposed to be from
|
||||||
|
// proposer.
|
||||||
|
errNotFromProposer = errors.New("message does not come from proposer")
|
||||||
|
// errIgnored is returned when a message was ignored.
|
||||||
|
errIgnored = errors.New("message is ignored")
|
||||||
|
// errFutureMessage is returned when current view is earlier than the
|
||||||
|
// view of the received message.
|
||||||
|
errFutureMessage = errors.New("future message")
|
||||||
|
// errOldMessage is returned when the received message's view is earlier
|
||||||
|
// than current view.
|
||||||
|
errOldMessage = errors.New("old message")
|
||||||
|
// errInvalidMessage is returned when the message is malformed.
|
||||||
|
errInvalidMessage = errors.New("invalid message")
|
||||||
|
// errFailedDecodePreprepare is returned when the PRE-PREPARE message is malformed.
|
||||||
|
errFailedDecodePreprepare = errors.New("failed to decode PRE-PREPARE")
|
||||||
|
// errFailedDecodePrepare is returned when the PREPARE message is malformed.
|
||||||
|
errFailedDecodePrepare = errors.New("failed to decode PREPARE")
|
||||||
|
// errFailedDecodeCommit is returned when the COMMIT message is malformed.
|
||||||
|
errFailedDecodeCommit = errors.New("failed to decode COMMIT")
|
||||||
|
// errFailedDecodeMessageSet is returned when the message set is malformed.
|
||||||
|
errFailedDecodeMessageSet = errors.New("failed to decode message set")
|
||||||
|
)
|
||||||
28
consensus/istanbul/core/events.go
Normal file
28
consensus/istanbul/core/events.go
Normal file
|
|
@ -0,0 +1,28 @@
|
||||||
|
// Copyright 2017 The go-ethereum Authors
|
||||||
|
// This file is part of the go-ethereum library.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is free software: you can redistribute it and/or modify
|
||||||
|
// it under the terms of the GNU Lesser General Public License as published by
|
||||||
|
// the Free Software Foundation, either version 3 of the License, or
|
||||||
|
// (at your option) any later version.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is distributed in the hope that it will be useful,
|
||||||
|
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||||
|
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||||
|
// GNU Lesser General Public License for more details.
|
||||||
|
//
|
||||||
|
// You should have received a copy of the GNU Lesser General Public License
|
||||||
|
// along with the go-ethereum library. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
|
package core
|
||||||
|
|
||||||
|
import (
|
||||||
|
"github.com/ethereum/go-ethereum/consensus/istanbul"
|
||||||
|
)
|
||||||
|
|
||||||
|
type backlogEvent struct {
|
||||||
|
src istanbul.Validator
|
||||||
|
msg *message
|
||||||
|
}
|
||||||
|
|
||||||
|
type timeoutEvent struct{}
|
||||||
26
consensus/istanbul/core/final_committed.go
Normal file
26
consensus/istanbul/core/final_committed.go
Normal file
|
|
@ -0,0 +1,26 @@
|
||||||
|
// Copyright 2017 The go-ethereum Authors
|
||||||
|
// This file is part of the go-ethereum library.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is free software: you can redistribute it and/or modify
|
||||||
|
// it under the terms of the GNU Lesser General Public License as published by
|
||||||
|
// the Free Software Foundation, either version 3 of the License, or
|
||||||
|
// (at your option) any later version.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is distributed in the hope that it will be useful,
|
||||||
|
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||||
|
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||||
|
// GNU Lesser General Public License for more details.
|
||||||
|
//
|
||||||
|
// You should have received a copy of the GNU Lesser General Public License
|
||||||
|
// along with the go-ethereum library. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
|
package core
|
||||||
|
|
||||||
|
import "github.com/ethereum/go-ethereum/common"
|
||||||
|
|
||||||
|
func (c *core) handleFinalCommitted() error {
|
||||||
|
logger := c.logger.New("state", c.state)
|
||||||
|
logger.Trace("Received a final committed proposal")
|
||||||
|
c.startNewRound(common.Big0)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
202
consensus/istanbul/core/handler.go
Normal file
202
consensus/istanbul/core/handler.go
Normal file
|
|
@ -0,0 +1,202 @@
|
||||||
|
// Copyright 2017 The go-ethereum Authors
|
||||||
|
// This file is part of the go-ethereum library.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is free software: you can redistribute it and/or modify
|
||||||
|
// it under the terms of the GNU Lesser General Public License as published by
|
||||||
|
// the Free Software Foundation, either version 3 of the License, or
|
||||||
|
// (at your option) any later version.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is distributed in the hope that it will be useful,
|
||||||
|
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||||
|
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||||
|
// GNU Lesser General Public License for more details.
|
||||||
|
//
|
||||||
|
// You should have received a copy of the GNU Lesser General Public License
|
||||||
|
// along with the go-ethereum library. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
|
package core
|
||||||
|
|
||||||
|
import (
|
||||||
|
"github.com/ethereum/go-ethereum/common"
|
||||||
|
"github.com/ethereum/go-ethereum/consensus/istanbul"
|
||||||
|
)
|
||||||
|
|
||||||
|
// Start implements core.Engine.Start
|
||||||
|
func (c *core) Start() error {
|
||||||
|
// Start a new round from last sequence + 1
|
||||||
|
c.startNewRound(common.Big0)
|
||||||
|
|
||||||
|
// Tests will handle events itself, so we have to make subscribeEvents()
|
||||||
|
// be able to call in test.
|
||||||
|
c.subscribeEvents()
|
||||||
|
go c.handleEvents()
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Stop implements core.Engine.Stop
|
||||||
|
func (c *core) Stop() error {
|
||||||
|
c.stopTimer()
|
||||||
|
c.unsubscribeEvents()
|
||||||
|
|
||||||
|
// Make sure the handler goroutine exits
|
||||||
|
c.handlerWg.Wait()
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// ----------------------------------------------------------------------------
|
||||||
|
|
||||||
|
// Subscribe both internal and external events
|
||||||
|
func (c *core) subscribeEvents() {
|
||||||
|
c.events = c.backend.EventMux().Subscribe(
|
||||||
|
// external events
|
||||||
|
istanbul.RequestEvent{},
|
||||||
|
istanbul.MessageEvent{},
|
||||||
|
// internal events
|
||||||
|
backlogEvent{},
|
||||||
|
)
|
||||||
|
c.timeoutSub = c.backend.EventMux().Subscribe(
|
||||||
|
timeoutEvent{},
|
||||||
|
)
|
||||||
|
c.finalCommittedSub = c.backend.EventMux().Subscribe(
|
||||||
|
istanbul.FinalCommittedEvent{},
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Unsubscribe all events
|
||||||
|
func (c *core) unsubscribeEvents() {
|
||||||
|
c.events.Unsubscribe()
|
||||||
|
c.timeoutSub.Unsubscribe()
|
||||||
|
c.finalCommittedSub.Unsubscribe()
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) handleEvents() {
|
||||||
|
// Clear state
|
||||||
|
defer func() {
|
||||||
|
c.current = nil
|
||||||
|
c.handlerWg.Done()
|
||||||
|
}()
|
||||||
|
|
||||||
|
c.handlerWg.Add(1)
|
||||||
|
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case event, ok := <-c.events.Chan():
|
||||||
|
if !ok {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
// A real event arrived, process interesting content
|
||||||
|
switch ev := event.Data.(type) {
|
||||||
|
case istanbul.RequestEvent:
|
||||||
|
r := &istanbul.Request{
|
||||||
|
Proposal: ev.Proposal,
|
||||||
|
}
|
||||||
|
err := c.handleRequest(r)
|
||||||
|
if err == errFutureMessage {
|
||||||
|
c.storeRequestMsg(r)
|
||||||
|
}
|
||||||
|
case istanbul.MessageEvent:
|
||||||
|
if err := c.handleMsg(ev.Payload); err == nil {
|
||||||
|
c.backend.Gossip(c.valSet, ev.Payload)
|
||||||
|
}
|
||||||
|
case backlogEvent:
|
||||||
|
// No need to check signature for internal messages
|
||||||
|
if err := c.handleCheckedMsg(ev.msg, ev.src); err == nil {
|
||||||
|
p, err := ev.msg.Payload()
|
||||||
|
if err != nil {
|
||||||
|
c.logger.Warn("Get message payload failed", "err", err)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
c.backend.Gossip(c.valSet, p)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
case _, ok := <-c.timeoutSub.Chan():
|
||||||
|
if !ok {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
c.handleTimeoutMsg()
|
||||||
|
case event, ok := <-c.finalCommittedSub.Chan():
|
||||||
|
if !ok {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
switch event.Data.(type) {
|
||||||
|
case istanbul.FinalCommittedEvent:
|
||||||
|
c.handleFinalCommitted()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// sendEvent sends events to mux
|
||||||
|
func (c *core) sendEvent(ev interface{}) {
|
||||||
|
c.backend.EventMux().Post(ev)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) handleMsg(payload []byte) error {
|
||||||
|
logger := c.logger.New()
|
||||||
|
|
||||||
|
// Decode message and check its signature
|
||||||
|
msg := new(message)
|
||||||
|
if err := msg.FromPayload(payload, c.validateFn); err != nil {
|
||||||
|
logger.Error("Failed to decode message from payload", "err", err)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
// Only accept message if the address is valid
|
||||||
|
_, src := c.valSet.GetByAddress(msg.Address)
|
||||||
|
if src == nil {
|
||||||
|
logger.Error("Invalid address in message", "msg", msg)
|
||||||
|
return istanbul.ErrUnauthorizedAddress
|
||||||
|
}
|
||||||
|
|
||||||
|
return c.handleCheckedMsg(msg, src)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) handleCheckedMsg(msg *message, src istanbul.Validator) error {
|
||||||
|
logger := c.logger.New("address", c.address, "from", src)
|
||||||
|
|
||||||
|
// Store the message if it's a future message
|
||||||
|
testBacklog := func(err error) error {
|
||||||
|
if err == errFutureMessage {
|
||||||
|
c.storeBacklog(msg, src)
|
||||||
|
}
|
||||||
|
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
switch msg.Code {
|
||||||
|
case msgPreprepare:
|
||||||
|
return testBacklog(c.handlePreprepare(msg, src))
|
||||||
|
case msgPrepare:
|
||||||
|
return testBacklog(c.handlePrepare(msg, src))
|
||||||
|
case msgCommit:
|
||||||
|
return testBacklog(c.handleCommit(msg, src))
|
||||||
|
case msgRoundChange:
|
||||||
|
return testBacklog(c.handleRoundChange(msg, src))
|
||||||
|
default:
|
||||||
|
logger.Error("Invalid message", "msg", msg)
|
||||||
|
}
|
||||||
|
|
||||||
|
return errInvalidMessage
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) handleTimeoutMsg() {
|
||||||
|
// If we're not waiting for round change yet, we can try to catch up
|
||||||
|
// the max round with F+1 round change message. We only need to catch up
|
||||||
|
// if the max round is larger than current round.
|
||||||
|
if !c.waitingForRoundChange {
|
||||||
|
maxRound := c.roundChangeSet.MaxRound(c.valSet.F() + 1)
|
||||||
|
if maxRound != nil && maxRound.Cmp(c.current.Round()) > 0 {
|
||||||
|
c.sendRoundChange(maxRound)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
lastProposal, _ := c.backend.LastProposal()
|
||||||
|
if lastProposal != nil && lastProposal.Number().Cmp(c.current.Sequence()) >= 0 {
|
||||||
|
c.logger.Trace("round change timeout, catch up latest sequence", "number", lastProposal.Number().Uint64())
|
||||||
|
c.startNewRound(common.Big0)
|
||||||
|
} else {
|
||||||
|
c.sendNextRoundChange()
|
||||||
|
}
|
||||||
|
}
|
||||||
115
consensus/istanbul/core/message_set.go
Normal file
115
consensus/istanbul/core/message_set.go
Normal file
|
|
@ -0,0 +1,115 @@
|
||||||
|
// Copyright 2017 The go-ethereum Authors
|
||||||
|
// This file is part of the go-ethereum library.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is free software: you can redistribute it and/or modify
|
||||||
|
// it under the terms of the GNU Lesser General Public License as published by
|
||||||
|
// the Free Software Foundation, either version 3 of the License, or
|
||||||
|
// (at your option) any later version.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is distributed in the hope that it will be useful,
|
||||||
|
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||||
|
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||||
|
// GNU Lesser General Public License for more details.
|
||||||
|
//
|
||||||
|
// You should have received a copy of the GNU Lesser General Public License
|
||||||
|
// along with the go-ethereum library. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
|
package core
|
||||||
|
|
||||||
|
import (
|
||||||
|
"fmt"
|
||||||
|
"math/big"
|
||||||
|
"strings"
|
||||||
|
"sync"
|
||||||
|
|
||||||
|
"github.com/ethereum/go-ethereum/common"
|
||||||
|
"github.com/ethereum/go-ethereum/consensus/istanbul"
|
||||||
|
)
|
||||||
|
|
||||||
|
// Construct a new message set to accumulate messages for given sequence/view number.
|
||||||
|
func newMessageSet(valSet istanbul.ValidatorSet) *messageSet {
|
||||||
|
return &messageSet{
|
||||||
|
view: &istanbul.View{
|
||||||
|
Round: new(big.Int),
|
||||||
|
Sequence: new(big.Int),
|
||||||
|
},
|
||||||
|
messagesMu: new(sync.Mutex),
|
||||||
|
messages: make(map[common.Address]*message),
|
||||||
|
valSet: valSet,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// ----------------------------------------------------------------------------
|
||||||
|
|
||||||
|
type messageSet struct {
|
||||||
|
view *istanbul.View
|
||||||
|
valSet istanbul.ValidatorSet
|
||||||
|
messagesMu *sync.Mutex
|
||||||
|
messages map[common.Address]*message
|
||||||
|
}
|
||||||
|
|
||||||
|
func (ms *messageSet) View() *istanbul.View {
|
||||||
|
return ms.view
|
||||||
|
}
|
||||||
|
|
||||||
|
func (ms *messageSet) Add(msg *message) error {
|
||||||
|
ms.messagesMu.Lock()
|
||||||
|
defer ms.messagesMu.Unlock()
|
||||||
|
|
||||||
|
if err := ms.verify(msg); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
return ms.addVerifiedMessage(msg)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (ms *messageSet) Values() (result []*message) {
|
||||||
|
ms.messagesMu.Lock()
|
||||||
|
defer ms.messagesMu.Unlock()
|
||||||
|
|
||||||
|
for _, v := range ms.messages {
|
||||||
|
result = append(result, v)
|
||||||
|
}
|
||||||
|
|
||||||
|
return result
|
||||||
|
}
|
||||||
|
|
||||||
|
func (ms *messageSet) Size() int {
|
||||||
|
ms.messagesMu.Lock()
|
||||||
|
defer ms.messagesMu.Unlock()
|
||||||
|
return len(ms.messages)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (ms *messageSet) Get(addr common.Address) *message {
|
||||||
|
ms.messagesMu.Lock()
|
||||||
|
defer ms.messagesMu.Unlock()
|
||||||
|
return ms.messages[addr]
|
||||||
|
}
|
||||||
|
|
||||||
|
// ----------------------------------------------------------------------------
|
||||||
|
|
||||||
|
func (ms *messageSet) verify(msg *message) error {
|
||||||
|
// verify if the message comes from one of the validators
|
||||||
|
if _, v := ms.valSet.GetByAddress(msg.Address); v == nil {
|
||||||
|
return istanbul.ErrUnauthorizedAddress
|
||||||
|
}
|
||||||
|
|
||||||
|
// TODO: check view number and sequence number
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (ms *messageSet) addVerifiedMessage(msg *message) error {
|
||||||
|
ms.messages[msg.Address] = msg
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (ms *messageSet) String() string {
|
||||||
|
ms.messagesMu.Lock()
|
||||||
|
defer ms.messagesMu.Unlock()
|
||||||
|
addresses := make([]string, 0, len(ms.messages))
|
||||||
|
for _, v := range ms.messages {
|
||||||
|
addresses = append(addresses, v.Address.String())
|
||||||
|
}
|
||||||
|
return fmt.Sprintf("[%v]", strings.Join(addresses, ", "))
|
||||||
|
}
|
||||||
95
consensus/istanbul/core/prepare.go
Normal file
95
consensus/istanbul/core/prepare.go
Normal file
|
|
@ -0,0 +1,95 @@
|
||||||
|
// Copyright 2017 The go-ethereum Authors
|
||||||
|
// This file is part of the go-ethereum library.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is free software: you can redistribute it and/or modify
|
||||||
|
// it under the terms of the GNU Lesser General Public License as published by
|
||||||
|
// the Free Software Foundation, either version 3 of the License, or
|
||||||
|
// (at your option) any later version.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is distributed in the hope that it will be useful,
|
||||||
|
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||||
|
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||||
|
// GNU Lesser General Public License for more details.
|
||||||
|
//
|
||||||
|
// You should have received a copy of the GNU Lesser General Public License
|
||||||
|
// along with the go-ethereum library. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
|
package core
|
||||||
|
|
||||||
|
import (
|
||||||
|
"reflect"
|
||||||
|
|
||||||
|
"github.com/ethereum/go-ethereum/consensus/istanbul"
|
||||||
|
)
|
||||||
|
|
||||||
|
func (c *core) sendPrepare() {
|
||||||
|
logger := c.logger.New("state", c.state)
|
||||||
|
|
||||||
|
sub := c.current.Subject()
|
||||||
|
encodedSubject, err := Encode(sub)
|
||||||
|
if err != nil {
|
||||||
|
logger.Error("Failed to encode", "subject", sub)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
c.broadcast(&message{
|
||||||
|
Code: msgPrepare,
|
||||||
|
Msg: encodedSubject,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) handlePrepare(msg *message, src istanbul.Validator) error {
|
||||||
|
// Decode PREPARE message
|
||||||
|
var prepare *istanbul.Subject
|
||||||
|
err := msg.Decode(&prepare)
|
||||||
|
if err != nil {
|
||||||
|
return errFailedDecodePrepare
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := c.checkMessage(msgPrepare, prepare.View); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
// If it is locked, it can only process on the locked block.
|
||||||
|
// Passing verifyPrepare and checkMessage implies it is processing on the locked block since it was verified in the Preprepared state.
|
||||||
|
if err := c.verifyPrepare(prepare, src); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
c.acceptPrepare(msg, src)
|
||||||
|
|
||||||
|
// Change to Prepared state if we've received enough PREPARE messages or it is locked
|
||||||
|
// and we are in earlier state before Prepared state.
|
||||||
|
if ((c.current.IsHashLocked() && prepare.Digest == c.current.GetLockedHash()) || c.current.GetPrepareOrCommitSize() > 2*c.valSet.F()) &&
|
||||||
|
c.state.Cmp(StatePrepared) < 0 {
|
||||||
|
c.current.LockHash()
|
||||||
|
c.setState(StatePrepared)
|
||||||
|
c.sendCommit()
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// verifyPrepare verifies if the received PREPARE message is equivalent to our subject
|
||||||
|
func (c *core) verifyPrepare(prepare *istanbul.Subject, src istanbul.Validator) error {
|
||||||
|
logger := c.logger.New("from", src, "state", c.state)
|
||||||
|
|
||||||
|
sub := c.current.Subject()
|
||||||
|
if !reflect.DeepEqual(prepare, sub) {
|
||||||
|
logger.Warn("Inconsistent subjects between PREPARE and proposal", "expected", sub, "got", prepare)
|
||||||
|
return errInconsistentSubject
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) acceptPrepare(msg *message, src istanbul.Validator) error {
|
||||||
|
logger := c.logger.New("from", src, "state", c.state)
|
||||||
|
|
||||||
|
// Add the PREPARE message to current round state
|
||||||
|
if err := c.current.Prepares.Add(msg); err != nil {
|
||||||
|
logger.Error("Failed to add PREPARE message to round state", "msg", msg, "err", err)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
130
consensus/istanbul/core/preprepare.go
Normal file
130
consensus/istanbul/core/preprepare.go
Normal file
|
|
@ -0,0 +1,130 @@
|
||||||
|
// Copyright 2017 The go-ethereum Authors
|
||||||
|
// This file is part of the go-ethereum library.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is free software: you can redistribute it and/or modify
|
||||||
|
// it under the terms of the GNU Lesser General Public License as published by
|
||||||
|
// the Free Software Foundation, either version 3 of the License, or
|
||||||
|
// (at your option) any later version.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is distributed in the hope that it will be useful,
|
||||||
|
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||||
|
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||||
|
// GNU Lesser General Public License for more details.
|
||||||
|
//
|
||||||
|
// You should have received a copy of the GNU Lesser General Public License
|
||||||
|
// along with the go-ethereum library. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
|
package core
|
||||||
|
|
||||||
|
import (
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/ethereum/go-ethereum/consensus"
|
||||||
|
"github.com/ethereum/go-ethereum/consensus/istanbul"
|
||||||
|
)
|
||||||
|
|
||||||
|
func (c *core) sendPreprepare(request *istanbul.Request) {
|
||||||
|
logger := c.logger.New("state", c.state)
|
||||||
|
|
||||||
|
// If I'm the proposer and I have the same sequence with the proposal
|
||||||
|
if c.current.Sequence().Cmp(request.Proposal.Number()) == 0 && c.isProposer() {
|
||||||
|
curView := c.currentView()
|
||||||
|
preprepare, err := Encode(&istanbul.Preprepare{
|
||||||
|
View: curView,
|
||||||
|
Proposal: request.Proposal,
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
logger.Error("Failed to encode", "view", curView)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
c.broadcast(&message{
|
||||||
|
Code: msgPreprepare,
|
||||||
|
Msg: preprepare,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) handlePreprepare(msg *message, src istanbul.Validator) error {
|
||||||
|
logger := c.logger.New("from", src, "state", c.state)
|
||||||
|
|
||||||
|
// Decode PRE-PREPARE
|
||||||
|
var preprepare *istanbul.Preprepare
|
||||||
|
err := msg.Decode(&preprepare)
|
||||||
|
if err != nil {
|
||||||
|
return errFailedDecodePreprepare
|
||||||
|
}
|
||||||
|
|
||||||
|
// Ensure we have the same view with the PRE-PREPARE message
|
||||||
|
// If it is old message, see if we need to broadcast COMMIT
|
||||||
|
if err := c.checkMessage(msgPreprepare, preprepare.View); err != nil {
|
||||||
|
if err == errOldMessage {
|
||||||
|
// Get validator set for the given proposal
|
||||||
|
valSet := c.backend.ParentValidators(preprepare.Proposal).Copy()
|
||||||
|
previousProposer := c.backend.GetProposer(preprepare.Proposal.Number().Uint64() - 1)
|
||||||
|
valSet.CalcProposer(previousProposer, preprepare.View.Round.Uint64())
|
||||||
|
// Broadcast COMMIT if it is an existing block
|
||||||
|
// 1. The proposer needs to be a proposer matches the given (Sequence + Round)
|
||||||
|
// 2. The given block must exist
|
||||||
|
if valSet.IsProposer(src.Address()) && c.backend.HasPropsal(preprepare.Proposal.Hash(), preprepare.Proposal.Number()) {
|
||||||
|
c.sendCommitForOldBlock(preprepare.View, preprepare.Proposal.Hash())
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
// Check if the message comes from current proposer
|
||||||
|
if !c.valSet.IsProposer(src.Address()) {
|
||||||
|
logger.Warn("Ignore preprepare messages from non-proposer")
|
||||||
|
return errNotFromProposer
|
||||||
|
}
|
||||||
|
|
||||||
|
// Verify the proposal we received
|
||||||
|
if duration, err := c.backend.Verify(preprepare.Proposal); err != nil {
|
||||||
|
logger.Warn("Failed to verify proposal", "err", err, "duration", duration)
|
||||||
|
// if it's a future block, we will handle it again after the duration
|
||||||
|
if err == consensus.ErrFutureBlock {
|
||||||
|
c.stopFuturePreprepareTimer()
|
||||||
|
c.futurePreprepareTimer = time.AfterFunc(duration, func() {
|
||||||
|
c.sendEvent(backlogEvent{
|
||||||
|
src: src,
|
||||||
|
msg: msg,
|
||||||
|
})
|
||||||
|
})
|
||||||
|
} else {
|
||||||
|
c.sendNextRoundChange()
|
||||||
|
}
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
// Here is about to accept the PRE-PREPARE
|
||||||
|
if c.state == StateAcceptRequest {
|
||||||
|
// Send ROUND CHANGE if the locked proposal and the received proposal are different
|
||||||
|
if c.current.IsHashLocked() {
|
||||||
|
if preprepare.Proposal.Hash() == c.current.GetLockedHash() {
|
||||||
|
// Broadcast COMMIT and enters Prepared state directly
|
||||||
|
c.acceptPreprepare(preprepare)
|
||||||
|
c.setState(StatePrepared)
|
||||||
|
c.sendCommit()
|
||||||
|
} else {
|
||||||
|
// Send round change
|
||||||
|
c.sendNextRoundChange()
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
// Either
|
||||||
|
// 1. the locked proposal and the received proposal match
|
||||||
|
// 2. we have no locked proposal
|
||||||
|
c.acceptPreprepare(preprepare)
|
||||||
|
c.setState(StatePreprepared)
|
||||||
|
c.sendPrepare()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) acceptPreprepare(preprepare *istanbul.Preprepare) {
|
||||||
|
c.consensusTimestamp = time.Now()
|
||||||
|
c.current.SetPreprepare(preprepare)
|
||||||
|
}
|
||||||
99
consensus/istanbul/core/request.go
Normal file
99
consensus/istanbul/core/request.go
Normal file
|
|
@ -0,0 +1,99 @@
|
||||||
|
// Copyright 2017 The go-ethereum Authors
|
||||||
|
// This file is part of the go-ethereum library.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is free software: you can redistribute it and/or modify
|
||||||
|
// it under the terms of the GNU Lesser General Public License as published by
|
||||||
|
// the Free Software Foundation, either version 3 of the License, or
|
||||||
|
// (at your option) any later version.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is distributed in the hope that it will be useful,
|
||||||
|
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||||
|
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||||
|
// GNU Lesser General Public License for more details.
|
||||||
|
//
|
||||||
|
// You should have received a copy of the GNU Lesser General Public License
|
||||||
|
// along with the go-ethereum library. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
|
package core
|
||||||
|
|
||||||
|
import "github.com/ethereum/go-ethereum/consensus/istanbul"
|
||||||
|
|
||||||
|
func (c *core) handleRequest(request *istanbul.Request) error {
|
||||||
|
logger := c.logger.New("state", c.state, "seq", c.current.sequence)
|
||||||
|
|
||||||
|
if err := c.checkRequestMsg(request); err != nil {
|
||||||
|
if err == errInvalidMessage {
|
||||||
|
logger.Warn("invalid request")
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
logger.Warn("unexpected request", "err", err, "number", request.Proposal.Number(), "hash", request.Proposal.Hash())
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
logger.Trace("handleRequest", "number", request.Proposal.Number(), "hash", request.Proposal.Hash())
|
||||||
|
|
||||||
|
c.current.pendingRequest = request
|
||||||
|
if c.state == StateAcceptRequest {
|
||||||
|
c.sendPreprepare(request)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// check request state
|
||||||
|
// return errInvalidMessage if the message is invalid
|
||||||
|
// return errFutureMessage if the sequence of proposal is larger than current sequence
|
||||||
|
// return errOldMessage if the sequence of proposal is smaller than current sequence
|
||||||
|
func (c *core) checkRequestMsg(request *istanbul.Request) error {
|
||||||
|
if request == nil || request.Proposal == nil {
|
||||||
|
return errInvalidMessage
|
||||||
|
}
|
||||||
|
|
||||||
|
if c := c.current.sequence.Cmp(request.Proposal.Number()); c > 0 {
|
||||||
|
return errOldMessage
|
||||||
|
} else if c < 0 {
|
||||||
|
return errFutureMessage
|
||||||
|
} else {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) storeRequestMsg(request *istanbul.Request) {
|
||||||
|
logger := c.logger.New("state", c.state)
|
||||||
|
|
||||||
|
logger.Trace("Store future request", "number", request.Proposal.Number(), "hash", request.Proposal.Hash())
|
||||||
|
|
||||||
|
c.pendingRequestsMu.Lock()
|
||||||
|
defer c.pendingRequestsMu.Unlock()
|
||||||
|
|
||||||
|
c.pendingRequests.Push(request, float32(-request.Proposal.Number().Int64()))
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) processPendingRequests() {
|
||||||
|
c.pendingRequestsMu.Lock()
|
||||||
|
defer c.pendingRequestsMu.Unlock()
|
||||||
|
|
||||||
|
for !(c.pendingRequests.Empty()) {
|
||||||
|
m, prio := c.pendingRequests.Pop()
|
||||||
|
r, ok := m.(*istanbul.Request)
|
||||||
|
if !ok {
|
||||||
|
c.logger.Warn("Malformed request, skip", "msg", m)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
// Push back if it's a future message
|
||||||
|
err := c.checkRequestMsg(r)
|
||||||
|
if err != nil {
|
||||||
|
if err == errFutureMessage {
|
||||||
|
c.logger.Trace("Stop processing request", "number", r.Proposal.Number(), "hash", r.Proposal.Hash())
|
||||||
|
c.pendingRequests.Push(m, prio)
|
||||||
|
break
|
||||||
|
}
|
||||||
|
c.logger.Trace("Skip the pending request", "number", r.Proposal.Number(), "hash", r.Proposal.Hash(), "err", err)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
c.logger.Trace("Post pending request", "number", r.Proposal.Number(), "hash", r.Proposal.Hash())
|
||||||
|
|
||||||
|
go c.sendEvent(istanbul.RequestEvent{
|
||||||
|
Proposal: r.Proposal,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
172
consensus/istanbul/core/roundchange.go
Normal file
172
consensus/istanbul/core/roundchange.go
Normal file
|
|
@ -0,0 +1,172 @@
|
||||||
|
// Copyright 2017 The go-ethereum Authors
|
||||||
|
// This file is part of the go-ethereum library.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is free software: you can redistribute it and/or modify
|
||||||
|
// it under the terms of the GNU Lesser General Public License as published by
|
||||||
|
// the Free Software Foundation, either version 3 of the License, or
|
||||||
|
// (at your option) any later version.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is distributed in the hope that it will be useful,
|
||||||
|
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||||
|
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||||
|
// GNU Lesser General Public License for more details.
|
||||||
|
//
|
||||||
|
// You should have received a copy of the GNU Lesser General Public License
|
||||||
|
// along with the go-ethereum library. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
|
package core
|
||||||
|
|
||||||
|
import (
|
||||||
|
"math/big"
|
||||||
|
"sync"
|
||||||
|
|
||||||
|
"github.com/ethereum/go-ethereum/common"
|
||||||
|
"github.com/ethereum/go-ethereum/consensus/istanbul"
|
||||||
|
)
|
||||||
|
|
||||||
|
// sendNextRoundChange sends the ROUND CHANGE message with current round + 1
|
||||||
|
func (c *core) sendNextRoundChange() {
|
||||||
|
cv := c.currentView()
|
||||||
|
c.sendRoundChange(new(big.Int).Add(cv.Round, common.Big1))
|
||||||
|
}
|
||||||
|
|
||||||
|
// sendRoundChange sends the ROUND CHANGE message with the given round
|
||||||
|
func (c *core) sendRoundChange(round *big.Int) {
|
||||||
|
logger := c.logger.New("state", c.state)
|
||||||
|
|
||||||
|
cv := c.currentView()
|
||||||
|
if cv.Round.Cmp(round) >= 0 {
|
||||||
|
logger.Error("Cannot send out the round change", "current round", cv.Round, "target round", round)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
c.catchUpRound(&istanbul.View{
|
||||||
|
// The round number we'd like to transfer to.
|
||||||
|
Round: new(big.Int).Set(round),
|
||||||
|
Sequence: new(big.Int).Set(cv.Sequence),
|
||||||
|
})
|
||||||
|
|
||||||
|
// Now we have the new round number and sequence number
|
||||||
|
cv = c.currentView()
|
||||||
|
rc := &istanbul.Subject{
|
||||||
|
View: cv,
|
||||||
|
Digest: common.Hash{},
|
||||||
|
}
|
||||||
|
|
||||||
|
payload, err := Encode(rc)
|
||||||
|
if err != nil {
|
||||||
|
logger.Error("Failed to encode ROUND CHANGE", "rc", rc, "err", err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
c.broadcast(&message{
|
||||||
|
Code: msgRoundChange,
|
||||||
|
Msg: payload,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *core) handleRoundChange(msg *message, src istanbul.Validator) error {
|
||||||
|
logger := c.logger.New("state", c.state, "from", src.Address().Hex())
|
||||||
|
|
||||||
|
// Decode ROUND CHANGE message
|
||||||
|
var rc *istanbul.Subject
|
||||||
|
if err := msg.Decode(&rc); err != nil {
|
||||||
|
logger.Error("Failed to decode ROUND CHANGE", "err", err)
|
||||||
|
return errInvalidMessage
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := c.checkMessage(msgRoundChange, rc.View); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
cv := c.currentView()
|
||||||
|
roundView := rc.View
|
||||||
|
|
||||||
|
// Add the ROUND CHANGE message to its message set and return how many
|
||||||
|
// messages we've got with the same round number and sequence number.
|
||||||
|
num, err := c.roundChangeSet.Add(roundView.Round, msg)
|
||||||
|
if err != nil {
|
||||||
|
logger.Warn("Failed to add round change message", "from", src, "msg", msg, "err", err)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
// Once we received f+1 ROUND CHANGE messages, those messages form a weak certificate.
|
||||||
|
// If our round number is smaller than the certificate's round number, we would
|
||||||
|
// try to catch up the round number.
|
||||||
|
if c.waitingForRoundChange && num == c.valSet.F()+1 {
|
||||||
|
if cv.Round.Cmp(roundView.Round) < 0 {
|
||||||
|
c.sendRoundChange(roundView.Round)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
} else if num == 2*c.valSet.F()+1 && (c.waitingForRoundChange || cv.Round.Cmp(roundView.Round) < 0) {
|
||||||
|
// We've received 2f+1 ROUND CHANGE messages, start a new round immediately.
|
||||||
|
c.startNewRound(roundView.Round)
|
||||||
|
return nil
|
||||||
|
} else if cv.Round.Cmp(roundView.Round) < 0 {
|
||||||
|
// Only gossip the message with current round to other validators.
|
||||||
|
return errIgnored
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// ----------------------------------------------------------------------------
|
||||||
|
|
||||||
|
func newRoundChangeSet(valSet istanbul.ValidatorSet) *roundChangeSet {
|
||||||
|
return &roundChangeSet{
|
||||||
|
validatorSet: valSet,
|
||||||
|
roundChanges: make(map[uint64]*messageSet),
|
||||||
|
mu: new(sync.Mutex),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
type roundChangeSet struct {
|
||||||
|
validatorSet istanbul.ValidatorSet
|
||||||
|
roundChanges map[uint64]*messageSet
|
||||||
|
mu *sync.Mutex
|
||||||
|
}
|
||||||
|
|
||||||
|
// Add adds the round and message into round change set
|
||||||
|
func (rcs *roundChangeSet) Add(r *big.Int, msg *message) (int, error) {
|
||||||
|
rcs.mu.Lock()
|
||||||
|
defer rcs.mu.Unlock()
|
||||||
|
|
||||||
|
round := r.Uint64()
|
||||||
|
if rcs.roundChanges[round] == nil {
|
||||||
|
rcs.roundChanges[round] = newMessageSet(rcs.validatorSet)
|
||||||
|
}
|
||||||
|
err := rcs.roundChanges[round].Add(msg)
|
||||||
|
if err != nil {
|
||||||
|
return 0, err
|
||||||
|
}
|
||||||
|
return rcs.roundChanges[round].Size(), nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Clear deletes the messages with smaller round
|
||||||
|
func (rcs *roundChangeSet) Clear(round *big.Int) {
|
||||||
|
rcs.mu.Lock()
|
||||||
|
defer rcs.mu.Unlock()
|
||||||
|
|
||||||
|
for k, rms := range rcs.roundChanges {
|
||||||
|
if len(rms.Values()) == 0 || k < round.Uint64() {
|
||||||
|
delete(rcs.roundChanges, k)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// MaxRound returns the max round which the number of messages is equal or larger than num
|
||||||
|
func (rcs *roundChangeSet) MaxRound(num int) *big.Int {
|
||||||
|
rcs.mu.Lock()
|
||||||
|
defer rcs.mu.Unlock()
|
||||||
|
|
||||||
|
var maxRound *big.Int
|
||||||
|
for k, rms := range rcs.roundChanges {
|
||||||
|
if rms.Size() < num {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
r := big.NewInt(int64(k))
|
||||||
|
if maxRound == nil || maxRound.Cmp(r) < 0 {
|
||||||
|
maxRound = r
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return maxRound
|
||||||
|
}
|
||||||
221
consensus/istanbul/core/roundstate.go
Normal file
221
consensus/istanbul/core/roundstate.go
Normal file
|
|
@ -0,0 +1,221 @@
|
||||||
|
// Copyright 2017 The go-ethereum Authors
|
||||||
|
// This file is part of the go-ethereum library.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is free software: you can redistribute it and/or modify
|
||||||
|
// it under the terms of the GNU Lesser General Public License as published by
|
||||||
|
// the Free Software Foundation, either version 3 of the License, or
|
||||||
|
// (at your option) any later version.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is distributed in the hope that it will be useful,
|
||||||
|
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||||
|
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||||
|
// GNU Lesser General Public License for more details.
|
||||||
|
//
|
||||||
|
// You should have received a copy of the GNU Lesser General Public License
|
||||||
|
// along with the go-ethereum library. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
|
package core
|
||||||
|
|
||||||
|
import (
|
||||||
|
"io"
|
||||||
|
"math/big"
|
||||||
|
"sync"
|
||||||
|
|
||||||
|
"github.com/ethereum/go-ethereum/common"
|
||||||
|
"github.com/ethereum/go-ethereum/consensus/istanbul"
|
||||||
|
"github.com/ethereum/go-ethereum/rlp"
|
||||||
|
)
|
||||||
|
|
||||||
|
// newRoundState creates a new roundState instance with the given view and validatorSet
|
||||||
|
// lockedHash and preprepare are for round change when lock exists,
|
||||||
|
// we need to keep a reference of preprepare in order to propose locked proposal when there is a lock and itself is the proposer
|
||||||
|
func newRoundState(view *istanbul.View, validatorSet istanbul.ValidatorSet, lockedHash common.Hash, preprepare *istanbul.Preprepare, pendingRequest *istanbul.Request, hasBadProposal func(hash common.Hash) bool) *roundState {
|
||||||
|
return &roundState{
|
||||||
|
round: view.Round,
|
||||||
|
sequence: view.Sequence,
|
||||||
|
Preprepare: preprepare,
|
||||||
|
Prepares: newMessageSet(validatorSet),
|
||||||
|
Commits: newMessageSet(validatorSet),
|
||||||
|
lockedHash: lockedHash,
|
||||||
|
mu: new(sync.RWMutex),
|
||||||
|
pendingRequest: pendingRequest,
|
||||||
|
hasBadProposal: hasBadProposal,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// roundState stores the consensus state
|
||||||
|
type roundState struct {
|
||||||
|
round *big.Int
|
||||||
|
sequence *big.Int
|
||||||
|
Preprepare *istanbul.Preprepare
|
||||||
|
Prepares *messageSet
|
||||||
|
Commits *messageSet
|
||||||
|
lockedHash common.Hash
|
||||||
|
pendingRequest *istanbul.Request
|
||||||
|
|
||||||
|
mu *sync.RWMutex
|
||||||
|
hasBadProposal func(hash common.Hash) bool
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *roundState) GetPrepareOrCommitSize() int {
|
||||||
|
s.mu.RLock()
|
||||||
|
defer s.mu.RUnlock()
|
||||||
|
|
||||||
|
result := s.Prepares.Size() + s.Commits.Size()
|
||||||
|
|
||||||
|
// find duplicate one
|
||||||
|
for _, m := range s.Prepares.Values() {
|
||||||
|
if s.Commits.Get(m.Address) != nil {
|
||||||
|
result--
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return result
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *roundState) Subject() *istanbul.Subject {
|
||||||
|
s.mu.RLock()
|
||||||
|
defer s.mu.RUnlock()
|
||||||
|
|
||||||
|
if s.Preprepare == nil {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
return &istanbul.Subject{
|
||||||
|
View: &istanbul.View{
|
||||||
|
Round: new(big.Int).Set(s.round),
|
||||||
|
Sequence: new(big.Int).Set(s.sequence),
|
||||||
|
},
|
||||||
|
Digest: s.Preprepare.Proposal.Hash(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *roundState) SetPreprepare(preprepare *istanbul.Preprepare) {
|
||||||
|
s.mu.Lock()
|
||||||
|
defer s.mu.Unlock()
|
||||||
|
|
||||||
|
s.Preprepare = preprepare
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *roundState) Proposal() istanbul.Proposal {
|
||||||
|
s.mu.RLock()
|
||||||
|
defer s.mu.RUnlock()
|
||||||
|
|
||||||
|
if s.Preprepare != nil {
|
||||||
|
return s.Preprepare.Proposal
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *roundState) SetRound(r *big.Int) {
|
||||||
|
s.mu.Lock()
|
||||||
|
defer s.mu.Unlock()
|
||||||
|
|
||||||
|
s.round = new(big.Int).Set(r)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *roundState) Round() *big.Int {
|
||||||
|
s.mu.RLock()
|
||||||
|
defer s.mu.RUnlock()
|
||||||
|
|
||||||
|
return s.round
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *roundState) SetSequence(seq *big.Int) {
|
||||||
|
s.mu.Lock()
|
||||||
|
defer s.mu.Unlock()
|
||||||
|
|
||||||
|
s.sequence = seq
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *roundState) Sequence() *big.Int {
|
||||||
|
s.mu.RLock()
|
||||||
|
defer s.mu.RUnlock()
|
||||||
|
|
||||||
|
return s.sequence
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *roundState) LockHash() {
|
||||||
|
s.mu.Lock()
|
||||||
|
defer s.mu.Unlock()
|
||||||
|
|
||||||
|
if s.Preprepare != nil {
|
||||||
|
s.lockedHash = s.Preprepare.Proposal.Hash()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *roundState) UnlockHash() {
|
||||||
|
s.mu.Lock()
|
||||||
|
defer s.mu.Unlock()
|
||||||
|
|
||||||
|
s.lockedHash = common.Hash{}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *roundState) IsHashLocked() bool {
|
||||||
|
s.mu.RLock()
|
||||||
|
defer s.mu.RUnlock()
|
||||||
|
|
||||||
|
if common.EmptyHash(s.lockedHash) {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
return !s.hasBadProposal(s.GetLockedHash())
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *roundState) GetLockedHash() common.Hash {
|
||||||
|
s.mu.RLock()
|
||||||
|
defer s.mu.RUnlock()
|
||||||
|
|
||||||
|
return s.lockedHash
|
||||||
|
}
|
||||||
|
|
||||||
|
// The DecodeRLP method should read one value from the given
|
||||||
|
// Stream. It is not forbidden to read less or more, but it might
|
||||||
|
// be confusing.
|
||||||
|
func (s *roundState) DecodeRLP(stream *rlp.Stream) error {
|
||||||
|
var ss struct {
|
||||||
|
Round *big.Int
|
||||||
|
Sequence *big.Int
|
||||||
|
Preprepare *istanbul.Preprepare
|
||||||
|
Prepares *messageSet
|
||||||
|
Commits *messageSet
|
||||||
|
lockedHash common.Hash
|
||||||
|
pendingRequest *istanbul.Request
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := stream.Decode(&ss); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
s.round = ss.Round
|
||||||
|
s.sequence = ss.Sequence
|
||||||
|
s.Preprepare = ss.Preprepare
|
||||||
|
s.Prepares = ss.Prepares
|
||||||
|
s.Commits = ss.Commits
|
||||||
|
s.lockedHash = ss.lockedHash
|
||||||
|
s.pendingRequest = ss.pendingRequest
|
||||||
|
s.mu = new(sync.RWMutex)
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// EncodeRLP should write the RLP encoding of its receiver to w.
|
||||||
|
// If the implementation is a pointer method, it may also be
|
||||||
|
// called for nil pointers.
|
||||||
|
//
|
||||||
|
// Implementations should generate valid RLP. The data written is
|
||||||
|
// not verified at the moment, but a future version might. It is
|
||||||
|
// recommended to write only a single value but writing multiple
|
||||||
|
// values or no value at all is also permitted.
|
||||||
|
func (s *roundState) EncodeRLP(w io.Writer) error {
|
||||||
|
s.mu.RLock()
|
||||||
|
defer s.mu.RUnlock()
|
||||||
|
|
||||||
|
return rlp.Encode(w, []interface{}{
|
||||||
|
s.round,
|
||||||
|
s.sequence,
|
||||||
|
s.Preprepare,
|
||||||
|
s.Prepares,
|
||||||
|
s.Commits,
|
||||||
|
s.lockedHash,
|
||||||
|
s.pendingRequest,
|
||||||
|
})
|
||||||
|
}
|
||||||
164
consensus/istanbul/core/types.go
Normal file
164
consensus/istanbul/core/types.go
Normal file
|
|
@ -0,0 +1,164 @@
|
||||||
|
// Copyright 2017 The go-ethereum Authors
|
||||||
|
// This file is part of the go-ethereum library.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is free software: you can redistribute it and/or modify
|
||||||
|
// it under the terms of the GNU Lesser General Public License as published by
|
||||||
|
// the Free Software Foundation, either version 3 of the License, or
|
||||||
|
// (at your option) any later version.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is distributed in the hope that it will be useful,
|
||||||
|
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||||
|
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||||
|
// GNU Lesser General Public License for more details.
|
||||||
|
//
|
||||||
|
// You should have received a copy of the GNU Lesser General Public License
|
||||||
|
// along with the go-ethereum library. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
|
package core
|
||||||
|
|
||||||
|
import (
|
||||||
|
"fmt"
|
||||||
|
"io"
|
||||||
|
|
||||||
|
"github.com/ethereum/go-ethereum/common"
|
||||||
|
"github.com/ethereum/go-ethereum/rlp"
|
||||||
|
)
|
||||||
|
|
||||||
|
type Engine interface {
|
||||||
|
Start() error
|
||||||
|
Stop() error
|
||||||
|
}
|
||||||
|
|
||||||
|
type State uint64
|
||||||
|
|
||||||
|
const (
|
||||||
|
StateAcceptRequest State = iota
|
||||||
|
StatePreprepared
|
||||||
|
StatePrepared
|
||||||
|
StateCommitted
|
||||||
|
)
|
||||||
|
|
||||||
|
func (s State) String() string {
|
||||||
|
if s == StateAcceptRequest {
|
||||||
|
return "Accept request"
|
||||||
|
} else if s == StatePreprepared {
|
||||||
|
return "Preprepared"
|
||||||
|
} else if s == StatePrepared {
|
||||||
|
return "Prepared"
|
||||||
|
} else if s == StateCommitted {
|
||||||
|
return "Committed"
|
||||||
|
} else {
|
||||||
|
return "Unknown"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Cmp compares s and y and returns:
|
||||||
|
// -1 if s is the previous state of y
|
||||||
|
// 0 if s and y are the same state
|
||||||
|
// +1 if s is the next state of y
|
||||||
|
func (s State) Cmp(y State) int {
|
||||||
|
if uint64(s) < uint64(y) {
|
||||||
|
return -1
|
||||||
|
}
|
||||||
|
if uint64(s) > uint64(y) {
|
||||||
|
return 1
|
||||||
|
}
|
||||||
|
return 0
|
||||||
|
}
|
||||||
|
|
||||||
|
const (
|
||||||
|
msgPreprepare uint64 = iota
|
||||||
|
msgPrepare
|
||||||
|
msgCommit
|
||||||
|
msgRoundChange
|
||||||
|
msgAll
|
||||||
|
)
|
||||||
|
|
||||||
|
type message struct {
|
||||||
|
Code uint64
|
||||||
|
Msg []byte
|
||||||
|
Address common.Address
|
||||||
|
Signature []byte
|
||||||
|
CommittedSeal []byte
|
||||||
|
}
|
||||||
|
|
||||||
|
// ==============================================
|
||||||
|
//
|
||||||
|
// define the functions that needs to be provided for rlp Encoder/Decoder.
|
||||||
|
|
||||||
|
// EncodeRLP serializes m into the Ethereum RLP format.
|
||||||
|
func (m *message) EncodeRLP(w io.Writer) error {
|
||||||
|
return rlp.Encode(w, []interface{}{m.Code, m.Msg, m.Address, m.Signature, m.CommittedSeal})
|
||||||
|
}
|
||||||
|
|
||||||
|
// DecodeRLP implements rlp.Decoder, and load the consensus fields from a RLP stream.
|
||||||
|
func (m *message) DecodeRLP(s *rlp.Stream) error {
|
||||||
|
var msg struct {
|
||||||
|
Code uint64
|
||||||
|
Msg []byte
|
||||||
|
Address common.Address
|
||||||
|
Signature []byte
|
||||||
|
CommittedSeal []byte
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := s.Decode(&msg); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
m.Code, m.Msg, m.Address, m.Signature, m.CommittedSeal = msg.Code, msg.Msg, msg.Address, msg.Signature, msg.CommittedSeal
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// ==============================================
|
||||||
|
//
|
||||||
|
// define the functions that needs to be provided for core.
|
||||||
|
|
||||||
|
func (m *message) FromPayload(b []byte, validateFn func([]byte, []byte) (common.Address, error)) error {
|
||||||
|
// Decode message
|
||||||
|
err := rlp.DecodeBytes(b, &m)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
// Validate message (on a message without Signature)
|
||||||
|
if validateFn != nil {
|
||||||
|
var payload []byte
|
||||||
|
payload, err = m.PayloadNoSig()
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
_, err = validateFn(payload, m.Signature)
|
||||||
|
}
|
||||||
|
// Still return the message even the err is not nil
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
func (m *message) Payload() ([]byte, error) {
|
||||||
|
return rlp.EncodeToBytes(m)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (m *message) PayloadNoSig() ([]byte, error) {
|
||||||
|
return rlp.EncodeToBytes(&message{
|
||||||
|
Code: m.Code,
|
||||||
|
Msg: m.Msg,
|
||||||
|
Address: m.Address,
|
||||||
|
Signature: []byte{},
|
||||||
|
CommittedSeal: m.CommittedSeal,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
func (m *message) Decode(val interface{}) error {
|
||||||
|
return rlp.DecodeBytes(m.Msg, val)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (m *message) String() string {
|
||||||
|
return fmt.Sprintf("{Code: %v, Address: %v}", m.Code, m.Address.String())
|
||||||
|
}
|
||||||
|
|
||||||
|
// ==============================================
|
||||||
|
//
|
||||||
|
// helper functions
|
||||||
|
|
||||||
|
func Encode(val interface{}) ([]byte, error) {
|
||||||
|
return rlp.EncodeToBytes(val)
|
||||||
|
}
|
||||||
Loading…
Reference in a new issue