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.

622 lines
18 KiB

mempool no gossip back (#2778) Closes #1798 This is done by making every mempool tx maintain a list of peers who its received the tx from. Instead of using the 20byte peer ID, it instead uses a local map from peerID to uint16 counter, so every peer adds 2 bytes. (Word aligned to probably make it 8 bytes) This also required resetting the callback function on every CheckTx. This likely has performance ramifications for instruction caching. The actual setting operation isn't costly with the removal of defers in this PR. * Make the mempool not gossip txs back to peers its received it from * Fix adversarial memleak * Don't break interface * Update changelog * Forgot to add a mtx * forgot a mutex * Update mempool/reactor.go Co-Authored-By: ValarDragon <ValarDragon@users.noreply.github.com> * Update mempool/mempool.go Co-Authored-By: ValarDragon <ValarDragon@users.noreply.github.com> * Use unknown peer ID Co-Authored-By: ValarDragon <ValarDragon@users.noreply.github.com> * fix compilation * use next wait chan logic when skipping * Minor fixes * Add TxInfo * Add reverse map * Make activeID's auto-reserve 0 * 0 -> UnknownPeerID Co-Authored-By: ValarDragon <ValarDragon@users.noreply.github.com> * Switch to making the normal case set a callback on the reqres object The recheck case is still done via the global callback, and stats are also set via global callback * fix merge conflict * Addres comments * Add cache tests * add cache tests * minor fixes * update metrics in reqResCb and reformat code * goimport -w mempool/reactor.go * mempool: update memTx senders I had to introduce txsMap for quick mempoolTx lookups. * change senders type from []uint16 to sync.Map Fixes DATA RACE: ``` Read at 0x00c0013fcd3a by goroutine 183: github.com/tendermint/tendermint/mempool.(*MempoolReactor).broadcastTxRoutine() /go/src/github.com/tendermint/tendermint/mempool/reactor.go:195 +0x3c7 Previous write at 0x00c0013fcd3a by D[2019-02-27|10:10:49.058] Read PacketMsg switch=3 peer=35bc1e3558c182927b31987eeff3feb3d58a0fc5@127.0.0.1 :46552 conn=MConn{pipe} packet="PacketMsg{30:2B06579D0A143EB78F3D3299DE8213A51D4E11FB05ACE4D6A14F T:1}" goroutine 190: github.com/tendermint/tendermint/mempool.(*Mempool).CheckTxWithInfo() /go/src/github.com/tendermint/tendermint/mempool/mempool.go:387 +0xdc1 github.com/tendermint/tendermint/mempool.(*MempoolReactor).Receive() /go/src/github.com/tendermint/tendermint/mempool/reactor.go:134 +0xb04 github.com/tendermint/tendermint/p2p.createMConnection.func1() /go/src/github.com/tendermint/tendermint/p2p/peer.go:374 +0x25b github.com/tendermint/tendermint/p2p/conn.(*MConnection).recvRoutine() /go/src/github.com/tendermint/tendermint/p2p/conn/connection.go:599 +0xcce Goroutine 183 (running) created at: D[2019-02-27|10:10:49.058] Send switch=2 peer=1efafad5443abeea4b7a8155218e4369525d987e@127.0.0.1:46193 channel=48 conn=MConn{pipe} m sgBytes=2B06579D0A146194480ADAE00C2836ED7125FEE65C1D9DD51049 github.com/tendermint/tendermint/mempool.(*MempoolReactor).AddPeer() /go/src/github.com/tendermint/tendermint/mempool/reactor.go:105 +0x1b1 github.com/tendermint/tendermint/p2p.(*Switch).startInitPeer() /go/src/github.com/tendermint/tendermint/p2p/switch.go:683 +0x13b github.com/tendermint/tendermint/p2p.(*Switch).addPeer() /go/src/github.com/tendermint/tendermint/p2p/switch.go:650 +0x585 github.com/tendermint/tendermint/p2p.(*Switch).addPeerWithConnection() /go/src/github.com/tendermint/tendermint/p2p/test_util.go:145 +0x939 github.com/tendermint/tendermint/p2p.Connect2Switches.func2() /go/src/github.com/tendermint/tendermint/p2p/test_util.go:109 +0x50 I[2019-02-27|10:10:49.058] Added good transaction validator=0 tx=43B4D1F0F03460BD262835C4AA560DB860CFBBE85BD02386D83DAC38C67B3AD7 res="&{CheckTx:gas_w anted:1 }" height=0 total=375 Goroutine 190 (running) created at: github.com/tendermint/tendermint/p2p/conn.(*MConnection).OnStart() /go/src/github.com/tendermint/tendermint/p2p/conn/connection.go:210 +0x313 github.com/tendermint/tendermint/libs/common.(*BaseService).Start() /go/src/github.com/tendermint/tendermint/libs/common/service.go:139 +0x4df github.com/tendermint/tendermint/p2p.(*peer).OnStart() /go/src/github.com/tendermint/tendermint/p2p/peer.go:179 +0x56 github.com/tendermint/tendermint/libs/common.(*BaseService).Start() /go/src/github.com/tendermint/tendermint/libs/common/service.go:139 +0x4df github.com/tendermint/tendermint/p2p.(*peer).Start() <autogenerated>:1 +0x43 github.com/tendermint/tendermint/p2p.(*Switch).startInitPeer() ``` * explain the choice of a map DS for senders * extract ids pool/mapper to a separate struct * fix literal copies lock value from senders: sync.Map contains sync.Mutex * use sync.Map#LoadOrStore instead of Load * fixes after Ismail's review * rename resCbNormal to resCbFirstTime
6 years ago
mempool: move interface into mempool package (#3524) ## Description Refs #2659 Breaking changes in the mempool package: [mempool] #2659 Mempool now an interface old Mempool renamed to CListMempool NewMempool renamed to NewCListMempool Option renamed to CListOption MempoolReactor renamed to Reactor NewMempoolReactor renamed to NewReactor unexpose TxID method TxInfo.PeerID renamed to SenderID unexpose MempoolReactor.Mempool Breaking changes in the state package: [state] #2659 Mempool interface moved to mempool package MockMempool moved to top-level mock package and renamed to Mempool Non Breaking changes in the node package: [node] #2659 Add Mempool method, which allows you to access mempool ## Commits * move Mempool interface into mempool package Refs #2659 Breaking changes in the mempool package: - Mempool now an interface - old Mempool renamed to CListMempool Breaking changes to state package: - MockMempool moved to mempool/mock package and renamed to Mempool - Mempool interface moved to mempool package * assert CListMempool impl Mempool * gofmt code * rename MempoolReactor to Reactor - combine everything into one interface - rename TxInfo.PeerID to TxInfo.SenderID - unexpose MempoolReactor.Mempool * move mempool mock into top-level mock package * add a fixme TxsFront should not be a part of the Mempool interface because it leaks implementation details. Instead, we need to come up with general interface for querying the mempool so the MempoolReactor can fetch and broadcast txs to peers. * change node#Mempool to return interface * save commit = new reactor arch * Revert "save commit = new reactor arch" This reverts commit 1bfceacd9d65a720574683a7f22771e69af9af4d. * require CListMempool in mempool.Reactor * add two changelog entries * fixes after my own review * quote interfaces, structs and functions * fixes after Ismail's review * make node's mempool an interface * make InitWAL/CloseWAL methods a part of Mempool interface * fix merge conflicts * make node's mempool an interface
6 years ago
mempool: move interface into mempool package (#3524) ## Description Refs #2659 Breaking changes in the mempool package: [mempool] #2659 Mempool now an interface old Mempool renamed to CListMempool NewMempool renamed to NewCListMempool Option renamed to CListOption MempoolReactor renamed to Reactor NewMempoolReactor renamed to NewReactor unexpose TxID method TxInfo.PeerID renamed to SenderID unexpose MempoolReactor.Mempool Breaking changes in the state package: [state] #2659 Mempool interface moved to mempool package MockMempool moved to top-level mock package and renamed to Mempool Non Breaking changes in the node package: [node] #2659 Add Mempool method, which allows you to access mempool ## Commits * move Mempool interface into mempool package Refs #2659 Breaking changes in the mempool package: - Mempool now an interface - old Mempool renamed to CListMempool Breaking changes to state package: - MockMempool moved to mempool/mock package and renamed to Mempool - Mempool interface moved to mempool package * assert CListMempool impl Mempool * gofmt code * rename MempoolReactor to Reactor - combine everything into one interface - rename TxInfo.PeerID to TxInfo.SenderID - unexpose MempoolReactor.Mempool * move mempool mock into top-level mock package * add a fixme TxsFront should not be a part of the Mempool interface because it leaks implementation details. Instead, we need to come up with general interface for querying the mempool so the MempoolReactor can fetch and broadcast txs to peers. * change node#Mempool to return interface * save commit = new reactor arch * Revert "save commit = new reactor arch" This reverts commit 1bfceacd9d65a720574683a7f22771e69af9af4d. * require CListMempool in mempool.Reactor * add two changelog entries * fixes after my own review * quote interfaces, structs and functions * fixes after Ismail's review * make node's mempool an interface * make InitWAL/CloseWAL methods a part of Mempool interface * fix merge conflicts * make node's mempool an interface
6 years ago
8 years ago
lint: Enable Golint (#4212) * Fix many golint errors * Fix golint errors in the 'lite' package * Don't export Pool.store * Fix typo * Revert unwanted changes * Fix errors in counter package * Fix linter errors in kvstore package * Fix linter error in example package * Fix error in tests package * Fix linter errors in v2 package * Fix linter errors in consensus package * Fix linter errors in evidence package * Fix linter error in fail package * Fix linter errors in query package * Fix linter errors in core package * Fix linter errors in node package * Fix linter errors in mempool package * Fix linter error in conn package * Fix linter errors in pex package * Rename PEXReactor export to Reactor * Fix linter errors in trust package * Fix linter errors in upnp package * Fix linter errors in p2p package * Fix linter errors in proxy package * Fix linter errors in mock_test package * Fix linter error in client_test package * Fix linter errors in coretypes package * Fix linter errors in coregrpc package * Fix linter errors in rpcserver package * Fix linter errors in rpctypes package * Fix linter errors in rpctest package * Fix linter error in json2wal script * Fix linter error in wal2json script * Fix linter errors in kv package * Fix linter error in state package * Fix linter error in grpc_client * Fix linter errors in types package * Fix linter error in version package * Fix remaining errors * Address review comments * Fix broken tests * Reconcile package coregrpc * Fix golangci bot error * Fix new golint errors * Fix broken reference * Enable golint linter * minor changes to bring golint into line * fix failing test * fix pex reactor naming * address PR comments
5 years ago
lint: Enable Golint (#4212) * Fix many golint errors * Fix golint errors in the 'lite' package * Don't export Pool.store * Fix typo * Revert unwanted changes * Fix errors in counter package * Fix linter errors in kvstore package * Fix linter error in example package * Fix error in tests package * Fix linter errors in v2 package * Fix linter errors in consensus package * Fix linter errors in evidence package * Fix linter error in fail package * Fix linter errors in query package * Fix linter errors in core package * Fix linter errors in node package * Fix linter errors in mempool package * Fix linter error in conn package * Fix linter errors in pex package * Rename PEXReactor export to Reactor * Fix linter errors in trust package * Fix linter errors in upnp package * Fix linter errors in p2p package * Fix linter errors in proxy package * Fix linter errors in mock_test package * Fix linter error in client_test package * Fix linter errors in coretypes package * Fix linter errors in coregrpc package * Fix linter errors in rpcserver package * Fix linter errors in rpctypes package * Fix linter errors in rpctest package * Fix linter error in json2wal script * Fix linter error in wal2json script * Fix linter errors in kv package * Fix linter error in state package * Fix linter error in grpc_client * Fix linter errors in types package * Fix linter error in version package * Fix remaining errors * Address review comments * Fix broken tests * Reconcile package coregrpc * Fix golangci bot error * Fix new golint errors * Fix broken reference * Enable golint linter * minor changes to bring golint into line * fix failing test * fix pex reactor naming * address PR comments
5 years ago
lint: Enable Golint (#4212) * Fix many golint errors * Fix golint errors in the 'lite' package * Don't export Pool.store * Fix typo * Revert unwanted changes * Fix errors in counter package * Fix linter errors in kvstore package * Fix linter error in example package * Fix error in tests package * Fix linter errors in v2 package * Fix linter errors in consensus package * Fix linter errors in evidence package * Fix linter error in fail package * Fix linter errors in query package * Fix linter errors in core package * Fix linter errors in node package * Fix linter errors in mempool package * Fix linter error in conn package * Fix linter errors in pex package * Rename PEXReactor export to Reactor * Fix linter errors in trust package * Fix linter errors in upnp package * Fix linter errors in p2p package * Fix linter errors in proxy package * Fix linter errors in mock_test package * Fix linter error in client_test package * Fix linter errors in coretypes package * Fix linter errors in coregrpc package * Fix linter errors in rpcserver package * Fix linter errors in rpctypes package * Fix linter errors in rpctest package * Fix linter error in json2wal script * Fix linter error in wal2json script * Fix linter errors in kv package * Fix linter error in state package * Fix linter error in grpc_client * Fix linter errors in types package * Fix linter error in version package * Fix remaining errors * Address review comments * Fix broken tests * Reconcile package coregrpc * Fix golangci bot error * Fix new golint errors * Fix broken reference * Enable golint linter * minor changes to bring golint into line * fix failing test * fix pex reactor naming * address PR comments
5 years ago
lint: Enable Golint (#4212) * Fix many golint errors * Fix golint errors in the 'lite' package * Don't export Pool.store * Fix typo * Revert unwanted changes * Fix errors in counter package * Fix linter errors in kvstore package * Fix linter error in example package * Fix error in tests package * Fix linter errors in v2 package * Fix linter errors in consensus package * Fix linter errors in evidence package * Fix linter error in fail package * Fix linter errors in query package * Fix linter errors in core package * Fix linter errors in node package * Fix linter errors in mempool package * Fix linter error in conn package * Fix linter errors in pex package * Rename PEXReactor export to Reactor * Fix linter errors in trust package * Fix linter errors in upnp package * Fix linter errors in p2p package * Fix linter errors in proxy package * Fix linter errors in mock_test package * Fix linter error in client_test package * Fix linter errors in coretypes package * Fix linter errors in coregrpc package * Fix linter errors in rpcserver package * Fix linter errors in rpctypes package * Fix linter errors in rpctest package * Fix linter error in json2wal script * Fix linter error in wal2json script * Fix linter errors in kv package * Fix linter error in state package * Fix linter error in grpc_client * Fix linter errors in types package * Fix linter error in version package * Fix remaining errors * Address review comments * Fix broken tests * Reconcile package coregrpc * Fix golangci bot error * Fix new golint errors * Fix broken reference * Enable golint linter * minor changes to bring golint into line * fix failing test * fix pex reactor naming * address PR comments
5 years ago
lint: Enable Golint (#4212) * Fix many golint errors * Fix golint errors in the 'lite' package * Don't export Pool.store * Fix typo * Revert unwanted changes * Fix errors in counter package * Fix linter errors in kvstore package * Fix linter error in example package * Fix error in tests package * Fix linter errors in v2 package * Fix linter errors in consensus package * Fix linter errors in evidence package * Fix linter error in fail package * Fix linter errors in query package * Fix linter errors in core package * Fix linter errors in node package * Fix linter errors in mempool package * Fix linter error in conn package * Fix linter errors in pex package * Rename PEXReactor export to Reactor * Fix linter errors in trust package * Fix linter errors in upnp package * Fix linter errors in p2p package * Fix linter errors in proxy package * Fix linter errors in mock_test package * Fix linter error in client_test package * Fix linter errors in coretypes package * Fix linter errors in coregrpc package * Fix linter errors in rpcserver package * Fix linter errors in rpctypes package * Fix linter errors in rpctest package * Fix linter error in json2wal script * Fix linter error in wal2json script * Fix linter errors in kv package * Fix linter error in state package * Fix linter error in grpc_client * Fix linter errors in types package * Fix linter error in version package * Fix remaining errors * Address review comments * Fix broken tests * Reconcile package coregrpc * Fix golangci bot error * Fix new golint errors * Fix broken reference * Enable golint linter * minor changes to bring golint into line * fix failing test * fix pex reactor naming * address PR comments
5 years ago
8 years ago
lint: Enable Golint (#4212) * Fix many golint errors * Fix golint errors in the 'lite' package * Don't export Pool.store * Fix typo * Revert unwanted changes * Fix errors in counter package * Fix linter errors in kvstore package * Fix linter error in example package * Fix error in tests package * Fix linter errors in v2 package * Fix linter errors in consensus package * Fix linter errors in evidence package * Fix linter error in fail package * Fix linter errors in query package * Fix linter errors in core package * Fix linter errors in node package * Fix linter errors in mempool package * Fix linter error in conn package * Fix linter errors in pex package * Rename PEXReactor export to Reactor * Fix linter errors in trust package * Fix linter errors in upnp package * Fix linter errors in p2p package * Fix linter errors in proxy package * Fix linter errors in mock_test package * Fix linter error in client_test package * Fix linter errors in coretypes package * Fix linter errors in coregrpc package * Fix linter errors in rpcserver package * Fix linter errors in rpctypes package * Fix linter errors in rpctest package * Fix linter error in json2wal script * Fix linter error in wal2json script * Fix linter errors in kv package * Fix linter error in state package * Fix linter error in grpc_client * Fix linter errors in types package * Fix linter error in version package * Fix remaining errors * Address review comments * Fix broken tests * Reconcile package coregrpc * Fix golangci bot error * Fix new golint errors * Fix broken reference * Enable golint linter * minor changes to bring golint into line * fix failing test * fix pex reactor naming * address PR comments
5 years ago
lint: Enable Golint (#4212) * Fix many golint errors * Fix golint errors in the 'lite' package * Don't export Pool.store * Fix typo * Revert unwanted changes * Fix errors in counter package * Fix linter errors in kvstore package * Fix linter error in example package * Fix error in tests package * Fix linter errors in v2 package * Fix linter errors in consensus package * Fix linter errors in evidence package * Fix linter error in fail package * Fix linter errors in query package * Fix linter errors in core package * Fix linter errors in node package * Fix linter errors in mempool package * Fix linter error in conn package * Fix linter errors in pex package * Rename PEXReactor export to Reactor * Fix linter errors in trust package * Fix linter errors in upnp package * Fix linter errors in p2p package * Fix linter errors in proxy package * Fix linter errors in mock_test package * Fix linter error in client_test package * Fix linter errors in coretypes package * Fix linter errors in coregrpc package * Fix linter errors in rpcserver package * Fix linter errors in rpctypes package * Fix linter errors in rpctest package * Fix linter error in json2wal script * Fix linter error in wal2json script * Fix linter errors in kv package * Fix linter error in state package * Fix linter error in grpc_client * Fix linter errors in types package * Fix linter error in version package * Fix remaining errors * Address review comments * Fix broken tests * Reconcile package coregrpc * Fix golangci bot error * Fix new golint errors * Fix broken reference * Enable golint linter * minor changes to bring golint into line * fix failing test * fix pex reactor naming * address PR comments
5 years ago
lint: Enable Golint (#4212) * Fix many golint errors * Fix golint errors in the 'lite' package * Don't export Pool.store * Fix typo * Revert unwanted changes * Fix errors in counter package * Fix linter errors in kvstore package * Fix linter error in example package * Fix error in tests package * Fix linter errors in v2 package * Fix linter errors in consensus package * Fix linter errors in evidence package * Fix linter error in fail package * Fix linter errors in query package * Fix linter errors in core package * Fix linter errors in node package * Fix linter errors in mempool package * Fix linter error in conn package * Fix linter errors in pex package * Rename PEXReactor export to Reactor * Fix linter errors in trust package * Fix linter errors in upnp package * Fix linter errors in p2p package * Fix linter errors in proxy package * Fix linter errors in mock_test package * Fix linter error in client_test package * Fix linter errors in coretypes package * Fix linter errors in coregrpc package * Fix linter errors in rpcserver package * Fix linter errors in rpctypes package * Fix linter errors in rpctest package * Fix linter error in json2wal script * Fix linter error in wal2json script * Fix linter errors in kv package * Fix linter error in state package * Fix linter error in grpc_client * Fix linter errors in types package * Fix linter error in version package * Fix remaining errors * Address review comments * Fix broken tests * Reconcile package coregrpc * Fix golangci bot error * Fix new golint errors * Fix broken reference * Enable golint linter * minor changes to bring golint into line * fix failing test * fix pex reactor naming * address PR comments
5 years ago
lint: Enable Golint (#4212) * Fix many golint errors * Fix golint errors in the 'lite' package * Don't export Pool.store * Fix typo * Revert unwanted changes * Fix errors in counter package * Fix linter errors in kvstore package * Fix linter error in example package * Fix error in tests package * Fix linter errors in v2 package * Fix linter errors in consensus package * Fix linter errors in evidence package * Fix linter error in fail package * Fix linter errors in query package * Fix linter errors in core package * Fix linter errors in node package * Fix linter errors in mempool package * Fix linter error in conn package * Fix linter errors in pex package * Rename PEXReactor export to Reactor * Fix linter errors in trust package * Fix linter errors in upnp package * Fix linter errors in p2p package * Fix linter errors in proxy package * Fix linter errors in mock_test package * Fix linter error in client_test package * Fix linter errors in coretypes package * Fix linter errors in coregrpc package * Fix linter errors in rpcserver package * Fix linter errors in rpctypes package * Fix linter errors in rpctest package * Fix linter error in json2wal script * Fix linter error in wal2json script * Fix linter errors in kv package * Fix linter error in state package * Fix linter error in grpc_client * Fix linter errors in types package * Fix linter error in version package * Fix remaining errors * Address review comments * Fix broken tests * Reconcile package coregrpc * Fix golangci bot error * Fix new golint errors * Fix broken reference * Enable golint linter * minor changes to bring golint into line * fix failing test * fix pex reactor naming * address PR comments
5 years ago
  1. package v0
  2. import (
  3. "context"
  4. "crypto/rand"
  5. "encoding/binary"
  6. "fmt"
  7. mrand "math/rand"
  8. "os"
  9. "testing"
  10. "time"
  11. "github.com/gogo/protobuf/proto"
  12. gogotypes "github.com/gogo/protobuf/types"
  13. "github.com/stretchr/testify/assert"
  14. "github.com/stretchr/testify/require"
  15. "github.com/tendermint/tendermint/abci/example/counter"
  16. "github.com/tendermint/tendermint/abci/example/kvstore"
  17. abciserver "github.com/tendermint/tendermint/abci/server"
  18. abci "github.com/tendermint/tendermint/abci/types"
  19. cfg "github.com/tendermint/tendermint/config"
  20. "github.com/tendermint/tendermint/internal/mempool"
  21. "github.com/tendermint/tendermint/libs/log"
  22. tmrand "github.com/tendermint/tendermint/libs/rand"
  23. "github.com/tendermint/tendermint/libs/service"
  24. pubmempool "github.com/tendermint/tendermint/pkg/mempool"
  25. "github.com/tendermint/tendermint/proxy"
  26. "github.com/tendermint/tendermint/types"
  27. )
  28. // A cleanupFunc cleans up any config / test files created for a particular
  29. // test.
  30. type cleanupFunc func()
  31. func newMempoolWithApp(cc proxy.ClientCreator) (*CListMempool, cleanupFunc) {
  32. return newMempoolWithAppAndConfig(cc, cfg.ResetTestRoot("mempool_test"))
  33. }
  34. func newMempoolWithAppAndConfig(cc proxy.ClientCreator, config *cfg.Config) (*CListMempool, cleanupFunc) {
  35. appConnMem, _ := cc.NewABCIClient()
  36. appConnMem.SetLogger(log.TestingLogger().With("module", "abci-client", "connection", "mempool"))
  37. err := appConnMem.Start()
  38. if err != nil {
  39. panic(err)
  40. }
  41. mp := NewCListMempool(config.Mempool, appConnMem, 0)
  42. mp.SetLogger(log.TestingLogger())
  43. return mp, func() { os.RemoveAll(config.RootDir) }
  44. }
  45. func ensureNoFire(t *testing.T, ch <-chan struct{}, timeoutMS int) {
  46. timer := time.NewTimer(time.Duration(timeoutMS) * time.Millisecond)
  47. select {
  48. case <-ch:
  49. t.Fatal("Expected not to fire")
  50. case <-timer.C:
  51. }
  52. }
  53. func ensureFire(t *testing.T, ch <-chan struct{}, timeoutMS int) {
  54. timer := time.NewTimer(time.Duration(timeoutMS) * time.Millisecond)
  55. select {
  56. case <-ch:
  57. case <-timer.C:
  58. t.Fatal("Expected to fire")
  59. }
  60. }
  61. func checkTxs(t *testing.T, mp mempool.Mempool, count int, peerID uint16) types.Txs {
  62. txs := make(types.Txs, count)
  63. txInfo := mempool.TxInfo{SenderID: peerID}
  64. for i := 0; i < count; i++ {
  65. txBytes := make([]byte, 20)
  66. txs[i] = txBytes
  67. _, err := rand.Read(txBytes)
  68. if err != nil {
  69. t.Error(err)
  70. }
  71. if err := mp.CheckTx(context.Background(), txBytes, nil, txInfo); err != nil {
  72. // Skip invalid txs.
  73. // TestMempoolFilters will fail otherwise. It asserts a number of txs
  74. // returned.
  75. if pubmempool.IsPreCheckError(err) {
  76. continue
  77. }
  78. t.Fatalf("CheckTx failed: %v while checking #%d tx", err, i)
  79. }
  80. }
  81. return txs
  82. }
  83. func TestReapMaxBytesMaxGas(t *testing.T) {
  84. app := kvstore.NewApplication()
  85. cc := proxy.NewLocalClientCreator(app)
  86. mp, cleanup := newMempoolWithApp(cc)
  87. defer cleanup()
  88. // Ensure gas calculation behaves as expected
  89. checkTxs(t, mp, 1, mempool.UnknownPeerID)
  90. tx0 := mp.TxsFront().Value.(*mempoolTx)
  91. // assert that kv store has gas wanted = 1.
  92. require.Equal(t, app.CheckTx(abci.RequestCheckTx{Tx: tx0.tx}).GasWanted, int64(1), "KVStore had a gas value neq to 1")
  93. require.Equal(t, tx0.gasWanted, int64(1), "transactions gas was set incorrectly")
  94. // ensure each tx is 20 bytes long
  95. require.Equal(t, len(tx0.tx), 20, "Tx is longer than 20 bytes")
  96. mp.Flush()
  97. // each table driven test creates numTxsToCreate txs with checkTx, and at the end clears all remaining txs.
  98. // each tx has 20 bytes
  99. tests := []struct {
  100. numTxsToCreate int
  101. maxBytes int64
  102. maxGas int64
  103. expectedNumTxs int
  104. }{
  105. {20, -1, -1, 20},
  106. {20, -1, 0, 0},
  107. {20, -1, 10, 10},
  108. {20, -1, 30, 20},
  109. {20, 0, -1, 0},
  110. {20, 0, 10, 0},
  111. {20, 10, 10, 0},
  112. {20, 24, 10, 1},
  113. {20, 240, 5, 5},
  114. {20, 240, -1, 10},
  115. {20, 240, 10, 10},
  116. {20, 240, 15, 10},
  117. {20, 20000, -1, 20},
  118. {20, 20000, 5, 5},
  119. {20, 20000, 30, 20},
  120. }
  121. for tcIndex, tt := range tests {
  122. checkTxs(t, mp, tt.numTxsToCreate, mempool.UnknownPeerID)
  123. got := mp.ReapMaxBytesMaxGas(tt.maxBytes, tt.maxGas)
  124. assert.Equal(t, tt.expectedNumTxs, len(got), "Got %d txs, expected %d, tc #%d",
  125. len(got), tt.expectedNumTxs, tcIndex)
  126. mp.Flush()
  127. }
  128. }
  129. func TestMempoolFilters(t *testing.T) {
  130. app := kvstore.NewApplication()
  131. cc := proxy.NewLocalClientCreator(app)
  132. mp, cleanup := newMempoolWithApp(cc)
  133. defer cleanup()
  134. emptyTxArr := []types.Tx{[]byte{}}
  135. nopPreFilter := func(tx types.Tx) error { return nil }
  136. nopPostFilter := func(tx types.Tx, res *abci.ResponseCheckTx) error { return nil }
  137. // each table driven test creates numTxsToCreate txs with checkTx, and at the end clears all remaining txs.
  138. // each tx has 20 bytes
  139. tests := []struct {
  140. numTxsToCreate int
  141. preFilter mempool.PreCheckFunc
  142. postFilter mempool.PostCheckFunc
  143. expectedNumTxs int
  144. }{
  145. {10, nopPreFilter, nopPostFilter, 10},
  146. {10, mempool.PreCheckMaxBytes(10), nopPostFilter, 0},
  147. {10, mempool.PreCheckMaxBytes(22), nopPostFilter, 10},
  148. {10, nopPreFilter, mempool.PostCheckMaxGas(-1), 10},
  149. {10, nopPreFilter, mempool.PostCheckMaxGas(0), 0},
  150. {10, nopPreFilter, mempool.PostCheckMaxGas(1), 10},
  151. {10, nopPreFilter, mempool.PostCheckMaxGas(3000), 10},
  152. {10, mempool.PreCheckMaxBytes(10), mempool.PostCheckMaxGas(20), 0},
  153. {10, mempool.PreCheckMaxBytes(30), mempool.PostCheckMaxGas(20), 10},
  154. {10, mempool.PreCheckMaxBytes(22), mempool.PostCheckMaxGas(1), 10},
  155. {10, mempool.PreCheckMaxBytes(22), mempool.PostCheckMaxGas(0), 0},
  156. }
  157. for tcIndex, tt := range tests {
  158. err := mp.Update(1, emptyTxArr, abciResponses(len(emptyTxArr), abci.CodeTypeOK), tt.preFilter, tt.postFilter)
  159. require.NoError(t, err)
  160. checkTxs(t, mp, tt.numTxsToCreate, mempool.UnknownPeerID)
  161. require.Equal(t, tt.expectedNumTxs, mp.Size(), "mempool had the incorrect size, on test case %d", tcIndex)
  162. mp.Flush()
  163. }
  164. }
  165. func TestMempoolUpdate(t *testing.T) {
  166. app := kvstore.NewApplication()
  167. cc := proxy.NewLocalClientCreator(app)
  168. mp, cleanup := newMempoolWithApp(cc)
  169. defer cleanup()
  170. // 1. Adds valid txs to the cache
  171. {
  172. err := mp.Update(1, []types.Tx{[]byte{0x01}}, abciResponses(1, abci.CodeTypeOK), nil, nil)
  173. require.NoError(t, err)
  174. err = mp.CheckTx(context.Background(), []byte{0x01}, nil, mempool.TxInfo{})
  175. require.NoError(t, err)
  176. }
  177. // 2. Removes valid txs from the mempool
  178. {
  179. err := mp.CheckTx(context.Background(), []byte{0x02}, nil, mempool.TxInfo{})
  180. require.NoError(t, err)
  181. err = mp.Update(1, []types.Tx{[]byte{0x02}}, abciResponses(1, abci.CodeTypeOK), nil, nil)
  182. require.NoError(t, err)
  183. assert.Zero(t, mp.Size())
  184. }
  185. // 3. Removes invalid transactions from the cache and the mempool (if present)
  186. {
  187. err := mp.CheckTx(context.Background(), []byte{0x03}, nil, mempool.TxInfo{})
  188. require.NoError(t, err)
  189. err = mp.Update(1, []types.Tx{[]byte{0x03}}, abciResponses(1, 1), nil, nil)
  190. require.NoError(t, err)
  191. assert.Zero(t, mp.Size())
  192. err = mp.CheckTx(context.Background(), []byte{0x03}, nil, mempool.TxInfo{})
  193. require.NoError(t, err)
  194. }
  195. }
  196. func TestMempool_KeepInvalidTxsInCache(t *testing.T) {
  197. app := counter.NewApplication(true)
  198. cc := proxy.NewLocalClientCreator(app)
  199. wcfg := cfg.DefaultConfig()
  200. wcfg.Mempool.KeepInvalidTxsInCache = true
  201. mp, cleanup := newMempoolWithAppAndConfig(cc, wcfg)
  202. defer cleanup()
  203. // 1. An invalid transaction must remain in the cache after Update
  204. {
  205. a := make([]byte, 8)
  206. binary.BigEndian.PutUint64(a, 0)
  207. b := make([]byte, 8)
  208. binary.BigEndian.PutUint64(b, 1)
  209. err := mp.CheckTx(context.Background(), b, nil, mempool.TxInfo{})
  210. require.NoError(t, err)
  211. // simulate new block
  212. _ = app.DeliverTx(abci.RequestDeliverTx{Tx: a})
  213. _ = app.DeliverTx(abci.RequestDeliverTx{Tx: b})
  214. err = mp.Update(1, []types.Tx{a, b},
  215. []*abci.ResponseDeliverTx{{Code: abci.CodeTypeOK}, {Code: 2}}, nil, nil)
  216. require.NoError(t, err)
  217. // a must be added to the cache
  218. err = mp.CheckTx(context.Background(), a, nil, mempool.TxInfo{})
  219. require.NoError(t, err)
  220. // b must remain in the cache
  221. err = mp.CheckTx(context.Background(), b, nil, mempool.TxInfo{})
  222. require.NoError(t, err)
  223. }
  224. // 2. An invalid transaction must remain in the cache
  225. {
  226. a := make([]byte, 8)
  227. binary.BigEndian.PutUint64(a, 0)
  228. // remove a from the cache to test (2)
  229. mp.cache.Remove(a)
  230. err := mp.CheckTx(context.Background(), a, nil, mempool.TxInfo{})
  231. require.NoError(t, err)
  232. }
  233. }
  234. func TestTxsAvailable(t *testing.T) {
  235. app := kvstore.NewApplication()
  236. cc := proxy.NewLocalClientCreator(app)
  237. mp, cleanup := newMempoolWithApp(cc)
  238. defer cleanup()
  239. mp.EnableTxsAvailable()
  240. timeoutMS := 500
  241. // with no txs, it shouldnt fire
  242. ensureNoFire(t, mp.TxsAvailable(), timeoutMS)
  243. // send a bunch of txs, it should only fire once
  244. txs := checkTxs(t, mp, 100, mempool.UnknownPeerID)
  245. ensureFire(t, mp.TxsAvailable(), timeoutMS)
  246. ensureNoFire(t, mp.TxsAvailable(), timeoutMS)
  247. // call update with half the txs.
  248. // it should fire once now for the new height
  249. // since there are still txs left
  250. committedTxs, txs := txs[:50], txs[50:]
  251. if err := mp.Update(1, committedTxs, abciResponses(len(committedTxs), abci.CodeTypeOK), nil, nil); err != nil {
  252. t.Error(err)
  253. }
  254. ensureFire(t, mp.TxsAvailable(), timeoutMS)
  255. ensureNoFire(t, mp.TxsAvailable(), timeoutMS)
  256. // send a bunch more txs. we already fired for this height so it shouldnt fire again
  257. moreTxs := checkTxs(t, mp, 50, mempool.UnknownPeerID)
  258. ensureNoFire(t, mp.TxsAvailable(), timeoutMS)
  259. // now call update with all the txs. it should not fire as there are no txs left
  260. committedTxs = append(txs, moreTxs...) //nolint: gocritic
  261. if err := mp.Update(2, committedTxs, abciResponses(len(committedTxs), abci.CodeTypeOK), nil, nil); err != nil {
  262. t.Error(err)
  263. }
  264. ensureNoFire(t, mp.TxsAvailable(), timeoutMS)
  265. // send a bunch more txs, it should only fire once
  266. checkTxs(t, mp, 100, mempool.UnknownPeerID)
  267. ensureFire(t, mp.TxsAvailable(), timeoutMS)
  268. ensureNoFire(t, mp.TxsAvailable(), timeoutMS)
  269. }
  270. func TestSerialReap(t *testing.T) {
  271. app := counter.NewApplication(true)
  272. cc := proxy.NewLocalClientCreator(app)
  273. mp, cleanup := newMempoolWithApp(cc)
  274. defer cleanup()
  275. appConnCon, _ := cc.NewABCIClient()
  276. appConnCon.SetLogger(log.TestingLogger().With("module", "abci-client", "connection", "consensus"))
  277. err := appConnCon.Start()
  278. require.Nil(t, err)
  279. cacheMap := make(map[string]struct{})
  280. deliverTxsRange := func(start, end int) {
  281. // Deliver some txs.
  282. for i := start; i < end; i++ {
  283. // This will succeed
  284. txBytes := make([]byte, 8)
  285. binary.BigEndian.PutUint64(txBytes, uint64(i))
  286. err := mp.CheckTx(context.Background(), txBytes, nil, mempool.TxInfo{})
  287. _, cached := cacheMap[string(txBytes)]
  288. if cached {
  289. require.NotNil(t, err, "expected error for cached tx")
  290. } else {
  291. require.Nil(t, err, "expected no err for uncached tx")
  292. }
  293. cacheMap[string(txBytes)] = struct{}{}
  294. // Duplicates are cached and should return error
  295. err = mp.CheckTx(context.Background(), txBytes, nil, mempool.TxInfo{})
  296. require.NotNil(t, err, "Expected error after CheckTx on duplicated tx")
  297. }
  298. }
  299. reapCheck := func(exp int) {
  300. txs := mp.ReapMaxBytesMaxGas(-1, -1)
  301. require.Equal(t, len(txs), exp, fmt.Sprintf("Expected to reap %v txs but got %v", exp, len(txs)))
  302. }
  303. updateRange := func(start, end int) {
  304. txs := make([]types.Tx, 0)
  305. for i := start; i < end; i++ {
  306. txBytes := make([]byte, 8)
  307. binary.BigEndian.PutUint64(txBytes, uint64(i))
  308. txs = append(txs, txBytes)
  309. }
  310. if err := mp.Update(0, txs, abciResponses(len(txs), abci.CodeTypeOK), nil, nil); err != nil {
  311. t.Error(err)
  312. }
  313. }
  314. commitRange := func(start, end int) {
  315. ctx := context.Background()
  316. // Deliver some txs.
  317. for i := start; i < end; i++ {
  318. txBytes := make([]byte, 8)
  319. binary.BigEndian.PutUint64(txBytes, uint64(i))
  320. res, err := appConnCon.DeliverTxSync(ctx, abci.RequestDeliverTx{Tx: txBytes})
  321. if err != nil {
  322. t.Errorf("client error committing tx: %v", err)
  323. }
  324. if res.IsErr() {
  325. t.Errorf("error committing tx. Code:%v result:%X log:%v",
  326. res.Code, res.Data, res.Log)
  327. }
  328. }
  329. res, err := appConnCon.CommitSync(ctx)
  330. if err != nil {
  331. t.Errorf("client error committing: %v", err)
  332. }
  333. if len(res.Data) != 8 {
  334. t.Errorf("error committing. Hash:%X", res.Data)
  335. }
  336. }
  337. //----------------------------------------
  338. // Deliver some txs.
  339. deliverTxsRange(0, 100)
  340. // Reap the txs.
  341. reapCheck(100)
  342. // Reap again. We should get the same amount
  343. reapCheck(100)
  344. // Deliver 0 to 999, we should reap 900 new txs
  345. // because 100 were already counted.
  346. deliverTxsRange(0, 1000)
  347. // Reap the txs.
  348. reapCheck(1000)
  349. // Reap again. We should get the same amount
  350. reapCheck(1000)
  351. // Commit from the conensus AppConn
  352. commitRange(0, 500)
  353. updateRange(0, 500)
  354. // We should have 500 left.
  355. reapCheck(500)
  356. // Deliver 100 invalid txs and 100 valid txs
  357. deliverTxsRange(900, 1100)
  358. // We should have 600 now.
  359. reapCheck(600)
  360. }
  361. func TestMempool_CheckTxChecksTxSize(t *testing.T) {
  362. app := kvstore.NewApplication()
  363. cc := proxy.NewLocalClientCreator(app)
  364. mempl, cleanup := newMempoolWithApp(cc)
  365. defer cleanup()
  366. maxTxSize := mempl.config.MaxTxBytes
  367. testCases := []struct {
  368. len int
  369. err bool
  370. }{
  371. // check small txs. no error
  372. 0: {10, false},
  373. 1: {1000, false},
  374. 2: {1000000, false},
  375. // check around maxTxSize
  376. 3: {maxTxSize - 1, false},
  377. 4: {maxTxSize, false},
  378. 5: {maxTxSize + 1, true},
  379. }
  380. for i, testCase := range testCases {
  381. caseString := fmt.Sprintf("case %d, len %d", i, testCase.len)
  382. tx := tmrand.Bytes(testCase.len)
  383. err := mempl.CheckTx(context.Background(), tx, nil, mempool.TxInfo{})
  384. bv := gogotypes.BytesValue{Value: tx}
  385. bz, err2 := bv.Marshal()
  386. require.NoError(t, err2)
  387. require.Equal(t, len(bz), proto.Size(&bv), caseString)
  388. if !testCase.err {
  389. require.NoError(t, err, caseString)
  390. } else {
  391. require.Equal(t, err, pubmempool.ErrTxTooLarge{
  392. Max: maxTxSize,
  393. Actual: testCase.len,
  394. }, caseString)
  395. }
  396. }
  397. }
  398. func TestMempoolTxsBytes(t *testing.T) {
  399. app := kvstore.NewApplication()
  400. cc := proxy.NewLocalClientCreator(app)
  401. config := cfg.ResetTestRoot("mempool_test")
  402. config.Mempool.MaxTxsBytes = 10
  403. mp, cleanup := newMempoolWithAppAndConfig(cc, config)
  404. defer cleanup()
  405. // 1. zero by default
  406. assert.EqualValues(t, 0, mp.SizeBytes())
  407. // 2. len(tx) after CheckTx
  408. err := mp.CheckTx(context.Background(), []byte{0x01}, nil, mempool.TxInfo{})
  409. require.NoError(t, err)
  410. assert.EqualValues(t, 1, mp.SizeBytes())
  411. // 3. zero again after tx is removed by Update
  412. err = mp.Update(1, []types.Tx{[]byte{0x01}}, abciResponses(1, abci.CodeTypeOK), nil, nil)
  413. require.NoError(t, err)
  414. assert.EqualValues(t, 0, mp.SizeBytes())
  415. // 4. zero after Flush
  416. err = mp.CheckTx(context.Background(), []byte{0x02, 0x03}, nil, mempool.TxInfo{})
  417. require.NoError(t, err)
  418. assert.EqualValues(t, 2, mp.SizeBytes())
  419. mp.Flush()
  420. assert.EqualValues(t, 0, mp.SizeBytes())
  421. // 5. ErrMempoolIsFull is returned when/if MaxTxsBytes limit is reached.
  422. err = mp.CheckTx(
  423. context.Background(),
  424. []byte{0x04, 0x04, 0x04, 0x04, 0x04, 0x04, 0x04, 0x04, 0x04, 0x04},
  425. nil,
  426. mempool.TxInfo{},
  427. )
  428. require.NoError(t, err)
  429. err = mp.CheckTx(context.Background(), []byte{0x05}, nil, mempool.TxInfo{})
  430. if assert.Error(t, err) {
  431. assert.IsType(t, pubmempool.ErrMempoolIsFull{}, err)
  432. }
  433. // 6. zero after tx is rechecked and removed due to not being valid anymore
  434. app2 := counter.NewApplication(true)
  435. cc = proxy.NewLocalClientCreator(app2)
  436. mp, cleanup = newMempoolWithApp(cc)
  437. defer cleanup()
  438. txBytes := make([]byte, 8)
  439. binary.BigEndian.PutUint64(txBytes, uint64(0))
  440. err = mp.CheckTx(context.Background(), txBytes, nil, mempool.TxInfo{})
  441. require.NoError(t, err)
  442. assert.EqualValues(t, 8, mp.SizeBytes())
  443. appConnCon, _ := cc.NewABCIClient()
  444. appConnCon.SetLogger(log.TestingLogger().With("module", "abci-client", "connection", "consensus"))
  445. err = appConnCon.Start()
  446. require.Nil(t, err)
  447. t.Cleanup(func() {
  448. if err := appConnCon.Stop(); err != nil {
  449. t.Error(err)
  450. }
  451. })
  452. ctx := context.Background()
  453. res, err := appConnCon.DeliverTxSync(ctx, abci.RequestDeliverTx{Tx: txBytes})
  454. require.NoError(t, err)
  455. require.EqualValues(t, 0, res.Code)
  456. res2, err := appConnCon.CommitSync(ctx)
  457. require.NoError(t, err)
  458. require.NotEmpty(t, res2.Data)
  459. // Pretend like we committed nothing so txBytes gets rechecked and removed.
  460. err = mp.Update(1, []types.Tx{}, abciResponses(0, abci.CodeTypeOK), nil, nil)
  461. require.NoError(t, err)
  462. assert.EqualValues(t, 0, mp.SizeBytes())
  463. // 7. Test RemoveTxByKey function
  464. err = mp.CheckTx(context.Background(), []byte{0x06}, nil, mempool.TxInfo{})
  465. require.NoError(t, err)
  466. assert.EqualValues(t, 1, mp.SizeBytes())
  467. mp.RemoveTxByKey(mempool.TxKey([]byte{0x07}), true)
  468. assert.EqualValues(t, 1, mp.SizeBytes())
  469. mp.RemoveTxByKey(mempool.TxKey([]byte{0x06}), true)
  470. assert.EqualValues(t, 0, mp.SizeBytes())
  471. }
  472. // This will non-deterministically catch some concurrency failures like
  473. // https://github.com/tendermint/tendermint/issues/3509
  474. // TODO: all of the tests should probably also run using the remote proxy app
  475. // since otherwise we're not actually testing the concurrency of the mempool here!
  476. func TestMempoolRemoteAppConcurrency(t *testing.T) {
  477. sockPath := fmt.Sprintf("unix:///tmp/echo_%v.sock", tmrand.Str(6))
  478. app := kvstore.NewApplication()
  479. cc, server := newRemoteApp(t, sockPath, app)
  480. t.Cleanup(func() {
  481. if err := server.Stop(); err != nil {
  482. t.Error(err)
  483. }
  484. })
  485. config := cfg.ResetTestRoot("mempool_test")
  486. mp, cleanup := newMempoolWithAppAndConfig(cc, config)
  487. defer cleanup()
  488. // generate small number of txs
  489. nTxs := 10
  490. txLen := 200
  491. txs := make([]types.Tx, nTxs)
  492. for i := 0; i < nTxs; i++ {
  493. txs[i] = tmrand.Bytes(txLen)
  494. }
  495. // simulate a group of peers sending them over and over
  496. N := config.Mempool.Size
  497. maxPeers := 5
  498. for i := 0; i < N; i++ {
  499. peerID := mrand.Intn(maxPeers)
  500. txNum := mrand.Intn(nTxs)
  501. tx := txs[txNum]
  502. // this will err with ErrTxInCache many times ...
  503. mp.CheckTx(context.Background(), tx, nil, mempool.TxInfo{SenderID: uint16(peerID)}) //nolint: errcheck // will error
  504. }
  505. err := mp.FlushAppConn()
  506. require.NoError(t, err)
  507. }
  508. // caller must close server
  509. func newRemoteApp(
  510. t *testing.T,
  511. addr string,
  512. app abci.Application,
  513. ) (
  514. clientCreator proxy.ClientCreator,
  515. server service.Service,
  516. ) {
  517. clientCreator = proxy.NewRemoteClientCreator(addr, "socket", true)
  518. // Start server
  519. server = abciserver.NewSocketServer(addr, app)
  520. server.SetLogger(log.TestingLogger().With("module", "abci-server"))
  521. if err := server.Start(); err != nil {
  522. t.Fatalf("Error starting socket server: %v", err.Error())
  523. }
  524. return clientCreator, server
  525. }
  526. func abciResponses(n int, code uint32) []*abci.ResponseDeliverTx {
  527. responses := make([]*abci.ResponseDeliverTx, 0, n)
  528. for i := 0; i < n; i++ {
  529. responses = append(responses, &abci.ResponseDeliverTx{Code: code})
  530. }
  531. return responses
  532. }