From e5c5db7f5bc39f36d1390bae90864c428b4215f4 Mon Sep 17 00:00:00 2001 From: Jasmina Malicevic Date: Fri, 17 Jun 2022 11:46:54 +0200 Subject: [PATCH] Updated mocks --- config/config.go | 6 + consensus/byzantine_test.go | 25 ++- consensus/common_test.go | 28 ++- consensus/reactor_test.go | 24 ++- consensus/replay_stubs.go | 5 + mempool/mock/mempool.go | 8 +- mempool/v0/clist_mempool_test.go | 14 +- mempool/v1/mempool_bench_test.go | 5 +- mempool/v1/mempool_test.go | 25 ++- mempool/v1/reactor_test.go | 231 ++++++++++++------------ mempool_bak/mock/mempool.go | 5 + node/node.go | 75 ++++++-- node/node_test.go | 59 ++++-- test/fuzz/mempool/checktx.go | 35 ---- test/maverick/consensus/replay_stubs.go | 1 + test/maverick/node/node.go | 70 +++++-- 16 files changed, 376 insertions(+), 240 deletions(-) delete mode 100644 test/fuzz/mempool/checktx.go diff --git a/config/config.go b/config/config.go index 8b1a54ae0..c5dfbb7f3 100644 --- a/config/config.go +++ b/config/config.go @@ -23,6 +23,10 @@ const ( // DefaultLogLevel defines a default log level as INFO. DefaultLogLevel = "info" + + // Mempool versions. V1 is prioritized mempool, v0 is regular mempool. + MempoolV0 = "v0" + MempoolV1 = "v1" ) // NOTE: Most of the structs & relevant comments + the @@ -676,6 +680,7 @@ func DefaultFuzzConnConfig() *FuzzConnConfig { // MempoolConfig defines the configuration options for the Tendermint mempool type MempoolConfig struct { + Version string `mapstructure:"version"` RootDir string `mapstructure:"home"` Recheck bool `mapstructure:"recheck"` Broadcast bool `mapstructure:"broadcast"` @@ -720,6 +725,7 @@ type MempoolConfig struct { // DefaultMempoolConfig returns a default configuration for the Tendermint mempool func DefaultMempoolConfig() *MempoolConfig { return &MempoolConfig{ + Version: MempoolV0, Recheck: true, Broadcast: true, WalPath: "", diff --git a/consensus/byzantine_test.go b/consensus/byzantine_test.go index d71f9003e..dd3779c8a 100644 --- a/consensus/byzantine_test.go +++ b/consensus/byzantine_test.go @@ -21,6 +21,10 @@ import ( "github.com/tendermint/tendermint/libs/service" tmsync "github.com/tendermint/tendermint/libs/sync" mempl "github.com/tendermint/tendermint/mempool" + + cfg "github.com/tendermint/tendermint/config" + mempoolv0 "github.com/tendermint/tendermint/mempool/v0" + mempoolv1 "github.com/tendermint/tendermint/mempool/v1" "github.com/tendermint/tendermint/p2p" tmproto "github.com/tendermint/tendermint/proto/tendermint/types" sm "github.com/tendermint/tendermint/state" @@ -60,11 +64,28 @@ func TestByzantinePrevoteEquivocation(t *testing.T) { // one for mempool, one for consensus mtx := new(tmsync.Mutex) - proxyAppConnMem := abcicli.NewLocalClient(mtx, app) + proxyAppConnCon := abcicli.NewLocalClient(mtx, app) // Make Mempool - mempool := mempl.NewTxMempool(log.TestingLogger().With("module", "mempool"), thisConfig.Mempool, proxyAppConnMem, 0) + var mempool mempl.Mempool + + switch thisConfig.Mempool.Version { + case cfg.MempoolV0: + mempool = mempoolv0.NewCListMempool(config.Mempool, + proxyAppConnCon, + state.LastBlockHeight, + mempoolv0.WithPreCheck(sm.TxPreCheck(state)), + mempoolv0.WithPostCheck(sm.TxPostCheck(state))) + case cfg.MempoolV1: + mempool = mempoolv1.NewTxMempool(logger, + config.Mempool, + proxyAppConnCon, + state.LastBlockHeight, + mempoolv1.WithPreCheck(sm.TxPreCheck(state)), + mempoolv1.WithPostCheck(sm.TxPostCheck(state)), + ) + } if thisConfig.Consensus.WaitForTxs() { mempool.EnableTxsAvailable() diff --git a/consensus/common_test.go b/consensus/common_test.go index 1c3660e80..5d6a905e2 100644 --- a/consensus/common_test.go +++ b/consensus/common_test.go @@ -31,6 +31,8 @@ import ( tmpubsub "github.com/tendermint/tendermint/libs/pubsub" tmsync "github.com/tendermint/tendermint/libs/sync" mempl "github.com/tendermint/tendermint/mempool" + mempoolv0 "github.com/tendermint/tendermint/mempool/v0" + mempoolv1 "github.com/tendermint/tendermint/mempool/v1" "github.com/tendermint/tendermint/p2p" "github.com/tendermint/tendermint/privval" tmproto "github.com/tendermint/tendermint/proto/tendermint/types" @@ -390,12 +392,34 @@ func newStateWithConfigAndBlockStore( // one for mempool, one for consensus mtx := new(tmsync.Mutex) - proxyAppConnMem := abcicli.NewLocalClient(mtx, app) + proxyAppConnCon := abcicli.NewLocalClient(mtx, app) // Make Mempool - mempool := mempl.NewTxMempool(log.TestingLogger().With("module", "mempool"), thisConfig.Mempool, proxyAppConnMem, 0) + memplMetrics := mempl.NopMetrics() + // Make Mempool + var mempool mempl.Mempool + + switch config.Mempool.Version { + case cfg.MempoolV0: + mempool = mempoolv0.NewCListMempool(config.Mempool, + proxyAppConnCon, + state.LastBlockHeight, + mempoolv0.WithMetrics(memplMetrics), + mempoolv0.WithPreCheck(sm.TxPreCheck(state)), + mempoolv0.WithPostCheck(sm.TxPostCheck(state))) + case cfg.MempoolV1: + logger := consensusLogger() + mempool = mempoolv1.NewTxMempool(logger, + config.Mempool, + proxyAppConnCon, + state.LastBlockHeight, + mempoolv1.WithMetrics(memplMetrics), + mempoolv1.WithPreCheck(sm.TxPreCheck(state)), + mempoolv1.WithPostCheck(sm.TxPostCheck(state)), + ) + } if thisConfig.Consensus.WaitForTxs() { mempool.EnableTxsAvailable() } diff --git a/consensus/reactor_test.go b/consensus/reactor_test.go index 0e2741021..f20158fb7 100644 --- a/consensus/reactor_test.go +++ b/consensus/reactor_test.go @@ -29,6 +29,8 @@ import ( "github.com/tendermint/tendermint/libs/log" tmsync "github.com/tendermint/tendermint/libs/sync" mempl "github.com/tendermint/tendermint/mempool" + mempoolv0 "github.com/tendermint/tendermint/mempool/v0" + mempoolv1 "github.com/tendermint/tendermint/mempool/v1" "github.com/tendermint/tendermint/p2p" p2pmock "github.com/tendermint/tendermint/p2p/mock" tmproto "github.com/tendermint/tendermint/proto/tendermint/types" @@ -154,12 +156,30 @@ func TestReactorWithEvidence(t *testing.T) { // one for mempool, one for consensus mtx := new(tmsync.Mutex) - proxyAppConnMem := abcicli.NewLocalClient(mtx, app) + memplMetrics := mempl.PrometheusMetrics("node_test_1") proxyAppConnCon := abcicli.NewLocalClient(mtx, app) // Make Mempool - mempool := mempl.NewTxMempool(log.TestingLogger().With("module", "mempool"), thisConfig.Mempool, proxyAppConnMem, 0) + var mempool mempl.Mempool + switch config.Mempool.Version { + case cfg.MempoolV0: + mempool = mempoolv0.NewCListMempool(config.Mempool, + proxyAppConnCon, + state.LastBlockHeight, + mempoolv0.WithMetrics(memplMetrics), + mempoolv0.WithPreCheck(sm.TxPreCheck(state)), + mempoolv0.WithPostCheck(sm.TxPostCheck(state))) + case cfg.MempoolV1: + mempool = mempoolv1.NewTxMempool(logger, + config.Mempool, + proxyAppConnCon, + state.LastBlockHeight, + mempoolv1.WithMetrics(memplMetrics), + mempoolv1.WithPreCheck(sm.TxPreCheck(state)), + mempoolv1.WithPostCheck(sm.TxPostCheck(state)), + ) + } if thisConfig.Consensus.WaitForTxs() { mempool.EnableTxsAvailable() } diff --git a/consensus/replay_stubs.go b/consensus/replay_stubs.go index ed5aa5626..9b4b3b062 100644 --- a/consensus/replay_stubs.go +++ b/consensus/replay_stubs.go @@ -22,6 +22,11 @@ func (emptyMempool) SizeBytes() int64 { return 0 } func (emptyMempool) CheckTx(_ types.Tx, _ func(*abci.Response), _ mempl.TxInfo) error { return nil } + +func (txmp emptyMempool) RemoveTxByKey(txKey types.TxKey) error { + return nil +} + func (emptyMempool) ReapMaxBytesMaxGas(_, _ int64) types.Txs { return types.Txs{} } func (emptyMempool) ReapMaxTxs(n int) types.Txs { return types.Txs{} } func (emptyMempool) Update( diff --git a/mempool/mock/mempool.go b/mempool/mock/mempool.go index 8e6f0c7bf..3f293381f 100644 --- a/mempool/mock/mempool.go +++ b/mempool/mock/mempool.go @@ -1,11 +1,9 @@ package mock import ( - "context" - abci "github.com/tendermint/tendermint/abci/types" - "github.com/tendermint/tendermint/internal/libs/clist" - "github.com/tendermint/tendermint/internal/mempool" + "github.com/tendermint/tendermint/libs/clist" + "github.com/tendermint/tendermint/mempool" "github.com/tendermint/tendermint/types" ) @@ -17,7 +15,7 @@ var _ mempool.Mempool = Mempool{} func (Mempool) Lock() {} func (Mempool) Unlock() {} func (Mempool) Size() int { return 0 } -func (Mempool) CheckTx(_ context.Context, _ types.Tx, _ func(*abci.Response), _ mempool.TxInfo) error { +func (Mempool) CheckTx(_ types.Tx, _ func(*abci.Response), _ mempool.TxInfo) error { return nil } func (Mempool) RemoveTxByKey(txKey types.TxKey) error { return nil } diff --git a/mempool/v0/clist_mempool_test.go b/mempool/v0/clist_mempool_test.go index b5dc54cd0..e569867ec 100644 --- a/mempool/v0/clist_mempool_test.go +++ b/mempool/v0/clist_mempool_test.go @@ -232,11 +232,11 @@ func TestMempoolUpdateDoesNotPanicWhenApplicationMissedTx(t *testing.T) { mockClient.On("FlushAsync", mock.Anything).Return(abciclient.NewReqRes(abci.ToRequestFlush()), nil) mockClient.On("SetResponseCallback", mock.MatchedBy(func(cb abciclient.Callback) bool { callback = cb; return true })) - cc := func() (abciclient.Client, error) { - return mockClient, nil - } - - mp, cleanup, err := newMempoolWithApp(cc) + // cc := func() (abciclient.Client, error) { + // return mockClient, nil + // } + app := kvstore.NewApplication() + mp, cleanup, err := newMempoolWithApp(proxy.NewLocalClientCreator(app)) require.NoError(t, err) defer cleanup() @@ -625,7 +625,7 @@ func TestMempoolTxsBytes(t *testing.T) { func TestMempoolRemoteAppConcurrency(t *testing.T) { sockPath := fmt.Sprintf("unix:///tmp/echo_%v.sock", tmrand.Str(6)) app := kvstore.NewApplication() - cc, server := newRemoteApp(t, sockPath, app) + _, server := newRemoteApp(t, sockPath, app) t.Cleanup(func() { if err := server.Stop(); err != nil { t.Error(err) @@ -634,7 +634,7 @@ func TestMempoolRemoteAppConcurrency(t *testing.T) { cfg := config.ResetTestRoot("mempool_test") - mp, cleanup := newMempoolWithAppAndConfig(cc, cfg) + mp, cleanup := newMempoolWithAppAndConfig(proxy.NewRemoteClientCreator(sockPath, "socket", true), cfg) defer cleanup() // generate small number of txs diff --git a/mempool/v1/mempool_bench_test.go b/mempool/v1/mempool_bench_test.go index ca23f1479..bad8ec8ab 100644 --- a/mempool/v1/mempool_bench_test.go +++ b/mempool/v1/mempool_bench_test.go @@ -1,14 +1,13 @@ package v1 import ( - "context" "fmt" "math/rand" "testing" "time" "github.com/stretchr/testify/require" - "github.com/tendermint/tendermint/internal/mempool" + "github.com/tendermint/tendermint/mempool" ) func BenchmarkTxMempool_CheckTx(b *testing.B) { @@ -27,6 +26,6 @@ func BenchmarkTxMempool_CheckTx(b *testing.B) { tx := []byte(fmt.Sprintf("%X=%d", prefix, priority)) b.StartTimer() - require.NoError(b, txmp.CheckTx(context.Background(), tx, nil, mempool.TxInfo{})) + require.NoError(b, txmp.CheckTx(tx, nil, mempool.TxInfo{})) } } diff --git a/mempool/v1/mempool_test.go b/mempool/v1/mempool_test.go index ffac223ba..385a782b1 100644 --- a/mempool/v1/mempool_test.go +++ b/mempool/v1/mempool_test.go @@ -2,7 +2,6 @@ package v1 import ( "bytes" - "context" "errors" "fmt" "math/rand" @@ -16,13 +15,13 @@ import ( "github.com/stretchr/testify/require" - abciclient "github.com/tendermint/tendermint/abci/client" "github.com/tendermint/tendermint/abci/example/code" "github.com/tendermint/tendermint/abci/example/kvstore" abci "github.com/tendermint/tendermint/abci/types" "github.com/tendermint/tendermint/config" - "github.com/tendermint/tendermint/internal/mempool" "github.com/tendermint/tendermint/libs/log" + "github.com/tendermint/tendermint/mempool" + "github.com/tendermint/tendermint/proxy" "github.com/tendermint/tendermint/types" ) @@ -77,12 +76,12 @@ func setup(t testing.TB, cacheSize int, options ...TxMempoolOption) *TxMempool { t.Helper() app := &application{kvstore.NewApplication()} - cc := abciclient.NewLocalCreator(app) + cc := proxy.NewLocalClientCreator(app) cfg := config.ResetTestRoot(strings.ReplaceAll(t.Name(), "/", "|")) cfg.Mempool.CacheSize = cacheSize - appConnMem, err := cc() + appConnMem, err := cc.NewABCIClient() require.NoError(t, err) require.NoError(t, appConnMem.Start()) @@ -111,7 +110,7 @@ func checkTxs(t *testing.T, txmp *TxMempool, numTxs int, peerID uint16) []testTx tx: []byte(fmt.Sprintf("sender-%d-%d=%X=%d", i, peerID, prefix, priority)), priority: priority, } - require.NoError(t, txmp.CheckTx(context.Background(), txs[i].tx, nil, txInfo)) + require.NoError(t, txmp.CheckTx(txs[i].tx, nil, txInfo)) } return txs @@ -327,13 +326,13 @@ func TestTxMempool_CheckTxExceedsMaxSize(t *testing.T) { _, err := rng.Read(tx) require.NoError(t, err) - require.Error(t, txmp.CheckTx(context.Background(), tx, nil, mempool.TxInfo{SenderID: 0})) + require.Error(t, txmp.CheckTx(tx, nil, mempool.TxInfo{SenderID: 0})) tx = make([]byte, txmp.config.MaxTxBytes-1) _, err = rng.Read(tx) require.NoError(t, err) - require.NoError(t, txmp.CheckTx(context.Background(), tx, nil, mempool.TxInfo{SenderID: 0})) + require.NoError(t, txmp.CheckTx(tx, nil, mempool.TxInfo{SenderID: 0})) } func TestTxMempool_CheckTxSamePeer(t *testing.T) { @@ -347,8 +346,8 @@ func TestTxMempool_CheckTxSamePeer(t *testing.T) { tx := []byte(fmt.Sprintf("sender-0=%X=%d", prefix, 50)) - require.NoError(t, txmp.CheckTx(context.Background(), tx, nil, mempool.TxInfo{SenderID: peerID})) - require.Error(t, txmp.CheckTx(context.Background(), tx, nil, mempool.TxInfo{SenderID: peerID})) + require.NoError(t, txmp.CheckTx(tx, nil, mempool.TxInfo{SenderID: peerID})) + require.Error(t, txmp.CheckTx(tx, nil, mempool.TxInfo{SenderID: peerID})) } func TestTxMempool_CheckTxSameSender(t *testing.T) { @@ -367,9 +366,9 @@ func TestTxMempool_CheckTxSameSender(t *testing.T) { tx1 := []byte(fmt.Sprintf("sender-0=%X=%d", prefix1, 50)) tx2 := []byte(fmt.Sprintf("sender-0=%X=%d", prefix2, 50)) - require.NoError(t, txmp.CheckTx(context.Background(), tx1, nil, mempool.TxInfo{SenderID: peerID})) + require.NoError(t, txmp.CheckTx(tx1, nil, mempool.TxInfo{SenderID: peerID})) require.Equal(t, 1, txmp.Size()) - require.NoError(t, txmp.CheckTx(context.Background(), tx2, nil, mempool.TxInfo{SenderID: peerID})) + require.NoError(t, txmp.CheckTx(tx2, nil, mempool.TxInfo{SenderID: peerID})) require.Equal(t, 1, txmp.Size()) } @@ -522,7 +521,7 @@ func TestTxMempool_CheckTxPostCheckError(t *testing.T) { } require.Equal(t, expectedErrString, checkTxRes.CheckTx.MempoolError) } - require.NoError(t, txmp.CheckTx(context.Background(), tx, callback, mempool.TxInfo{SenderID: 0})) + require.NoError(t, txmp.CheckTx(tx, callback, mempool.TxInfo{SenderID: 0})) }) } } diff --git a/mempool/v1/reactor_test.go b/mempool/v1/reactor_test.go index 0454ad9c5..3f7c3df46 100644 --- a/mempool/v1/reactor_test.go +++ b/mempool/v1/reactor_test.go @@ -1,145 +1,146 @@ package v1 -import ( - "os" - "strings" - "sync" - "testing" +// import ( +// "os" +// "strings" +// "sync" +// "testing" - "github.com/stretchr/testify/require" - "github.com/tendermint/tendermint/abci/example/kvstore" - "github.com/tendermint/tendermint/config" - tmsync "github.com/tendermint/tendermint/internal/libs/sync" - "github.com/tendermint/tendermint/internal/mempool" - "github.com/tendermint/tendermint/internal/p2p" - "github.com/tendermint/tendermint/internal/p2p/p2ptest" - "github.com/tendermint/tendermint/libs/log" - protomem "github.com/tendermint/tendermint/proto/tendermint/mempool" - "github.com/tendermint/tendermint/types" -) +// "github.com/stretchr/testify/require" +// "github.com/tendermint/tendermint/abci/example/kvstore" +// "github.com/tendermint/tendermint/config" +// "github.com/tendermint/tendermint/libs/log" +// tmsync "github.com/tendermint/tendermint/libs/sync" +// "github.com/tendermint/tendermint/mempool" +// "github.com/tendermint/tendermint/p2p" -type reactorTestSuite struct { - network *p2ptest.Network - logger log.Logger +// // "github.com/tendermint/tendermint/p2p" +// protomem "github.com/tendermint/tendermint/proto/tendermint/mempool" +// "github.com/tendermint/tendermint/types" +// ) - reactors map[types.NodeID]*Reactor - mempoolChannels map[types.NodeID]*p2p.Channel - mempools map[types.NodeID]*TxMempool - kvstores map[types.NodeID]*kvstore.Application +// type reactorTestSuite struct { +// network *p2ptest.Network +// logger log.Logger - peerChans map[types.NodeID]chan p2p.PeerUpdate - peerUpdates map[types.NodeID]*p2p.PeerUpdates +// reactors map[types.NodeID]*Reactor +// mempoolChannels map[types.NodeID]*p2p.Channel +// mempools map[types.NodeID]*TxMempool +// kvstores map[types.NodeID]*kvstore.Application - nodes []types.NodeID -} +// peerChans map[types.NodeID]chan p2p.PeerUpdate +// peerUpdates map[types.NodeID]*p2p.PeerUpdates -func setupReactors(t *testing.T, numNodes int, chBuf uint) *reactorTestSuite { - t.Helper() +// nodes []types.NodeID +// } - cfg, err := config.ResetTestRoot(strings.ReplaceAll(t.Name(), "/", "|")) - require.NoError(t, err) - t.Cleanup(func() { os.RemoveAll(cfg.RootDir) }) +// func setupReactors(t *testing.T, numNodes int, chBuf uint) *reactorTestSuite { +// t.Helper() - rts := &reactorTestSuite{ - logger: log.TestingLogger().With("testCase", t.Name()), - network: p2ptest.MakeNetwork(t, p2ptest.NetworkOptions{NumNodes: numNodes}), - reactors: make(map[types.NodeID]*Reactor, numNodes), - mempoolChannels: make(map[types.NodeID]*p2p.Channel, numNodes), - mempools: make(map[types.NodeID]*TxMempool, numNodes), - kvstores: make(map[types.NodeID]*kvstore.Application, numNodes), - peerChans: make(map[types.NodeID]chan p2p.PeerUpdate, numNodes), - peerUpdates: make(map[types.NodeID]*p2p.PeerUpdates, numNodes), - } +// cfg, err := config.ResetTestRoot(strings.ReplaceAll(t.Name(), "/", "|")) +// require.NoError(t, err) +// t.Cleanup(func() { os.RemoveAll(cfg.RootDir) }) - chDesc := p2p.ChannelDescriptor{ID: byte(mempool.MempoolChannel)} - rts.mempoolChannels = rts.network.MakeChannelsNoCleanup(t, chDesc, new(protomem.Message), int(chBuf)) +// rts := &reactorTestSuite{ +// logger: log.TestingLogger().With("testCase", t.Name()), +// network: p2ptest.MakeNetwork(t, p2ptest.NetworkOptions{NumNodes: numNodes}), +// reactors: make(map[types.NodeID]*Reactor, numNodes), +// mempoolChannels: make(map[types.NodeID]*p2p.Channel, numNodes), +// mempools: make(map[types.NodeID]*TxMempool, numNodes), +// kvstores: make(map[types.NodeID]*kvstore.Application, numNodes), +// peerChans: make(map[types.NodeID]chan p2p.PeerUpdate, numNodes), +// peerUpdates: make(map[types.NodeID]*p2p.PeerUpdates, numNodes), +// } - for nodeID := range rts.network.Nodes { - rts.kvstores[nodeID] = kvstore.NewApplication() +// chDesc := p2p.ChannelDescriptor{ID: byte(mempool.MempoolChannel)} +// rts.mempoolChannels = rts.network.MakeChannelsNoCleanup(t, chDesc, new(protomem.Message), int(chBuf)) - mempool := setup(t, 0) - rts.mempools[nodeID] = mempool +// for nodeID := range rts.network.Nodes { +// rts.kvstores[nodeID] = kvstore.NewApplication() - rts.peerChans[nodeID] = make(chan p2p.PeerUpdate) - rts.peerUpdates[nodeID] = p2p.NewPeerUpdates(rts.peerChans[nodeID], 1) - rts.network.Nodes[nodeID].PeerManager.Register(rts.peerUpdates[nodeID]) +// mempool := setup(t, 0) +// rts.mempools[nodeID] = mempool - rts.reactors[nodeID] = NewReactor( - rts.logger.With("nodeID", nodeID), - cfg.Mempool, - mempool, - rts.mempoolChannels[nodeID], - rts.peerUpdates[nodeID], - ) +// rts.peerChans[nodeID] = make(chan p2p.PeerUpdate) +// rts.peerUpdates[nodeID] = p2p.NewPeerUpdates(rts.peerChans[nodeID], 1) +// rts.network.Nodes[nodeID].PeerManager.Register(rts.peerUpdates[nodeID]) - rts.nodes = append(rts.nodes, nodeID) +// rts.reactors[nodeID] = NewReactor( +// rts.logger.With("nodeID", nodeID), +// cfg.Mempool, +// mempool, +// rts.mempoolChannels[nodeID], +// rts.peerUpdates[nodeID], +// ) - require.NoError(t, rts.reactors[nodeID].Start()) - require.True(t, rts.reactors[nodeID].IsRunning()) - } +// rts.nodes = append(rts.nodes, nodeID) - require.Len(t, rts.reactors, numNodes) +// require.NoError(t, rts.reactors[nodeID].Start()) +// require.True(t, rts.reactors[nodeID].IsRunning()) +// } - t.Cleanup(func() { - for nodeID := range rts.reactors { - if rts.reactors[nodeID].IsRunning() { - require.NoError(t, rts.reactors[nodeID].Stop()) - require.False(t, rts.reactors[nodeID].IsRunning()) - } - } - }) +// require.Len(t, rts.reactors, numNodes) - return rts -} +// t.Cleanup(func() { +// for nodeID := range rts.reactors { +// if rts.reactors[nodeID].IsRunning() { +// require.NoError(t, rts.reactors[nodeID].Stop()) +// require.False(t, rts.reactors[nodeID].IsRunning()) +// } +// } +// }) -func (rts *reactorTestSuite) start(t *testing.T) { - t.Helper() - rts.network.Start(t) - require.Len(t, - rts.network.RandomNode().PeerManager.Peers(), - len(rts.nodes)-1, - "network does not have expected number of nodes") -} +// return rts +// } -func TestReactorBroadcastDoesNotPanic(t *testing.T) { - numNodes := 2 - rts := setupReactors(t, numNodes, 0) +// func (rts *reactorTestSuite) start(t *testing.T) { +// t.Helper() +// rts.network.Start(t) +// require.Len(t, +// rts.network.RandomNode().PeerManager.Peers(), +// len(rts.nodes)-1, +// "network does not have expected number of nodes") +// } - observePanic := func(r interface{}) { - t.Fatal("panic detected in reactor") - } +// func TestReactorBroadcastDoesNotPanic(t *testing.T) { +// numNodes := 2 +// rts := setupReactors(t, numNodes, 0) - primary := rts.nodes[0] - secondary := rts.nodes[1] - primaryReactor := rts.reactors[primary] - primaryMempool := primaryReactor.mempool - secondaryReactor := rts.reactors[secondary] +// observePanic := func(r interface{}) { +// t.Fatal("panic detected in reactor") +// } - primaryReactor.observePanic = observePanic - secondaryReactor.observePanic = observePanic +// primary := rts.nodes[0] +// secondary := rts.nodes[1] +// primaryReactor := rts.reactors[primary] +// primaryMempool := primaryReactor.mempool +// secondaryReactor := rts.reactors[secondary] - firstTx := &WrappedTx{} - primaryMempool.insertTx(firstTx) +// primaryReactor.observePanic = observePanic +// secondaryReactor.observePanic = observePanic - // run the router - rts.start(t) +// firstTx := &WrappedTx{} +// primaryMempool.insertTx(firstTx) - closer := tmsync.NewCloser() - primaryReactor.peerWG.Add(1) - go primaryReactor.broadcastTxRoutine(secondary, closer) +// // run the router +// rts.start(t) - wg := &sync.WaitGroup{} - for i := 0; i < 50; i++ { - next := &WrappedTx{} - wg.Add(1) - go func() { - defer wg.Done() - primaryMempool.insertTx(next) - }() - } +// closer := tmsync.NewCloser() +// primaryReactor.peerWG.Add(1) +// go primaryReactor.broadcastTxRoutine(secondary, closer) - err := primaryReactor.Stop() - require.NoError(t, err) - primaryReactor.peerWG.Wait() - wg.Wait() -} +// wg := &sync.WaitGroup{} +// for i := 0; i < 50; i++ { +// next := &WrappedTx{} +// wg.Add(1) +// go func() { +// defer wg.Done() +// primaryMempool.insertTx(next) +// }() +// } + +// err := primaryReactor.Stop() +// require.NoError(t, err) +// primaryReactor.peerWG.Wait() +// wg.Wait() +// } diff --git a/mempool_bak/mock/mempool.go b/mempool_bak/mock/mempool.go index cd8df2198..95c33b7ec 100644 --- a/mempool_bak/mock/mempool.go +++ b/mempool_bak/mock/mempool.go @@ -41,3 +41,8 @@ func (Mempool) TxsWaitChan() <-chan struct{} { return nil } func (Mempool) InitWAL() error { return nil } func (Mempool) CloseWAL() {} + +// RemoveTxByKey implements mempool.Mempool +func (Mempool) RemoveTxByKey(txKey types.TxKey) error { + return nil +} diff --git a/node/node.go b/node/node.go index 52dad73aa..b1ebe8a50 100644 --- a/node/node.go +++ b/node/node.go @@ -23,12 +23,15 @@ import ( cs "github.com/tendermint/tendermint/consensus" "github.com/tendermint/tendermint/crypto" "github.com/tendermint/tendermint/evidence" + tmjson "github.com/tendermint/tendermint/libs/json" "github.com/tendermint/tendermint/libs/log" tmpubsub "github.com/tendermint/tendermint/libs/pubsub" "github.com/tendermint/tendermint/libs/service" "github.com/tendermint/tendermint/light" mempl "github.com/tendermint/tendermint/mempool" + mempoolv0 "github.com/tendermint/tendermint/mempool/v0" + mempoolv1 "github.com/tendermint/tendermint/mempool/v1" "github.com/tendermint/tendermint/p2p" "github.com/tendermint/tendermint/p2p/pex" "github.com/tendermint/tendermint/privval" @@ -209,7 +212,7 @@ type Node struct { stateStore sm.Store blockStore *store.BlockStore // store the blockchain to disk bcReactor p2p.Reactor // for fast-syncing - mempoolReactor *mempl.Reactor // for gossipping transactions + mempoolReactor p2p.Reactor // for gossipping transactions mempool mempl.Mempool stateSync bool // whether the node should state sync on startup stateSyncReactor *statesync.Reactor // for hosting and restoring state sync snapshots @@ -367,24 +370,56 @@ func createMempoolAndMempoolReactor( state sm.State, memplMetrics *mempl.Metrics, logger log.Logger, -) (*mempl.Reactor, *mempl.TxMempool) { - mempool := mempl.NewTxMempool( - logger.With("module", "mempool"), - config.Mempool, - proxyApp.Mempool(), - state.LastBlockHeight, - mempl.WithMetrics(memplMetrics), - mempl.WithPreCheck(sm.TxPreCheck(state)), - mempl.WithPostCheck(sm.TxPostCheck(state)), - ) +) (mempl.Mempool, p2p.Reactor) { - mempoolReactor := mempl.NewReactor(config.Mempool, mempool) + switch config.Mempool.Version { + case cfg.MempoolV1: + mp := mempoolv1.NewTxMempool( + logger, + config.Mempool, + proxyApp.Mempool(), + state.LastBlockHeight, + mempoolv1.WithMetrics(memplMetrics), + mempoolv1.WithPreCheck(sm.TxPreCheck(state)), + mempoolv1.WithPostCheck(sm.TxPostCheck(state)), + ) - if config.Consensus.WaitForTxs() { - mempool.EnableTxsAvailable() + reactor := mempoolv1.NewReactor( + config.Mempool, + mp, + ) + if config.Consensus.WaitForTxs() { + mp.EnableTxsAvailable() + } + + return mp, reactor + + case cfg.MempoolV0: + mp := mempoolv0.NewCListMempool( + config.Mempool, + proxyApp.Mempool(), + state.LastBlockHeight, + mempoolv0.WithMetrics(memplMetrics), + mempoolv0.WithPreCheck(sm.TxPreCheck(state)), + mempoolv0.WithPostCheck(sm.TxPostCheck(state)), + ) + + mp.SetLogger(logger) + mp.SetLogger(logger) + + reactor := mempoolv0.NewReactor( + config.Mempool, + mp, + ) + if config.Consensus.WaitForTxs() { + mp.EnableTxsAvailable() + } + + return mp, reactor + + default: + return nil, nil } - - return mempoolReactor, mempool } func createEvidenceReactor(config *cfg.Config, dbProvider DBProvider, @@ -430,7 +465,7 @@ func createConsensusReactor(config *cfg.Config, state sm.State, blockExec *sm.BlockExecutor, blockStore sm.BlockStore, - mempool *mempl.TxMempool, + mempool mempl.Mempool, evidencePool *evidence.Pool, privValidator types.PrivValidator, csMetrics *cs.Metrics, @@ -532,7 +567,7 @@ func createSwitch(config *cfg.Config, transport p2p.Transport, p2pMetrics *p2p.Metrics, peerFilters []p2p.PeerFilterFunc, - mempoolReactor *mempl.Reactor, + mempoolReactor p2p.Reactor, bcReactor p2p.Reactor, stateSyncReactor *statesync.Reactor, consensusReactor *cs.Reactor, @@ -757,7 +792,7 @@ func NewNode(config *cfg.Config, csMetrics, p2pMetrics, memplMetrics, smMetrics := metricsProvider(genDoc.ChainID) // Make MempoolReactor - mempoolReactor, mempool := createMempoolAndMempoolReactor(config, proxyApp, state, memplMetrics, logger) + mempool, mempoolReactor := createMempoolAndMempoolReactor(config, proxyApp, state, memplMetrics, logger) // Make Evidence Reactor evidenceReactor, evidencePool, err := createEvidenceReactor(config, dbProvider, stateDB, blockStore, logger) @@ -1217,7 +1252,7 @@ func (n *Node) ConsensusReactor() *cs.Reactor { } // MempoolReactor returns the Node's mempool reactor. -func (n *Node) MempoolReactor() *mempl.Reactor { +func (n *Node) MempoolReactor() p2p.Reactor { return n.mempoolReactor } diff --git a/node/node_test.go b/node/node_test.go index 3eb707a73..efcb492a6 100644 --- a/node/node_test.go +++ b/node/node_test.go @@ -21,6 +21,8 @@ import ( "github.com/tendermint/tendermint/libs/log" tmrand "github.com/tendermint/tendermint/libs/rand" mempl "github.com/tendermint/tendermint/mempool" + mempoolv0 "github.com/tendermint/tendermint/mempool/v0" + mempoolv1 "github.com/tendermint/tendermint/mempool/v1" "github.com/tendermint/tendermint/p2p" "github.com/tendermint/tendermint/p2p/conn" p2pmock "github.com/tendermint/tendermint/p2p/mock" @@ -243,15 +245,26 @@ func TestCreateProposalBlock(t *testing.T) { // Make Mempool memplMetrics := mempl.PrometheusMetrics("node_test_1") - mempool := mempl.NewTxMempool( - logger, - config.Mempool, - proxyApp.Mempool(), - state.LastBlockHeight, - mempl.WithMetrics(memplMetrics), - mempl.WithPreCheck(sm.TxPreCheck(state)), - mempl.WithPostCheck(sm.TxPostCheck(state)), - ) + var mempool mempl.Mempool + + switch config.Mempool.Version { + case cfg.MempoolV0: + mempool = mempoolv0.NewCListMempool(config.Mempool, + proxyApp.Mempool(), + state.LastBlockHeight, + mempoolv0.WithMetrics(memplMetrics), + mempoolv0.WithPreCheck(sm.TxPreCheck(state)), + mempoolv0.WithPostCheck(sm.TxPostCheck(state))) + case cfg.MempoolV1: + mempool = mempoolv1.NewTxMempool(logger, + config.Mempool, + proxyApp.Mempool(), + state.LastBlockHeight, + mempoolv1.WithMetrics(memplMetrics), + mempoolv1.WithPreCheck(sm.TxPreCheck(state)), + mempoolv1.WithPostCheck(sm.TxPostCheck(state)), + ) + } // Make EvidencePool evidenceDB := dbm.NewMemDB() @@ -335,15 +348,25 @@ func TestMaxProposalBlockSize(t *testing.T) { // Make Mempool memplMetrics := mempl.PrometheusMetrics("node_test_2") - mempool := mempl.NewTxMempool( - logger, - config.Mempool, - proxyApp.Mempool(), - state.LastBlockHeight, - mempl.WithMetrics(memplMetrics), - mempl.WithPreCheck(sm.TxPreCheck(state)), - mempl.WithPostCheck(sm.TxPostCheck(state)), - ) + var mempool mempl.Mempool + switch config.Mempool.Version { + case cfg.MempoolV0: + mempool = mempoolv0.NewCListMempool(config.Mempool, + proxyApp.Mempool(), + state.LastBlockHeight, + mempoolv0.WithMetrics(memplMetrics), + mempoolv0.WithPreCheck(sm.TxPreCheck(state)), + mempoolv0.WithPostCheck(sm.TxPostCheck(state))) + case cfg.MempoolV1: + mempool = mempoolv1.NewTxMempool(logger, + config.Mempool, + proxyApp.Mempool(), + state.LastBlockHeight, + mempoolv1.WithMetrics(memplMetrics), + mempoolv1.WithPreCheck(sm.TxPreCheck(state)), + mempoolv1.WithPostCheck(sm.TxPostCheck(state)), + ) + } // fill the mempool with one txs just below the maximum size txLength := int(types.MaxDataBytesNoEvidence(maxBytes, 1)) diff --git a/test/fuzz/mempool/checktx.go b/test/fuzz/mempool/checktx.go deleted file mode 100644 index 30c93d007..000000000 --- a/test/fuzz/mempool/checktx.go +++ /dev/null @@ -1,35 +0,0 @@ -package checktx - -import ( - "github.com/tendermint/tendermint/abci/example/kvstore" - "github.com/tendermint/tendermint/config" - "github.com/tendermint/tendermint/libs/log" - mempl "github.com/tendermint/tendermint/mempool" - "github.com/tendermint/tendermint/proxy" -) - -var mempool mempl.Mempool - -func init() { - app := kvstore.NewApplication() - cc := proxy.NewLocalClientCreator(app) - appConnMem, _ := cc.NewABCIClient() - err := appConnMem.Start() - if err != nil { - panic(err) - } - - cfg := config.DefaultMempoolConfig() - cfg.Broadcast = false - - mempool = mempl.NewTxMempool(log.NewNopLogger(), cfg, appConnMem, 0) -} - -func Fuzz(data []byte) int { - err := mempool.CheckTx(data, nil, mempl.TxInfo{}) - if err != nil { - return 0 - } - - return 1 -} diff --git a/test/maverick/consensus/replay_stubs.go b/test/maverick/consensus/replay_stubs.go index ed5aa5626..5f3badd18 100644 --- a/test/maverick/consensus/replay_stubs.go +++ b/test/maverick/consensus/replay_stubs.go @@ -22,6 +22,7 @@ func (emptyMempool) SizeBytes() int64 { return 0 } func (emptyMempool) CheckTx(_ types.Tx, _ func(*abci.Response), _ mempl.TxInfo) error { return nil } +func (emptyMempool) RemoveTxByKey(txKey types.TxKey) error { return nil } func (emptyMempool) ReapMaxBytesMaxGas(_, _ int64) types.Txs { return types.Txs{} } func (emptyMempool) ReapMaxTxs(n int) types.Txs { return types.Txs{} } func (emptyMempool) Update( diff --git a/test/maverick/node/node.go b/test/maverick/node/node.go index 33ce51946..919554c6d 100644 --- a/test/maverick/node/node.go +++ b/test/maverick/node/node.go @@ -32,6 +32,8 @@ import ( "github.com/tendermint/tendermint/libs/service" "github.com/tendermint/tendermint/light" mempl "github.com/tendermint/tendermint/mempool" + mempoolv0 "github.com/tendermint/tendermint/mempool/v0" + mempoolv1 "github.com/tendermint/tendermint/mempool/v1" "github.com/tendermint/tendermint/p2p" "github.com/tendermint/tendermint/p2p/pex" "github.com/tendermint/tendermint/privval" @@ -238,7 +240,7 @@ type Node struct { stateStore sm.Store blockStore *store.BlockStore // store the blockchain to disk bcReactor p2p.Reactor // for fast-syncing - mempoolReactor *mempl.Reactor // for gossipping transactions + mempoolReactor p2p.Reactor // for gossipping transactions mempool mempl.Mempool stateSync bool // whether the node should state sync on startup stateSyncReactor *statesync.Reactor // for hosting and restoring state sync snapshots @@ -378,24 +380,56 @@ func onlyValidatorIsUs(state sm.State, pubKey crypto.PubKey) bool { } func createMempoolAndMempoolReactor(config *cfg.Config, proxyApp proxy.AppConns, - state sm.State, memplMetrics *mempl.Metrics, logger log.Logger) (*mempl.Reactor, *mempl.TxMempool) { + state sm.State, memplMetrics *mempl.Metrics, logger log.Logger) (p2p.Reactor, mempl.Mempool) { - mempool := mempl.NewTxMempool( - logger.With("module", "mempool"), - config.Mempool, - proxyApp.Mempool(), - state.LastBlockHeight, - mempl.WithMetrics(memplMetrics), - mempl.WithPreCheck(sm.TxPreCheck(state)), - mempl.WithPostCheck(sm.TxPostCheck(state)), - ) + switch config.Mempool.Version { + case cfg.MempoolV1: + mp := mempoolv1.NewTxMempool( + logger, + config.Mempool, + proxyApp.Mempool(), + state.LastBlockHeight, + mempoolv1.WithMetrics(memplMetrics), + mempoolv1.WithPreCheck(sm.TxPreCheck(state)), + mempoolv1.WithPostCheck(sm.TxPostCheck(state)), + ) - mempoolReactor := mempl.NewReactor(config.Mempool, mempool) + reactor := mempoolv1.NewReactor( + config.Mempool, + mp, + ) + if config.Consensus.WaitForTxs() { + mp.EnableTxsAvailable() + } - if config.Consensus.WaitForTxs() { - mempool.EnableTxsAvailable() + return reactor, mp + + case cfg.MempoolV0: + mp := mempoolv0.NewCListMempool( + config.Mempool, + proxyApp.Mempool(), + state.LastBlockHeight, + mempoolv0.WithMetrics(memplMetrics), + mempoolv0.WithPreCheck(sm.TxPreCheck(state)), + mempoolv0.WithPostCheck(sm.TxPostCheck(state)), + ) + + mp.SetLogger(logger) + mp.SetLogger(logger) + + reactor := mempoolv0.NewReactor( + config.Mempool, + mp, + ) + if config.Consensus.WaitForTxs() { + mp.EnableTxsAvailable() + } + + return reactor, mp + + default: + return nil, nil } - return mempoolReactor, mempool } func createEvidenceReactor(config *cfg.Config, dbProvider DBProvider, @@ -441,7 +475,7 @@ func createConsensusReactor(config *cfg.Config, state sm.State, blockExec *sm.BlockExecutor, blockStore sm.BlockStore, - mempool *mempl.TxMempool, + mempool mempl.Mempool, evidencePool *evidence.Pool, privValidator types.PrivValidator, csMetrics *consensus.Metrics, @@ -545,7 +579,7 @@ func createSwitch(config *cfg.Config, transport p2p.Transport, p2pMetrics *p2p.Metrics, peerFilters []p2p.PeerFilterFunc, - mempoolReactor *mempl.Reactor, + mempoolReactor p2p.Reactor, bcReactor p2p.Reactor, stateSyncReactor *statesync.Reactor, consensusReactor *cs.Reactor, @@ -1212,7 +1246,7 @@ func (n *Node) ConsensusReactor() *cs.Reactor { } // MempoolReactor returns the Node's mempool reactor. -func (n *Node) MempoolReactor() *mempl.Reactor { +func (n *Node) MempoolReactor() p2p.Reactor { return n.mempoolReactor }