tendermint/mempool/reactor.go

154 lines
3.2 KiB
Go
Raw Normal View History

2014-09-14 15:37:32 -07:00
package mempool
import (
"bytes"
"fmt"
"reflect"
2014-09-14 15:37:32 -07:00
"sync/atomic"
2015-04-01 17:30:16 -07:00
"github.com/tendermint/tendermint/binary"
. "github.com/tendermint/tendermint/common"
"github.com/tendermint/tendermint/events"
2015-04-01 17:30:16 -07:00
"github.com/tendermint/tendermint/p2p"
"github.com/tendermint/tendermint/types"
2014-09-14 15:37:32 -07:00
)
var (
MempoolChannel = byte(0x30)
2014-09-14 15:37:32 -07:00
)
// MempoolReactor handles mempool tx broadcasting amongst peers.
type MempoolReactor struct {
sw *p2p.Switch
2014-09-14 15:37:32 -07:00
quit chan struct{}
started uint32
stopped uint32
Mempool *Mempool
2015-04-15 23:40:27 -07:00
evsw events.Fireable
2014-09-14 15:37:32 -07:00
}
2014-10-22 17:20:44 -07:00
func NewMempoolReactor(mempool *Mempool) *MempoolReactor {
2014-09-14 15:37:32 -07:00
memR := &MempoolReactor{
quit: make(chan struct{}),
Mempool: mempool,
2014-09-14 15:37:32 -07:00
}
return memR
}
2014-10-22 17:20:44 -07:00
// Implements Reactor
func (memR *MempoolReactor) Start(sw *p2p.Switch) {
2014-09-14 15:37:32 -07:00
if atomic.CompareAndSwapUint32(&memR.started, 0, 1) {
2014-10-22 17:20:44 -07:00
memR.sw = sw
2015-07-19 14:49:13 -07:00
log.Notice("Starting MempoolReactor")
2014-09-14 15:37:32 -07:00
}
}
2014-10-22 17:20:44 -07:00
// Implements Reactor
2014-09-14 15:37:32 -07:00
func (memR *MempoolReactor) Stop() {
if atomic.CompareAndSwapUint32(&memR.stopped, 0, 1) {
2015-07-19 14:49:13 -07:00
log.Notice("Stopping MempoolReactor")
2014-09-14 15:37:32 -07:00
close(memR.quit)
}
}
2014-10-22 17:20:44 -07:00
// Implements Reactor
func (memR *MempoolReactor) GetChannels() []*p2p.ChannelDescriptor {
return []*p2p.ChannelDescriptor{
&p2p.ChannelDescriptor{
Id: MempoolChannel,
2014-10-22 17:20:44 -07:00
Priority: 5,
},
2014-09-14 15:37:32 -07:00
}
}
// Implements Reactor
func (pexR *MempoolReactor) AddPeer(peer *p2p.Peer) {
2014-09-14 15:37:32 -07:00
}
// Implements Reactor
2014-10-22 17:20:44 -07:00
func (pexR *MempoolReactor) RemovePeer(peer *p2p.Peer, reason interface{}) {
2014-09-14 15:37:32 -07:00
}
2014-10-22 17:20:44 -07:00
// Implements Reactor
func (memR *MempoolReactor) Receive(chId byte, src *p2p.Peer, msgBytes []byte) {
2015-07-13 16:00:01 -07:00
_, msg, err := DecodeMessage(msgBytes)
if err != nil {
2014-12-29 18:39:19 -08:00
log.Warn("Error decoding message", "error", err)
return
}
2015-07-19 14:49:13 -07:00
log.Notice("MempoolReactor received message", "msg", msg)
2014-09-14 15:37:32 -07:00
2015-07-13 16:00:01 -07:00
switch msg := msg.(type) {
2014-09-14 15:37:32 -07:00
case *TxMessage:
err := memR.Mempool.AddTx(msg.Tx)
2014-09-14 15:37:32 -07:00
if err != nil {
// Bad, seen, or conflicting tx.
2015-07-19 14:49:13 -07:00
log.Info("Could not add tx", "tx", msg.Tx)
2014-09-14 15:37:32 -07:00
return
} else {
2015-07-19 14:49:13 -07:00
log.Info("Added valid tx", "tx", msg.Tx)
2014-09-14 15:37:32 -07:00
}
// Share tx.
// We use a simple shotgun approach for now.
// TODO: improve efficiency
for _, peer := range memR.sw.Peers().List() {
if peer.Key == src.Key {
continue
}
peer.TrySend(MempoolChannel, msg)
2014-09-14 15:37:32 -07:00
}
default:
log.Warn(Fmt("Unknown message type %v", reflect.TypeOf(msg)))
2014-09-14 15:37:32 -07:00
}
}
func (memR *MempoolReactor) BroadcastTx(tx types.Tx) error {
err := memR.Mempool.AddTx(tx)
2014-10-22 17:20:44 -07:00
if err != nil {
return err
}
msg := &TxMessage{Tx: tx}
memR.sw.Broadcast(MempoolChannel, msg)
2014-10-22 17:20:44 -07:00
return nil
}
// implements events.Eventable
2015-04-15 23:40:27 -07:00
func (memR *MempoolReactor) SetFireable(evsw events.Fireable) {
memR.evsw = evsw
}
2014-09-14 15:37:32 -07:00
//-----------------------------------------------------------------------------
// Messages
const (
2015-04-14 15:57:16 -07:00
msgTypeTx = byte(0x01)
2014-09-14 15:37:32 -07:00
)
2015-04-14 15:57:16 -07:00
type MempoolMessage interface{}
var _ = binary.RegisterInterface(
struct{ MempoolMessage }{},
binary.ConcreteType{&TxMessage{}, msgTypeTx},
2015-04-14 15:57:16 -07:00
)
func DecodeMessage(bz []byte) (msgType byte, msg MempoolMessage, err error) {
2014-09-14 15:37:32 -07:00
msgType = bz[0]
2015-04-14 15:57:16 -07:00
n := new(int64)
r := bytes.NewReader(bz)
msg = binary.ReadBinary(struct{ MempoolMessage }{}, r, n, &err).(struct{ MempoolMessage }).MempoolMessage
2014-09-14 15:37:32 -07:00
return
}
//-------------------------------------
type TxMessage struct {
Tx types.Tx
2014-09-14 15:37:32 -07:00
}
func (m *TxMessage) String() string {
return fmt.Sprintf("[TxMessage %v]", m.Tx)
}