state/indexer: reconstruct indexer, move txindex into the indexer package (#6382)

This commit is contained in:
JayT106
2021-04-21 16:37:44 -04:00
committed by GitHub
parent 5bafedff17
commit 43eacd159f
14 changed files with 58 additions and 70 deletions
+1
View File
@@ -16,6 +16,7 @@ Friendly reminder: We have a [bug bounty program](https://hackerone.com/tendermi
- [rpc] \#6019 standardise RPC errors and return the correct status code (@bipulprasad & @cmwaters)
- [rpc] \#6168 Change default sorting to desc for `/tx_search` results (@melekes)
- [cli] \#6282 User must specify the node mode when using `tendermint init` (@cmwaters)
- [state/indexer] \#6382 reconstruct indexer, move txindex into the indexer package (@JayT106)
- Apps
- [ABCI] \#5447 Remove `SetOption` method from `ABCI.Client` interface
+9 -10
View File
@@ -44,9 +44,8 @@ import (
"github.com/tendermint/tendermint/state/indexer"
blockidxkv "github.com/tendermint/tendermint/state/indexer/block/kv"
blockidxnull "github.com/tendermint/tendermint/state/indexer/block/null"
"github.com/tendermint/tendermint/state/txindex"
"github.com/tendermint/tendermint/state/txindex/kv"
"github.com/tendermint/tendermint/state/txindex/null"
"github.com/tendermint/tendermint/state/indexer/tx/kv"
"github.com/tendermint/tendermint/state/indexer/tx/null"
"github.com/tendermint/tendermint/statesync"
"github.com/tendermint/tendermint/store"
"github.com/tendermint/tendermint/types"
@@ -226,9 +225,9 @@ type Node struct {
evidencePool *evidence.Pool // tracking evidence
proxyApp proxy.AppConns // connection to the application
rpcListeners []net.Listener // rpc servers
txIndexer txindex.TxIndexer
txIndexer indexer.TxIndexer
blockIndexer indexer.BlockIndexer
indexerService *txindex.IndexerService
indexerService *indexer.Service
prometheusSrv *http.Server
}
@@ -267,10 +266,10 @@ func createAndStartIndexerService(
dbProvider DBProvider,
eventBus *types.EventBus,
logger log.Logger,
) (*txindex.IndexerService, txindex.TxIndexer, indexer.BlockIndexer, error) {
) (*indexer.Service, indexer.TxIndexer, indexer.BlockIndexer, error) {
var (
txIndexer txindex.TxIndexer
txIndexer indexer.TxIndexer
blockIndexer indexer.BlockIndexer
)
@@ -288,7 +287,7 @@ func createAndStartIndexerService(
blockIndexer = &blockidxnull.BlockerIndexer{}
}
indexerService := txindex.NewIndexerService(txIndexer, blockIndexer, eventBus)
indexerService := indexer.NewIndexerService(txIndexer, blockIndexer, eventBus)
indexerService.SetLogger(logger.With("module", "txindex"))
if err := indexerService.Start(); err != nil {
@@ -1736,7 +1735,7 @@ func (n *Node) Config() *cfg.Config {
}
// TxIndexer returns the Node's TxIndexer.
func (n *Node) TxIndexer() txindex.TxIndexer {
func (n *Node) TxIndexer() indexer.TxIndexer {
return n.txIndexer
}
@@ -1760,7 +1759,7 @@ func (n *Node) NodeInfo() p2p.NodeInfo {
func makeNodeInfo(
config *cfg.Config,
nodeKey p2p.NodeKey,
txIndexer txindex.TxIndexer,
txIndexer indexer.TxIndexer,
genDoc *types.GenesisDoc,
state sm.State,
) (p2p.NodeInfo, error) {
+1 -2
View File
@@ -14,7 +14,6 @@ import (
ctypes "github.com/tendermint/tendermint/rpc/core/types"
sm "github.com/tendermint/tendermint/state"
"github.com/tendermint/tendermint/state/indexer"
"github.com/tendermint/tendermint/state/txindex"
"github.com/tendermint/tendermint/types"
)
@@ -83,7 +82,7 @@ type Environment struct {
// objects
PubKey crypto.PubKey
GenDoc *types.GenesisDoc // cache the genesis structure
TxIndexer txindex.TxIndexer
TxIndexer indexer.TxIndexer
BlockIndexer indexer.BlockIndexer
ConsensusReactor *consensus.Reactor
EventBus *types.EventBus // thread safe
+1 -1
View File
@@ -9,7 +9,7 @@ import (
tmquery "github.com/tendermint/tendermint/libs/pubsub/query"
ctypes "github.com/tendermint/tendermint/rpc/core/types"
rpctypes "github.com/tendermint/tendermint/rpc/jsonrpc/types"
"github.com/tendermint/tendermint/state/txindex/null"
"github.com/tendermint/tendermint/state/indexer/tx/null"
"github.com/tendermint/tendermint/types"
)
-22
View File
@@ -1,22 +0,0 @@
package indexer
import (
"context"
"github.com/tendermint/tendermint/libs/pubsub/query"
"github.com/tendermint/tendermint/types"
)
// BlockIndexer defines an interface contract for indexing block events.
type BlockIndexer interface {
// Has returns true if the given height has been indexed. An error is returned
// upon database query failure.
Has(height int64) (bool, error)
// Index indexes BeginBlock and EndBlock events for a given block by its height.
Index(types.EventDataNewBlockHeader) error
// Search performs a query for block heights that match a given BeginBlock
// and Endblock event search criteria.
Search(ctx context.Context, q *query.Query) ([]int64, error)
}
@@ -1,4 +1,4 @@
package txindex
package indexer
import (
"context"
@@ -6,10 +6,9 @@ import (
abci "github.com/tendermint/tendermint/abci/types"
"github.com/tendermint/tendermint/libs/pubsub/query"
"github.com/tendermint/tendermint/types"
)
// XXX/TODO: These types should be moved to the indexer package.
// TxIndexer interface defines methods to index and search transactions.
type TxIndexer interface {
// AddBatch analyzes, indexes and stores a batch of transactions.
@@ -26,6 +25,20 @@ type TxIndexer interface {
Search(ctx context.Context, q *query.Query) ([]*abci.TxResult, error)
}
// BlockIndexer defines an interface contract for indexing block events.
type BlockIndexer interface {
// Has returns true if the given height has been indexed. An error is returned
// upon database query failure.
Has(height int64) (bool, error)
// Index indexes BeginBlock and EndBlock events for a given block by its height.
Index(types.EventDataNewBlockHeader) error
// Search performs a query for block heights that match a given BeginBlock
// and Endblock event search criteria.
Search(ctx context.Context, q *query.Query) ([]int64, error)
}
// Batch groups together multiple Index operations to be performed at the same time.
// NOTE: Batch is NOT thread-safe and must not be modified after starting its execution.
type Batch struct {
@@ -1,10 +1,9 @@
package txindex
package indexer
import (
"context"
"github.com/tendermint/tendermint/libs/service"
"github.com/tendermint/tendermint/state/indexer"
"github.com/tendermint/tendermint/types"
)
@@ -14,31 +13,31 @@ const (
subscriber = "IndexerService"
)
// IndexerService connects event bus, transaction and block indexers together in
// Service connects event bus, transaction and block indexers together in
// order to index transactions and blocks coming from the event bus.
type IndexerService struct {
type Service struct {
service.BaseService
txIdxr TxIndexer
blockIdxr indexer.BlockIndexer
blockIdxr BlockIndexer
eventBus *types.EventBus
}
// NewIndexerService returns a new service instance.
func NewIndexerService(
txIdxr TxIndexer,
blockIdxr indexer.BlockIndexer,
blockIdxr BlockIndexer,
eventBus *types.EventBus,
) *IndexerService {
) *Service {
is := &IndexerService{txIdxr: txIdxr, blockIdxr: blockIdxr, eventBus: eventBus}
is := &Service{txIdxr: txIdxr, blockIdxr: blockIdxr, eventBus: eventBus}
is.BaseService = *service.NewBaseService(nil, "IndexerService", is)
return is
}
// OnStart implements service.Service by subscribing for all transactions
// and indexing them by events.
func (is *IndexerService) OnStart() error {
func (is *Service) OnStart() error {
// Use SubscribeUnbuffered here to ensure both subscriptions does not get
// canceled due to not pulling messages fast enough. Cause this might
// sometimes happen when there are no other subscribers.
@@ -93,7 +92,7 @@ func (is *IndexerService) OnStart() error {
}
// OnStop implements service.Service by unsubscribing from all transactions.
func (is *IndexerService) OnStop() {
func (is *Service) OnStop() {
if is.eventBus.IsRunning() {
_ = is.eventBus.UnsubscribeAll(context.Background(), subscriber)
}
@@ -1,4 +1,4 @@
package txindex_test
package indexer_test
import (
"testing"
@@ -9,9 +9,9 @@ import (
abci "github.com/tendermint/tendermint/abci/types"
"github.com/tendermint/tendermint/libs/log"
indexer "github.com/tendermint/tendermint/state/indexer"
blockidxkv "github.com/tendermint/tendermint/state/indexer/block/kv"
"github.com/tendermint/tendermint/state/txindex"
"github.com/tendermint/tendermint/state/txindex/kv"
"github.com/tendermint/tendermint/state/indexer/tx/kv"
"github.com/tendermint/tendermint/types"
)
@@ -32,7 +32,7 @@ func TestIndexerServiceIndexesBlocks(t *testing.T) {
txIndexer := kv.NewTxIndex(store)
blockIndexer := blockidxkv.New(db.NewPrefixDB(store, []byte("block_events")))
service := txindex.NewIndexerService(txIndexer, blockIndexer, eventBus)
service := indexer.NewIndexerService(txIndexer, blockIndexer, eventBus)
service.SetLogger(log.TestingLogger())
err = service.Start()
require.NoError(t, err)
@@ -13,12 +13,11 @@ import (
abci "github.com/tendermint/tendermint/abci/types"
"github.com/tendermint/tendermint/libs/pubsub/query"
"github.com/tendermint/tendermint/state/indexer"
"github.com/tendermint/tendermint/state/txindex"
indexer "github.com/tendermint/tendermint/state/indexer"
"github.com/tendermint/tendermint/types"
)
var _ txindex.TxIndexer = (*TxIndex)(nil)
var _ indexer.TxIndexer = (*TxIndex)(nil)
// TxIndex is the simplest possible indexer
// It is backed by two kv stores:
@@ -39,7 +38,7 @@ func NewTxIndex(store dbm.DB) *TxIndex {
// transaction is not found.
func (txi *TxIndex) Get(hash []byte) (*abci.TxResult, error) {
if len(hash) == 0 {
return nil, txindex.ErrorEmptyHash
return nil, indexer.ErrorEmptyHash
}
rawBytes, err := txi.store.Get(primaryKey(hash))
@@ -63,7 +62,7 @@ func (txi *TxIndex) Get(hash []byte) (*abci.TxResult, error) {
// key that indexed from the tx's events is a composite of the event type and
// the respective attribute's key delimited by a "." (eg. "account.number").
// Any event with an empty type is not indexed.
func (txi *TxIndex) AddBatch(b *txindex.Batch) error {
func (txi *TxIndex) AddBatch(b *indexer.Batch) error {
storeBatch := txi.store.NewBatch()
defer storeBatch.Close()
@@ -16,12 +16,12 @@ import (
abci "github.com/tendermint/tendermint/abci/types"
"github.com/tendermint/tendermint/libs/pubsub/query"
tmrand "github.com/tendermint/tendermint/libs/rand"
"github.com/tendermint/tendermint/state/txindex"
indexer "github.com/tendermint/tendermint/state/indexer"
"github.com/tendermint/tendermint/types"
)
func TestTxIndex(t *testing.T) {
indexer := NewTxIndex(db.NewMemDB())
txIndexer := NewTxIndex(db.NewMemDB())
tx := types.Tx("HELLO WORLD")
txResult := &abci.TxResult{
@@ -35,14 +35,14 @@ func TestTxIndex(t *testing.T) {
}
hash := tx.Hash()
batch := txindex.NewBatch(1)
batch := indexer.NewBatch(1)
if err := batch.Add(txResult); err != nil {
t.Error(err)
}
err := indexer.AddBatch(batch)
err := txIndexer.AddBatch(batch)
require.NoError(t, err)
loadedTxResult, err := indexer.Get(hash)
loadedTxResult, err := txIndexer.Get(hash)
require.NoError(t, err)
assert.True(t, proto.Equal(txResult, loadedTxResult))
@@ -58,10 +58,10 @@ func TestTxIndex(t *testing.T) {
}
hash2 := tx2.Hash()
err = indexer.Index(txResult2)
err = txIndexer.Index(txResult2)
require.NoError(t, err)
loadedTxResult2, err := indexer.Get(hash2)
loadedTxResult2, err := txIndexer.Get(hash2)
require.NoError(t, err)
assert.True(t, proto.Equal(txResult2, loadedTxResult2))
}
@@ -341,9 +341,9 @@ func benchmarkTxIndex(txsCount int64, b *testing.B) {
store, err := db.NewDB("tx_index", "goleveldb", dir)
require.NoError(b, err)
indexer := NewTxIndex(store)
txIndexer := NewTxIndex(store)
batch := txindex.NewBatch(txsCount)
batch := indexer.NewBatch(txsCount)
txIndex := uint32(0)
for i := int64(0); i < txsCount; i++ {
tx := tmrand.Bytes(250)
@@ -367,7 +367,7 @@ func benchmarkTxIndex(txsCount int64, b *testing.B) {
b.ResetTimer()
for n := 0; n < b.N; n++ {
err = indexer.AddBatch(batch)
err = txIndexer.AddBatch(batch)
}
if err != nil {
b.Fatal(err)
@@ -6,10 +6,10 @@ import (
abci "github.com/tendermint/tendermint/abci/types"
"github.com/tendermint/tendermint/libs/pubsub/query"
"github.com/tendermint/tendermint/state/txindex"
"github.com/tendermint/tendermint/state/indexer"
)
var _ txindex.TxIndexer = (*TxIndex)(nil)
var _ indexer.TxIndexer = (*TxIndex)(nil)
// TxIndex acts as a /dev/null.
type TxIndex struct{}
@@ -20,7 +20,7 @@ func (txi *TxIndex) Get(hash []byte) (*abci.TxResult, error) {
}
// AddBatch is a noop and always returns nil.
func (txi *TxIndex) AddBatch(batch *txindex.Batch) error {
func (txi *TxIndex) AddBatch(batch *indexer.Batch) error {
return nil
}