tendermint/evidence/reactor.go

169 lines
4.8 KiB
Go
Raw Normal View History

2017-11-19 20:22:25 -08:00
package evidence
2017-11-02 11:06:48 -07:00
import (
"fmt"
"reflect"
2017-11-02 17:26:07 -07:00
"time"
2017-11-02 11:06:48 -07:00
2018-04-09 05:14:33 -07:00
"github.com/tendermint/go-amino"
"github.com/tendermint/tmlibs/log"
2017-11-02 11:06:48 -07:00
"github.com/tendermint/tendermint/p2p"
"github.com/tendermint/tendermint/types"
)
const (
EvidenceChannel = byte(0x38)
2017-11-02 11:06:48 -07:00
2018-04-09 05:14:33 -07:00
maxMsgSize = 1048576 // 1MB TODO make it configurable
2017-11-02 17:26:07 -07:00
broadcastEvidenceIntervalS = 60 // broadcast uncommitted evidence this often
2017-11-02 11:06:48 -07:00
)
// EvidenceReactor handles evpool evidence broadcasting amongst peers.
type EvidenceReactor struct {
2017-11-02 11:06:48 -07:00
p2p.BaseReactor
evpool *EvidencePool
eventBus *types.EventBus
2017-11-02 11:06:48 -07:00
}
// NewEvidenceReactor returns a new EvidenceReactor with the given config and evpool.
2017-12-26 17:34:57 -08:00
func NewEvidenceReactor(evpool *EvidencePool) *EvidenceReactor {
evR := &EvidenceReactor{
2017-11-02 17:26:07 -07:00
evpool: evpool,
2017-11-02 11:06:48 -07:00
}
evR.BaseReactor = *p2p.NewBaseReactor("EvidenceReactor", evR)
2017-11-02 11:06:48 -07:00
return evR
}
// SetLogger sets the Logger on the reactor and the underlying Evidence.
func (evR *EvidenceReactor) SetLogger(l log.Logger) {
2017-11-02 11:06:48 -07:00
evR.Logger = l
2017-11-02 17:26:07 -07:00
evR.evpool.SetLogger(l)
}
// OnStart implements cmn.Service
func (evR *EvidenceReactor) OnStart() error {
2017-11-02 17:26:07 -07:00
if err := evR.BaseReactor.OnStart(); err != nil {
return err
}
go evR.broadcastRoutine()
return nil
2017-11-02 11:06:48 -07:00
}
// GetChannels implements Reactor.
// It returns the list of channels for this reactor.
func (evR *EvidenceReactor) GetChannels() []*p2p.ChannelDescriptor {
2017-11-02 11:06:48 -07:00
return []*p2p.ChannelDescriptor{
&p2p.ChannelDescriptor{
ID: EvidenceChannel,
2017-11-02 11:06:48 -07:00
Priority: 5,
},
}
}
// AddPeer implements Reactor.
func (evR *EvidenceReactor) AddPeer(peer p2p.Peer) {
// send the peer our high-priority evidence.
// the rest will be sent by the broadcastRoutine
2017-12-26 22:27:03 -08:00
evidences := evR.evpool.PriorityEvidence()
msg := &EvidenceListMessage{evidences}
2018-04-05 05:43:23 -07:00
success := peer.Send(EvidenceChannel, cdc.MustMarshalBinaryBare(msg))
2017-11-02 11:06:48 -07:00
if !success {
// TODO: remove peer ?
}
}
// RemovePeer implements Reactor.
func (evR *EvidenceReactor) RemovePeer(peer p2p.Peer, reason interface{}) {
// nothing to do
2017-11-02 11:06:48 -07:00
}
// Receive implements Reactor.
// It adds any received evidence to the evpool.
func (evR *EvidenceReactor) Receive(chID byte, src p2p.Peer, msgBytes []byte) {
2018-04-05 05:43:23 -07:00
msg, err := DecodeMessage(msgBytes)
2017-11-02 11:06:48 -07:00
if err != nil {
2018-03-04 01:42:45 -08:00
evR.Logger.Error("Error decoding message", "src", src, "chId", chID, "msg", msg, "err", err, "bytes", msgBytes)
evR.Switch.StopPeerForError(src, err)
2017-11-02 11:06:48 -07:00
return
}
evR.Logger.Debug("Receive", "src", src, "chId", chID, "msg", msg)
switch msg := msg.(type) {
case *EvidenceListMessage:
2017-11-02 11:06:48 -07:00
for _, ev := range msg.Evidence {
2017-11-02 17:26:07 -07:00
err := evR.evpool.AddEvidence(ev)
2017-11-02 11:06:48 -07:00
if err != nil {
evR.Logger.Info("Evidence is not valid", "evidence", msg.Evidence, "err", err)
2018-03-04 02:22:58 -08:00
// punish peer
evR.Switch.StopPeerForError(src, err)
2017-11-02 11:06:48 -07:00
}
}
default:
evR.Logger.Error(fmt.Sprintf("Unknown message type %v", reflect.TypeOf(msg)))
}
}
// SetEventSwitch implements events.Eventable.
func (evR *EvidenceReactor) SetEventBus(b *types.EventBus) {
evR.eventBus = b
2017-11-02 11:06:48 -07:00
}
2017-12-26 22:27:03 -08:00
// Broadcast new evidence to all peers.
// Broadcasts must be non-blocking so routine is always available to read off EvidenceChan.
func (evR *EvidenceReactor) broadcastRoutine() {
2017-11-02 17:26:07 -07:00
ticker := time.NewTicker(time.Second * broadcastEvidenceIntervalS)
for {
select {
2017-11-18 16:57:55 -08:00
case evidence := <-evR.evpool.EvidenceChan():
2017-11-02 17:26:07 -07:00
// broadcast some new evidence
2017-11-19 20:32:53 -08:00
msg := &EvidenceListMessage{[]types.Evidence{evidence}}
2018-04-05 05:43:23 -07:00
evR.Switch.Broadcast(EvidenceChannel, cdc.MustMarshalBinaryBare(msg))
2017-11-02 17:26:07 -07:00
// TODO: Broadcast runs asynchronously, so this should wait on the successChan
2017-11-02 17:26:07 -07:00
// in another routine before marking to be proper.
evR.evpool.evidenceStore.MarkEvidenceAsBroadcasted(evidence)
2017-11-02 17:26:07 -07:00
case <-ticker.C:
// broadcast all pending evidence
2017-11-19 20:32:53 -08:00
msg := &EvidenceListMessage{evR.evpool.PendingEvidence()}
2018-04-05 05:43:23 -07:00
evR.Switch.Broadcast(EvidenceChannel, cdc.MustMarshalBinaryBare(msg))
case <-evR.Quit():
2017-11-02 17:26:07 -07:00
return
}
}
}
2017-11-02 11:06:48 -07:00
//-----------------------------------------------------------------------------
// Messages
// EvidenceMessage is a message sent or received by the EvidenceReactor.
type EvidenceMessage interface{}
2017-11-02 11:06:48 -07:00
2018-04-05 05:43:23 -07:00
func RegisterEvidenceMessages(cdc *amino.Codec) {
cdc.RegisterInterface((*EvidenceMessage)(nil), nil)
cdc.RegisterConcrete(&EvidenceListMessage{},
2018-04-06 13:46:40 -07:00
"tendermint/evidence/EvidenceListMessage", nil)
2018-04-05 05:43:23 -07:00
}
2017-11-02 11:06:48 -07:00
// DecodeMessage decodes a byte-array into a EvidenceMessage.
2018-04-05 05:43:23 -07:00
func DecodeMessage(bz []byte) (msg EvidenceMessage, err error) {
2018-04-09 05:14:33 -07:00
if len(bz) > maxMsgSize {
return msg, fmt.Errorf("Msg exceeds max size (%d > %d)",
len(bz), maxMsgSize)
}
2018-04-05 05:43:23 -07:00
err = cdc.UnmarshalBinaryBare(bz, &msg)
2017-11-02 11:06:48 -07:00
return
}
//-------------------------------------
// EvidenceMessage contains a list of evidence.
type EvidenceListMessage struct {
2017-11-02 17:26:07 -07:00
Evidence []types.Evidence
2017-11-02 11:06:48 -07:00
}
// String returns a string representation of the EvidenceListMessage.
func (m *EvidenceListMessage) String() string {
return fmt.Sprintf("[EvidenceListMessage %v]", m.Evidence)
2017-11-02 11:06:48 -07:00
}