You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

191 lines
5.0 KiB

10 years ago
10 years ago
8 years ago
10 years ago
10 years ago
10 years ago
10 years ago
10 years ago
8 years ago
10 years ago
10 years ago
7 years ago
10 years ago
7 years ago
10 years ago
7 years ago
8 years ago
10 years ago
10 years ago
10 years ago
10 years ago
8 years ago
8 years ago
10 years ago
10 years ago
10 years ago
10 years ago
10 years ago
10 years ago
10 years ago
10 years ago
10 years ago
10 years ago
10 years ago
10 years ago
10 years ago
  1. package mempool
  2. import (
  3. "bytes"
  4. "fmt"
  5. "reflect"
  6. "time"
  7. abci "github.com/tendermint/abci/types"
  8. wire "github.com/tendermint/go-wire"
  9. "github.com/tendermint/tmlibs/clist"
  10. "github.com/tendermint/tmlibs/log"
  11. cfg "github.com/tendermint/tendermint/config"
  12. "github.com/tendermint/tendermint/p2p"
  13. "github.com/tendermint/tendermint/types"
  14. )
  15. const (
  16. MempoolChannel = byte(0x30)
  17. maxMempoolMessageSize = 1048576 // 1MB TODO make it configurable
  18. peerCatchupSleepIntervalMS = 100 // If peer is behind, sleep this amount
  19. )
  20. // MempoolReactor handles mempool tx broadcasting amongst peers.
  21. type MempoolReactor struct {
  22. p2p.BaseReactor
  23. config *cfg.MempoolConfig
  24. Mempool *Mempool
  25. }
  26. // NewMempoolReactor returns a new MempoolReactor with the given config and mempool.
  27. func NewMempoolReactor(config *cfg.MempoolConfig, mempool *Mempool) *MempoolReactor {
  28. memR := &MempoolReactor{
  29. config: config,
  30. Mempool: mempool,
  31. }
  32. memR.BaseReactor = *p2p.NewBaseReactor("MempoolReactor", memR)
  33. return memR
  34. }
  35. // SetLogger sets the Logger on the reactor and the underlying Mempool.
  36. func (memR *MempoolReactor) SetLogger(l log.Logger) {
  37. memR.Logger = l
  38. memR.Mempool.SetLogger(l)
  39. }
  40. // GetChannels implements Reactor.
  41. // It returns the list of channels for this reactor.
  42. func (memR *MempoolReactor) GetChannels() []*p2p.ChannelDescriptor {
  43. return []*p2p.ChannelDescriptor{
  44. {
  45. ID: MempoolChannel,
  46. Priority: 5,
  47. },
  48. }
  49. }
  50. // AddPeer implements Reactor.
  51. // It starts a broadcast routine ensuring all txs are forwarded to the given peer.
  52. func (memR *MempoolReactor) AddPeer(peer p2p.Peer) {
  53. go memR.broadcastTxRoutine(peer)
  54. }
  55. // RemovePeer implements Reactor.
  56. func (memR *MempoolReactor) RemovePeer(peer p2p.Peer, reason interface{}) {
  57. // broadcast routine checks if peer is gone and returns
  58. }
  59. // Receive implements Reactor.
  60. // It adds any received transactions to the mempool.
  61. func (memR *MempoolReactor) Receive(chID byte, src p2p.Peer, msgBytes []byte) {
  62. _, msg, err := DecodeMessage(msgBytes)
  63. if err != nil {
  64. memR.Logger.Error("Error decoding message", "err", err)
  65. return
  66. }
  67. memR.Logger.Debug("Receive", "src", src, "chId", chID, "msg", msg)
  68. switch msg := msg.(type) {
  69. case *TxMessage:
  70. err := memR.Mempool.CheckTx(msg.Tx, nil)
  71. if err != nil {
  72. memR.Logger.Info("Could not check tx", "tx", msg.Tx, "err", err)
  73. }
  74. // broadcasting happens from go routines per peer
  75. default:
  76. memR.Logger.Error(fmt.Sprintf("Unknown message type %v", reflect.TypeOf(msg)))
  77. }
  78. }
  79. // BroadcastTx is an alias for Mempool.CheckTx. Broadcasting itself happens in peer routines.
  80. func (memR *MempoolReactor) BroadcastTx(tx types.Tx, cb func(*abci.Response)) error {
  81. return memR.Mempool.CheckTx(tx, cb)
  82. }
  83. // PeerState describes the state of a peer.
  84. type PeerState interface {
  85. GetHeight() int64
  86. }
  87. // Send new mempool txs to peer.
  88. func (memR *MempoolReactor) broadcastTxRoutine(peer p2p.Peer) {
  89. if !memR.config.Broadcast {
  90. return
  91. }
  92. var next *clist.CElement
  93. for {
  94. // This happens because the CElement we were looking at got garbage
  95. // collected (removed). That is, .NextWait() returned nil. Go ahead and
  96. // start from the beginning.
  97. if next == nil {
  98. select {
  99. case <-memR.Mempool.TxsWaitChan(): // Wait until a tx is available
  100. if next = memR.Mempool.TxsFront(); next == nil {
  101. continue
  102. }
  103. case <-peer.QuitChan():
  104. return
  105. case <-memR.Quit:
  106. return
  107. }
  108. }
  109. memTx := next.Value.(*mempoolTx)
  110. // make sure the peer is up to date
  111. height := memTx.Height()
  112. if peerState_i := peer.Get(types.PeerStateKey); peerState_i != nil {
  113. peerState := peerState_i.(PeerState)
  114. if peerState.GetHeight() < height-1 { // Allow for a lag of 1 block
  115. time.Sleep(peerCatchupSleepIntervalMS * time.Millisecond)
  116. continue
  117. }
  118. }
  119. // send memTx
  120. msg := &TxMessage{Tx: memTx.tx}
  121. success := peer.Send(MempoolChannel, struct{ MempoolMessage }{msg})
  122. if !success {
  123. time.Sleep(peerCatchupSleepIntervalMS * time.Millisecond)
  124. continue
  125. }
  126. select {
  127. case <-next.NextWaitChan():
  128. // see the start of the for loop for nil check
  129. next = next.Next()
  130. case <-peer.QuitChan():
  131. return
  132. case <-memR.Quit:
  133. return
  134. }
  135. }
  136. }
  137. //-----------------------------------------------------------------------------
  138. // Messages
  139. const (
  140. msgTypeTx = byte(0x01)
  141. )
  142. // MempoolMessage is a message sent or received by the MempoolReactor.
  143. type MempoolMessage interface{}
  144. var _ = wire.RegisterInterface(
  145. struct{ MempoolMessage }{},
  146. wire.ConcreteType{&TxMessage{}, msgTypeTx},
  147. )
  148. // DecodeMessage decodes a byte-array into a MempoolMessage.
  149. func DecodeMessage(bz []byte) (msgType byte, msg MempoolMessage, err error) {
  150. msgType = bz[0]
  151. n := new(int)
  152. r := bytes.NewReader(bz)
  153. msg = wire.ReadBinary(struct{ MempoolMessage }{}, r, maxMempoolMessageSize, n, &err).(struct{ MempoolMessage }).MempoolMessage
  154. return
  155. }
  156. //-------------------------------------
  157. // TxMessage is a MempoolMessage containing a transaction.
  158. type TxMessage struct {
  159. Tx types.Tx
  160. }
  161. // String returns a string representation of the TxMessage.
  162. func (m *TxMessage) String() string {
  163. return fmt.Sprintf("[TxMessage %v]", m.Tx)
  164. }