This commit is contained in:
Aleksandr Bezobchuk
2022-06-14 10:24:27 -04:00
parent fee86dd29b
commit d33b440098
3 changed files with 105 additions and 104 deletions
+18 -17
View File
@@ -525,25 +525,26 @@ func TestTxMempool_CheckTxPostCheckError(t *testing.T) {
}
}
// // A cleanupFunc cleans up any config / test files created for a particular
// // test.
// type cleanupFunc func()
// A cleanupFunc cleans up any config / test files created for a particular
// test.
type cleanupFunc func()
// func newMempoolWithApp(cc proxy.ClientCreator) (*CListMempool, cleanupFunc) {
// return newMempoolWithAppAndConfig(cc, cfg.ResetTestRoot("mempool_test"))
// }
func newMempoolWithApp(cc proxy.ClientCreator) (*TxMempool, cleanupFunc) {
return newMempoolWithAppAndConfig(cc, config.ResetTestRoot("mempool_test"))
}
// func newMempoolWithAppAndConfig(cc proxy.ClientCreator, config *cfg.Config) (*CListMempool, cleanupFunc) {
// appConnMem, _ := cc.NewABCIClient()
// appConnMem.SetLogger(log.TestingLogger().With("module", "abci-client", "connection", "mempool"))
// err := appConnMem.Start()
// if err != nil {
// panic(err)
// }
// mempool := NewCListMempool(config.Mempool, appConnMem, 0)
// mempool.SetLogger(log.TestingLogger())
// return mempool, func() { os.RemoveAll(config.RootDir) }
// }
func newMempoolWithAppAndConfig(cc proxy.ClientCreator, config *config.Config) (*TxMempool, cleanupFunc) {
appConnMem, _ := cc.NewABCIClient()
appConnMem.SetLogger(log.TestingLogger().With("module", "abci-client", "connection", "mempool"))
err := appConnMem.Start()
if err != nil {
panic(err)
}
mempool := NewTxMempool(log.TestingLogger(), config.Mempool, appConnMem, 0)
return mempool, func() { os.RemoveAll(config.RootDir) }
}
// func ensureNoFire(t *testing.T, ch <-chan struct{}, timeoutMS int) {
// timer := time.NewTimer(time.Duration(timeoutMS) * time.Millisecond)
+4 -3
View File
@@ -12,9 +12,10 @@ type Mempool struct{}
var _ mempl.Mempool = Mempool{}
func (Mempool) Lock() {}
func (Mempool) Unlock() {}
func (Mempool) Size() int { return 0 }
func (Mempool) Lock() {}
func (Mempool) Unlock() {}
func (Mempool) Size() int { return 0 }
func (Mempool) SizeBytes() int64 { return 0 }
func (Mempool) CheckTx(_ types.Tx, _ func(*abci.Response), _ mempl.TxInfo) error {
return nil
}
+83 -84
View File
@@ -14,10 +14,8 @@ import (
"github.com/stretchr/testify/require"
"github.com/tendermint/tendermint/abci/example/kvstore"
abci "github.com/tendermint/tendermint/abci/types"
cfg "github.com/tendermint/tendermint/config"
"github.com/tendermint/tendermint/libs/log"
tmrand "github.com/tendermint/tendermint/libs/rand"
"github.com/tendermint/tendermint/p2p"
"github.com/tendermint/tendermint/p2p/mock"
memproto "github.com/tendermint/tendermint/proto/tendermint/mempool"
@@ -65,65 +63,66 @@ func TestReactorBroadcastTxsMessage(t *testing.T) {
waitForTxsOnReactors(t, txs, reactors)
}
// regression test for https://github.com/tendermint/tendermint/issues/5408
func TestReactorConcurrency(t *testing.T) {
config := cfg.TestConfig()
const N = 2
reactors := makeAndConnectReactors(config, N)
defer func() {
for _, r := range reactors {
if err := r.Stop(); err != nil {
assert.NoError(t, err)
}
}
}()
for _, r := range reactors {
for _, peer := range r.Switch.Peers().List() {
peer.Set(types.PeerStateKey, peerState{1})
}
}
var wg sync.WaitGroup
// // regression test for https://github.com/tendermint/tendermint/issues/5408
// func TestReactorConcurrency(t *testing.T) {
// config := cfg.TestConfig()
// const N = 2
// reactors := makeAndConnectReactors(config, N)
// defer func() {
// for _, r := range reactors {
// if err := r.Stop(); err != nil {
// assert.NoError(t, err)
// }
// }
// }()
// for _, r := range reactors {
// for _, peer := range r.Switch.Peers().List() {
// peer.Set(types.PeerStateKey, peerState{1})
// }
// }
// var wg sync.WaitGroup
const numTxs = 5
// const numTxs = 5
for i := 0; i < 1000; i++ {
wg.Add(2)
// for i := 0; i < 1000; i++ {
// wg.Add(2)
// 1. submit a bunch of txs
// 2. update the whole mempool
txs := checkTxs(t, reactors[0].mempool, numTxs, UnknownPeerID)
go func() {
defer wg.Done()
// // 1. submit a bunch of txs
// // 2. update the whole mempool
// txs := checkTxs(t, reactors[0].mempool, numTxs, UnknownPeerID)
// go func() {
// defer wg.Done()
reactors[0].mempool.Lock()
defer reactors[0].mempool.Unlock()
// reactors[0].mempool.Lock()
// defer reactors[0].mempool.Unlock()
deliverTxResponses := make([]*abci.ResponseDeliverTx, len(txs))
for i := range txs {
deliverTxResponses[i] = &abci.ResponseDeliverTx{Code: 0}
}
err := reactors[0].mempool.Update(1, txs, deliverTxResponses, nil, nil)
assert.NoError(t, err)
}()
// deliverTxResponses := make([]*abci.ResponseDeliverTx, len(txs))
// for i := range txs {
// deliverTxResponses[i] = &abci.ResponseDeliverTx{Code: 0}
// }
// 1. submit a bunch of txs
// 2. update none
_ = checkTxs(t, reactors[1].mempool, numTxs, UnknownPeerID)
go func() {
defer wg.Done()
// err := reactors[0].mempool.Update(1, txs, deliverTxResponses, nil, nil)
// assert.NoError(t, err)
// }()
reactors[1].mempool.Lock()
defer reactors[1].mempool.Unlock()
err := reactors[1].mempool.Update(1, []types.Tx{}, make([]*abci.ResponseDeliverTx, 0), nil, nil)
assert.NoError(t, err)
}()
// // 1. submit a bunch of txs
// // 2. update none
// _ = checkTxs(t, reactors[1].mempool, numTxs, UnknownPeerID)
// go func() {
// defer wg.Done()
// 1. flush the mempool
reactors[1].mempool.Flush()
}
// reactors[1].mempool.Lock()
// defer reactors[1].mempool.Unlock()
// err := reactors[1].mempool.Update(1, []types.Tx{}, make([]*abci.ResponseDeliverTx, 0), nil, nil)
// assert.NoError(t, err)
// }()
wg.Wait()
}
// // 1. flush the mempool
// reactors[1].mempool.Flush()
// }
// wg.Wait()
// }
// Send a bunch of txs to the first reactor's mempool, claiming it came from peer
// ensure peer gets no txs.
@@ -149,40 +148,40 @@ func TestReactorNoBroadcastToSender(t *testing.T) {
ensureNoTxs(t, reactors[peerID], 100*time.Millisecond)
}
func TestReactor_MaxTxBytes(t *testing.T) {
config := cfg.TestConfig()
// func TestReactor_MaxTxBytes(t *testing.T) {
// config := cfg.TestConfig()
const N = 2
reactors := makeAndConnectReactors(config, N)
defer func() {
for _, r := range reactors {
if err := r.Stop(); err != nil {
assert.NoError(t, err)
}
}
}()
for _, r := range reactors {
for _, peer := range r.Switch.Peers().List() {
peer.Set(types.PeerStateKey, peerState{1})
}
}
// const N = 2
// reactors := makeAndConnectReactors(config, N)
// defer func() {
// for _, r := range reactors {
// if err := r.Stop(); err != nil {
// assert.NoError(t, err)
// }
// }
// }()
// for _, r := range reactors {
// for _, peer := range r.Switch.Peers().List() {
// peer.Set(types.PeerStateKey, peerState{1})
// }
// }
// Broadcast a tx, which has the max size
// => ensure it's received by the second reactor.
tx1 := tmrand.Bytes(config.Mempool.MaxTxBytes)
err := reactors[0].mempool.CheckTx(tx1, nil, TxInfo{SenderID: UnknownPeerID})
require.NoError(t, err)
waitForTxsOnReactors(t, []types.Tx{tx1}, reactors)
// // Broadcast a tx, which has the max size
// // => ensure it's received by the second reactor.
// tx1 := tmrand.Bytes(config.Mempool.MaxTxBytes)
// err := reactors[0].mempool.CheckTx(tx1, nil, TxInfo{SenderID: UnknownPeerID})
// require.NoError(t, err)
// waitForTxsOnReactors(t, []types.Tx{tx1}, reactors)
reactors[0].mempool.Flush()
reactors[1].mempool.Flush()
// reactors[0].mempool.Flush()
// reactors[1].mempool.Flush()
// Broadcast a tx, which is beyond the max size
// => ensure it's not sent
tx2 := tmrand.Bytes(config.Mempool.MaxTxBytes + 1)
err = reactors[0].mempool.CheckTx(tx2, nil, TxInfo{SenderID: UnknownPeerID})
require.Error(t, err)
}
// // Broadcast a tx, which is beyond the max size
// // => ensure it's not sent
// tx2 := tmrand.Bytes(config.Mempool.MaxTxBytes + 1)
// err = reactors[0].mempool.CheckTx(tx2, nil, TxInfo{SenderID: UnknownPeerID})
// require.Error(t, err)
// }
func TestBroadcastTxForPeerStopsWhenPeerStops(t *testing.T) {
if testing.Short() {
@@ -318,7 +317,7 @@ func makeAndConnectReactors(config *cfg.Config, n int) []*Reactor {
return reactors
}
func waitForTxsOnReactors(t *testing.T, txs types.Txs, reactors []*Reactor) {
func waitForTxsOnReactors(t *testing.T, txs []testTx, reactors []*Reactor) {
// wait for the txs in all mempools
wg := new(sync.WaitGroup)
for i, reactor := range reactors {
@@ -343,7 +342,7 @@ func waitForTxsOnReactors(t *testing.T, txs types.Txs, reactors []*Reactor) {
}
}
func waitForTxsOnReactor(t *testing.T, txs types.Txs, reactor *Reactor, reactorIndex int) {
func waitForTxsOnReactor(t *testing.T, txs []testTx, reactor *Reactor, reactorIndex int) {
mempool := reactor.mempool
for mempool.Size() < len(txs) {
time.Sleep(time.Millisecond * 100)