Tendermint <-> Application refactor

This commit is contained in:
Jae Kwon
2015-12-01 20:12:01 -08:00
parent a8dc417cd9
commit ef43af19ab
401 changed files with 2401 additions and 76617 deletions
+184 -93
View File
@@ -1,84 +1,167 @@
/*
Mempool receives new transactions and applies them to the latest committed state.
If the transaction is acceptable, then it broadcasts the tx to peers.
When this node happens to be the next proposer, it simply uses the recently
modified state (and the associated transactions) to construct a proposal.
*/
package mempool
import (
"bytes"
"sync"
"sync/atomic"
sm "github.com/tendermint/tendermint/state"
"github.com/tendermint/go-clist"
. "github.com/tendermint/go-common"
"github.com/tendermint/tendermint/proxy"
"github.com/tendermint/tendermint/types"
tmsp "github.com/tendermint/tmsp/types"
)
/*
The mempool pushes new txs onto the proxyAppCtx.
It gets a stream of (req, res) tuples from the proxy.
The memool stores good txs in a concurrent linked-list.
Multiple concurrent go-routines can traverse this linked-list
safely by calling .NextWait() on each element.
So we have several go-routines:
1. Consensus calling Update() and Reap() synchronously
2. Many mempool reactor's peer routines calling AppendTx()
3. Many mempool reactor's peer routines traversing the txs linked list
4. Another goroutine calling GarbageCollectTxs() periodically
To manage these goroutines, there are three methods of locking.
1. Mutations to the linked-list is protected by an internal mtx (CList is goroutine-safe)
2. Mutations to the linked-list elements are atomic
3. AppendTx() calls can be paused upon Update() and Reap(), protected by .proxyMtx
Garbage collection of old elements from mempool.txs is handlde via
the DetachPrev() call, which makes old elements not reachable by
peer broadcastTxRoutine() automatically garbage collected.
*/
type Mempool struct {
mtx sync.Mutex
state *sm.State
txs []types.Tx // TODO: we need to add a map to facilitate replace-by-fee
proxyMtx sync.Mutex
proxyAppCtx proxy.AppContext
txs *clist.CList // concurrent linked-list of good txs
counter int64 // simple incrementing counter
height int // the last block Update()'d to
expected *clist.CElement // pointer to .txs for next response
}
func NewMempool(state *sm.State) *Mempool {
return &Mempool{
state: state,
func NewMempool(proxyAppCtx proxy.AppContext) *Mempool {
mempool := &Mempool{
proxyAppCtx: proxyAppCtx,
txs: clist.New(),
counter: 0,
height: 0,
expected: nil,
}
proxyAppCtx.SetResponseCallback(mempool.resCb)
return mempool
}
func (mem *Mempool) GetState() *sm.State {
return mem.state
// Return the first element of mem.txs for peer goroutines to call .NextWait() on.
// Blocks until txs has elements.
func (mem *Mempool) TxsFrontWait() *clist.CElement {
return mem.txs.FrontWait()
}
func (mem *Mempool) GetHeight() int {
mem.mtx.Lock()
defer mem.mtx.Unlock()
return mem.state.LastBlockHeight
}
// Try a new transaction in the mempool.
// Potentially blocking if we're blocking on Update() or Reap().
func (mem *Mempool) AppendTx(tx types.Tx) (err error) {
mem.proxyMtx.Lock()
defer mem.proxyMtx.Unlock()
// Apply tx to the state and remember it.
func (mem *Mempool) AddTx(tx types.Tx) (err error) {
mem.mtx.Lock()
defer mem.mtx.Unlock()
err = sm.ExecTx(mem.state, tx, nil)
if err != nil {
log.Info("AddTx() error", "tx", tx, "error", err)
if err = mem.proxyAppCtx.Error(); err != nil {
return err
} else {
log.Info("AddTx() success", "tx", tx)
mem.txs = append(mem.txs, tx)
return nil
}
mem.proxyAppCtx.AppendTxAsync(tx)
return nil
}
// TMSP callback function
// CONTRACT: No other goroutines mutate mem.expected concurrently.
func (mem *Mempool) resCb(req tmsp.Request, res tmsp.Response) {
switch res := res.(type) {
case tmsp.ResponseAppendTx:
reqAppendTx := req.(tmsp.RequestAppendTx)
if mem.expected == nil { // Normal operation
if res.RetCode == tmsp.RetCodeOK {
mem.counter++
memTx := &mempoolTx{
counter: mem.counter,
height: int64(mem.height),
tx: reqAppendTx.TxBytes,
}
mem.txs.PushBack(memTx)
} else {
// ignore bad transaction
// TODO: handle other retcodes
}
} else { // During Update()
// TODO Log sane warning if mem.expected is nil.
memTx := mem.expected.Value.(*mempoolTx)
if !bytes.Equal(reqAppendTx.TxBytes, memTx.tx) {
PanicSanity("Unexpected tx response from proxy")
}
if res.RetCode == tmsp.RetCodeOK {
// Good, nothing to do.
} else {
// TODO: handle other retcodes
// Tx became invalidated due to newly committed block.
// NOTE: Concurrent traversal of mem.txs via CElement.Next() still works.
mem.txs.Remove(mem.expected)
mem.expected.DetachPrev()
}
mem.expected = mem.expected.Next()
}
default:
// ignore other messages
}
}
func (mem *Mempool) GetProposalTxs() []types.Tx {
mem.mtx.Lock()
defer mem.mtx.Unlock()
log.Info("GetProposalTxs:", "txs", mem.txs)
return mem.txs
// Get the valid transactions run so far, and the hash of
// the application state that results from those transactions.
func (mem *Mempool) Reap() ([]types.Tx, []byte, error) {
mem.proxyMtx.Lock()
defer mem.proxyMtx.Unlock()
// First, get the hash of txs run so far
hash, err := mem.proxyAppCtx.GetHashSync()
if err != nil {
return nil, nil, err
}
// And collect all the transactions.
txs := mem.collectTxs()
return txs, hash, nil
}
// We use this to inform peer routines of how the mempool has been updated
type ResetInfo struct {
Height int
Included []Range
Invalid []Range
func (mem *Mempool) collectTxs() []types.Tx {
txs := make([]types.Tx, 0, mem.txs.Len())
for e := mem.txs.Front(); e != nil; e = e.Next() {
memTx := e.Value.(*mempoolTx)
txs = append(txs, memTx.tx)
}
return txs
}
type Range struct {
Start int
Length int
}
// "block" is the new block being committed.
// "state" is the result of state.AppendBlock("block").
// "block" is the new block that was committed.
// Txs that are present in "block" are discarded from mempool.
// Txs that have become invalid in the new "state" are also discarded.
func (mem *Mempool) ResetForBlockAndState(block *types.Block, state *sm.State) ResetInfo {
mem.mtx.Lock()
defer mem.mtx.Unlock()
mem.state = state.Copy()
// NOTE: this should be called *after* block is committed by consensus.
// CONTRACT: block is valid and next in sequence.
func (mem *Mempool) Update(block *types.Block) error {
mem.proxyMtx.Lock()
defer mem.proxyMtx.Unlock()
// Rollback mempool synchronously
// TODO: test that proxyAppCtx's state matches the block's
err := mem.proxyAppCtx.RollbackSync()
if err != nil {
return err
}
// First, create a lookup map of txns in new block.
blockTxsMap := make(map[string]struct{})
@@ -86,50 +169,58 @@ func (mem *Mempool) ResetForBlockAndState(block *types.Block, state *sm.State) R
blockTxsMap[string(tx)] = struct{}{}
}
// Now we filter all txs from mem.txs that are in blockTxsMap,
// and ExecTx on what remains. Only valid txs are kept.
// We track the ranges of txs included in the block and invalidated by it
// so we can tell peer routines
var ri = ResetInfo{Height: block.Height}
var validTxs []types.Tx
includedStart, invalidStart := -1, -1
for i, tx := range mem.txs {
if _, ok := blockTxsMap[string(tx)]; ok {
startRange(&includedStart, i) // start counting included txs
endRange(&invalidStart, i, &ri.Invalid) // stop counting invalid txs
log.Info("Filter out, already committed", "tx", tx)
} else {
endRange(&includedStart, i, &ri.Included) // stop counting included txs
err := sm.ExecTx(mem.state, tx, nil)
if err != nil {
startRange(&invalidStart, i) // start counting invalid txs
log.Info("Filter out, no longer valid", "tx", tx, "error", err)
} else {
endRange(&invalidStart, i, &ri.Invalid) // stop counting invalid txs
log.Info("Filter in, new, valid", "tx", tx)
validTxs = append(validTxs, tx)
}
// Remove transactions that are already in block.
// Return the remaining potentially good txs.
goodTxs := mem.filterTxs(block.Height, blockTxsMap)
// Set height and expected
mem.height = block.Height
mem.expected = mem.txs.Front()
// Push good txs to proxyAppCtx
// NOTE: resCb() may be called concurrently.
for _, tx := range goodTxs {
mem.proxyAppCtx.AppendTxAsync(tx)
if err := mem.proxyAppCtx.Error(); err != nil {
return err
}
}
endRange(&includedStart, len(mem.txs)-1, &ri.Included) // stop counting included txs
endRange(&invalidStart, len(mem.txs)-1, &ri.Invalid) // stop counting invalid txs
// We're done!
log.Info("New txs", "txs", validTxs, "oldTxs", mem.txs)
mem.txs = validTxs
return ri
// NOTE: Even though we return immediately without e.g.
// calling mem.proxyAppCtx.FlushSync(),
// New mempool txs will still have to wait until
// all goodTxs are re-processed.
// So we could make synchronous calls here to proxyAppCtx.
return nil
}
func startRange(start *int, i int) {
if *start < 0 {
*start = i
func (mem *Mempool) filterTxs(height int, blockTxsMap map[string]struct{}) []types.Tx {
goodTxs := make([]types.Tx, 0, mem.txs.Len())
for e := mem.txs.Front(); e != nil; e = e.Next() {
memTx := e.Value.(*mempoolTx)
if _, ok := blockTxsMap[string(memTx.tx)]; ok {
// Remove the tx since already in block.
mem.txs.Remove(e)
e.DetachPrev()
continue
}
// Good tx!
atomic.StoreInt64(&memTx.height, int64(height))
goodTxs = append(goodTxs, memTx.tx)
}
return goodTxs
}
func endRange(start *int, i int, ranger *[]Range) {
if *start >= 0 {
length := i - *start
*ranger = append(*ranger, Range{*start, length})
*start = -1
}
//--------------------------------------------------------------------------------
// A transaction that successfully ran
type mempoolTx struct {
counter int64 // a simple incrementing counter
height int64 // height that this tx had been validated in
tx types.Tx //
}
func (memTx *mempoolTx) Height() int {
return int(atomic.LoadInt64(&memTx.height))
}
+118
View File
@@ -0,0 +1,118 @@
package mempool
import (
"encoding/binary"
"testing"
"github.com/tendermint/tendermint/proxy"
"github.com/tendermint/tendermint/types"
"github.com/tendermint/tmsp/example"
tmsp "github.com/tendermint/tmsp/types"
)
func TestSerialReap(t *testing.T) {
app := example.NewCounterApplication()
appCtxMempool := app.Open()
appCtxMempool.SetOption("serial", "on")
proxyAppCtx := proxy.NewLocalAppContext(appCtxMempool)
mempool := NewMempool(proxyAppCtx)
// Create another AppContext for committing.
appCtxConsensus := app.Open()
appCtxConsensus.SetOption("serial", "on")
appendTxsRange := func(start, end int) {
// Append some txs.
for i := start; i < end; i++ {
// This will succeed
txBytes := make([]byte, 32)
_ = binary.PutVarint(txBytes, int64(i))
err := mempool.AppendTx(txBytes)
if err != nil {
t.Fatal("Error after AppendTx: %v", err)
}
// This will fail because not serial (incrementing)
// However, error should still be nil.
// It just won't show up on Reap().
err = mempool.AppendTx(txBytes)
if err != nil {
t.Fatal("Error after AppendTx: %v", err)
}
}
}
reapCheck := func(exp int) {
txs, _, err := mempool.Reap()
if err != nil {
t.Error("Error in mempool.Reap()", err)
}
if len(txs) != exp {
t.Fatalf("Expected to reap %v txs but got %v", exp, len(txs))
}
}
updateRange := func(start, end int) {
txs := make([]types.Tx, 0)
for i := start; i < end; i++ {
txBytes := make([]byte, 32)
_ = binary.PutVarint(txBytes, int64(i))
txs = append(txs, txBytes)
}
blockHeader := &types.Header{Height: 0}
blockData := &types.Data{Txs: txs}
block := &types.Block{Header: blockHeader, Data: blockData}
err := mempool.Update(block)
if err != nil {
t.Error("Error in mempool.Update()", err)
}
}
commitRange := func(start, end int) {
// Append some txs.
for i := start; i < end; i++ {
txBytes := make([]byte, 32)
_ = binary.PutVarint(txBytes, int64(i))
_, retCode := appCtxConsensus.AppendTx(txBytes)
if retCode != tmsp.RetCodeOK {
t.Error("Error committing tx", retCode)
}
}
retCode := appCtxConsensus.Commit()
if retCode != tmsp.RetCodeOK {
t.Error("Error committing range", retCode)
}
}
//----------------------------------------
// Append some txs.
appendTxsRange(0, 100)
// Reap the txs.
reapCheck(100)
// Reap again. We should get the same amount
reapCheck(100)
// Append 0 to 999, we should reap 900 txs
// because 100 were already counted.
appendTxsRange(0, 1000)
// Reap the txs.
reapCheck(1000)
// Reap again. We should get the same amount
reapCheck(1000)
// Commit from the conensus AppContext
commitRange(0, 500)
updateRange(0, 500)
// We should have 500 left.
reapCheck(500)
}
+40 -119
View File
@@ -2,35 +2,30 @@ package mempool
import (
"bytes"
"errors"
"fmt"
"reflect"
"time"
"github.com/tendermint/go-clist"
. "github.com/tendermint/go-common"
"github.com/tendermint/go-p2p"
"github.com/tendermint/go-wire"
"github.com/tendermint/tendermint/events"
sm "github.com/tendermint/tendermint/state"
"github.com/tendermint/tendermint/types"
)
var (
const (
MempoolChannel = byte(0x30)
checkExecutedTxsMilliseconds = 1 // check for new mempool txs to send to peer
txsToSendPerCheck = 64 // send up to this many txs from the mempool per check
newBlockChCapacity = 100 // queue to process this many ResetInfos per peer
maxMempoolMessageSize = 1048576 // 1MB TODO make it configurable
maxMempoolMessageSize = 1048576 // 1MB TODO make it configurable
peerCatchupSleepIntervalMS = 100 // If peer is behind, sleep this amount
)
// MempoolReactor handles mempool tx broadcasting amongst peers.
type MempoolReactor struct {
p2p.BaseReactor
Mempool *Mempool
evsw events.Fireable
Mempool *Mempool // TODO: un-expose
evsw events.Fireable
}
func NewMempoolReactor(mempool *Mempool) *MempoolReactor {
@@ -53,11 +48,7 @@ func (memR *MempoolReactor) GetChannels() []*p2p.ChannelDescriptor {
// Implements Reactor
func (memR *MempoolReactor) AddPeer(peer *p2p.Peer) {
// Each peer gets a go routine on which we broadcast transactions in the same order we applied them to our state.
newBlockChan := make(chan ResetInfo, newBlockChCapacity)
peer.Data.Set(types.PeerMempoolChKey, newBlockChan)
timer := time.NewTicker(time.Millisecond * time.Duration(checkExecutedTxsMilliseconds))
go memR.broadcastTxRoutine(timer.C, newBlockChan, peer)
go memR.broadcastTxRoutine(peer)
}
// Implements Reactor
@@ -76,7 +67,7 @@ func (memR *MempoolReactor) Receive(chID byte, src *p2p.Peer, msgBytes []byte) {
switch msg := msg.(type) {
case *TxMessage:
err := memR.Mempool.AddTx(msg.Tx)
err := memR.Mempool.AppendTx(msg.Tx)
if err != nil {
// Bad, seen, or conflicting tx.
log.Info("Could not add tx", "tx", msg.Tx)
@@ -90,30 +81,9 @@ func (memR *MempoolReactor) Receive(chID byte, src *p2p.Peer, msgBytes []byte) {
}
}
// "block" is the new block being committed.
// "state" is the result of state.AppendBlock("block").
// Txs that are present in "block" are discarded from mempool.
// Txs that have become invalid in the new "state" are also discarded.
func (memR *MempoolReactor) ResetForBlockAndState(block *types.Block, state *sm.State) {
ri := memR.Mempool.ResetForBlockAndState(block, state)
for _, peer := range memR.Switch.Peers().List() {
peerMempoolChI := peer.Data.Get(types.PeerMempoolChKey)
if peerMempoolChI == nil {
// peer was added to switch but not yet to the memR
continue
}
peerMempoolCh := peerMempoolChI.(chan ResetInfo)
select {
case peerMempoolCh <- ri:
default:
memR.Switch.StopPeerForError(peer, errors.New("Peer's mempool push channel full"))
}
}
}
// Just an alias for AddTx since broadcasting happens in peer routines
// Just an alias for AppendTx since broadcasting happens in peer routines
func (memR *MempoolReactor) BroadcastTx(tx types.Tx) error {
return memR.Mempool.AddTx(tx)
return memR.Mempool.AppendTx(tx)
}
type PeerState interface {
@@ -126,91 +96,42 @@ type Peer interface {
Get(string) interface{}
}
// send new mempool txs to peer, strictly in order we applied them to our state.
// new blocks take chunks out of the mempool, but we've already sent some txs to the peer.
// so we wait to hear that the peer has progressed to the new height, and then continue sending txs from where we left off
func (memR *MempoolReactor) broadcastTxRoutine(tickerChan <-chan time.Time, newBlockChan chan ResetInfo, peer Peer) {
var height = memR.Mempool.GetHeight()
var txsSent int // new txs sent for height. (reset every new height)
// Send new mempool txs to peer.
// TODO: Handle mempool or reactor shutdown?
// As is this routine may block forever if no new txs come in.
func (memR *MempoolReactor) broadcastTxRoutine(peer Peer) {
var next *clist.CElement
for {
select {
case <-tickerChan:
if !peer.IsRunning() {
return
}
// make sure the peer is up to date
if peerState_i := peer.Get(types.PeerStateKey); peerState_i != nil {
peerState := peerState_i.(PeerState)
if peerState.GetHeight() < height {
continue
}
} else {
if !memR.IsRunning() {
return // Quit!
}
if next == nil {
// This happens because the CElement we were looking at got
// garbage collected (removed). That is, .NextWait() returned nil.
// Go ahead and start from the beginning.
next = memR.Mempool.TxsFrontWait() // Wait until a tx is available
}
memTx := next.Value.(*mempoolTx)
// make sure the peer is up to date
height := memTx.Height()
if peerState_i := peer.Get(types.PeerStateKey); peerState_i != nil {
peerState := peerState_i.(PeerState)
if peerState.GetHeight() < height-1 { // Allow for a lag of 1 block
time.Sleep(peerCatchupSleepIntervalMS * time.Millisecond)
continue
}
// check the mempool for new transactions
newTxs := memR.getNewTxs(height)
txsSentLoop := 0
start := time.Now()
TX_LOOP:
for i := txsSent; i < len(newTxs) && txsSentLoop < txsToSendPerCheck; i++ {
tx := newTxs[i]
msg := &TxMessage{Tx: tx}
success := peer.Send(MempoolChannel, msg)
if !success {
break TX_LOOP
} else {
txsSentLoop += 1
}
}
if txsSentLoop > 0 {
txsSent += txsSentLoop
log.Info("Sent txs to peer", "txsSentLoop", txsSentLoop,
"took", time.Since(start), "txsSent", txsSent, "newTxs", len(newTxs))
}
case ri := <-newBlockChan:
height = ri.Height
// find out how many txs below what we've sent were included in a block and how many became invalid
included := tallyRangesUpTo(ri.Included, txsSent)
invalidated := tallyRangesUpTo(ri.Invalid, txsSent)
txsSent -= included + invalidated
}
}
}
// fetch new txs from the mempool
func (memR *MempoolReactor) getNewTxs(height int) (txs []types.Tx) {
memR.Mempool.mtx.Lock()
defer memR.Mempool.mtx.Unlock()
// if the mempool got ahead of us just return empty txs
if memR.Mempool.state.LastBlockHeight != height {
return
}
return memR.Mempool.txs
}
// return the size of ranges less than upTo
func tallyRangesUpTo(ranger []Range, upTo int) int {
totalUpTo := 0
for _, r := range ranger {
if r.Start >= upTo {
break
// send memTx
msg := &TxMessage{Tx: memTx.tx}
success := peer.Send(MempoolChannel, msg)
if !success {
time.Sleep(peerCatchupSleepIntervalMS * time.Millisecond)
continue
}
if r.Start+r.Length >= upTo {
totalUpTo += upTo - r.Start
break
}
totalUpTo += r.Length
next = next.NextWait()
continue
}
return totalUpTo
}
// implements events.Eventable