From 088a57ad5fb114c8bae4ad633a215ffe3a04ab7e Mon Sep 17 00:00:00 2001 From: Aleksandr Bezobchuk Date: Thu, 16 Jun 2022 15:15:55 -0400 Subject: [PATCH] updates --- node/node.go | 16 ++++++++++++++++ state/indexer/sink/pubsub/pubsub.go | 4 ++-- 2 files changed, 18 insertions(+), 2 deletions(-) diff --git a/node/node.go b/node/node.go index 18fde0062..086a7a133 100644 --- a/node/node.go +++ b/node/node.go @@ -41,6 +41,7 @@ import ( blockidxkv "github.com/tendermint/tendermint/state/indexer/block/kv" blockidxnull "github.com/tendermint/tendermint/state/indexer/block/null" "github.com/tendermint/tendermint/state/indexer/sink/psql" + "github.com/tendermint/tendermint/state/indexer/sink/pubsub" "github.com/tendermint/tendermint/state/txindex" "github.com/tendermint/tendermint/state/txindex/kv" "github.com/tendermint/tendermint/state/txindex/null" @@ -288,13 +289,28 @@ func createAndStartIndexerService( if config.TxIndex.PsqlConn == "" { return nil, nil, nil, errors.New(`no psql-conn is set for the "psql" indexer`) } + es, err := psql.NewEventSink(config.TxIndex.PsqlConn, chainID) if err != nil { return nil, nil, nil, fmt.Errorf("creating psql indexer: %w", err) } + txIndexer = es.TxIndexer() blockIndexer = es.BlockIndexer() + case "pubsub": + if config.TxIndex.PubsubProjectID == "" { + return nil, nil, nil, errors.New("no 'pubsub-project-id' is set for the 'pubsub' indexer") + } + + sink, err := pubsub.NewEventSink(config.TxIndex.PubsubProjectID, chainID) + if err != nil { + return nil, nil, nil, fmt.Errorf("creating pubsub indexer: %w", err) + } + + txIndexer = pubsub.NewTxIndexer(sink) + blockIndexer = pubsub.NewBlockIndexer(sink) + default: txIndexer = &null.TxIndex{} blockIndexer = &blockidxnull.BlockerIndexer{} diff --git a/state/indexer/sink/pubsub/pubsub.go b/state/indexer/sink/pubsub/pubsub.go index 94446a2bd..a7fc4c194 100644 --- a/state/indexer/sink/pubsub/pubsub.go +++ b/state/indexer/sink/pubsub/pubsub.go @@ -18,7 +18,7 @@ type EventSink struct { chainID string } -func NewEventSink(connStr, chainID string) (*EventSink, error) { +func NewEventSink(projectID, chainID string) (*EventSink, error) { if s := os.Getenv(credsEnvVar); len(s) == 0 { return nil, fmt.Errorf("missing '%s' environment variable", credsEnvVar) } @@ -26,7 +26,7 @@ func NewEventSink(connStr, chainID string) (*EventSink, error) { ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) defer cancel() - c, err := pubsub.NewClient(ctx, "project-id") + c, err := pubsub.NewClient(ctx, projectID) if err != nil { return nil, fmt.Errorf("failed to create a Google Cloud Pubsub client: %w", err) }