mirror of
https://github.com/tendermint/tendermint.git
synced 2026-09-19 22:44:24 +00:00
@@ -46,8 +46,11 @@ func New(cfg *config.RPCConfig, bs state.BlockStore, ss state.Store, es []indexe
|
||||
routes := rpc.Routes(*cfg, ss, bs, es, logger)
|
||||
eb := types.NewEventBus()
|
||||
eb.SetLogger(logger.With("module", "events"))
|
||||
is := indexer.NewIndexerService(es, eb)
|
||||
is.SetLogger(logger.With("module", "txindex"))
|
||||
is := indexer.NewService(indexer.ServiceArgs{
|
||||
Sinks: es,
|
||||
EventBus: eb,
|
||||
Logger: logger.With("module", "txindex"),
|
||||
})
|
||||
return &Inspector{
|
||||
routes: routes,
|
||||
config: cfg,
|
||||
|
||||
@@ -2,7 +2,9 @@ package indexer
|
||||
|
||||
import (
|
||||
"context"
|
||||
"time"
|
||||
|
||||
"github.com/tendermint/tendermint/libs/log"
|
||||
"github.com/tendermint/tendermint/libs/service"
|
||||
"github.com/tendermint/tendermint/types"
|
||||
)
|
||||
@@ -20,14 +22,30 @@ type Service struct {
|
||||
|
||||
eventSinks []EventSink
|
||||
eventBus *types.EventBus
|
||||
metrics *Metrics
|
||||
}
|
||||
|
||||
// NewService constructs a new indexer service from the given arguments.
|
||||
func NewService(args ServiceArgs) *Service {
|
||||
is := &Service{
|
||||
eventSinks: args.Sinks,
|
||||
eventBus: args.EventBus,
|
||||
metrics: args.Metrics,
|
||||
}
|
||||
if is.metrics == nil {
|
||||
is.metrics = NopMetrics()
|
||||
}
|
||||
is.BaseService = *service.NewBaseService(args.Logger, "IndexerService", is)
|
||||
return is
|
||||
}
|
||||
|
||||
// NewIndexerService returns a new service instance.
|
||||
// Deprecated: Use NewService instead.
|
||||
func NewIndexerService(es []EventSink, eventBus *types.EventBus) *Service {
|
||||
|
||||
is := &Service{eventSinks: es, eventBus: eventBus}
|
||||
is.BaseService = *service.NewBaseService(nil, "IndexerService", is)
|
||||
return is
|
||||
return NewService(ServiceArgs{
|
||||
Sinks: es,
|
||||
EventBus: eventBus,
|
||||
})
|
||||
}
|
||||
|
||||
// OnStart implements service.Service by subscribing for all transactions
|
||||
@@ -79,17 +97,23 @@ func (is *Service) OnStart() error {
|
||||
}
|
||||
|
||||
for _, sink := range is.eventSinks {
|
||||
start := time.Now()
|
||||
if err := sink.IndexBlockEvents(eventDataHeader); err != nil {
|
||||
is.Logger.Error("failed to index block", "height", height, "err", err)
|
||||
} else {
|
||||
is.metrics.BlockEventsSeconds.Observe(time.Since(start).Seconds())
|
||||
is.metrics.BlocksIndexed.Add(1)
|
||||
is.Logger.Debug("indexed block", "height", height, "sink", sink.Type())
|
||||
}
|
||||
|
||||
if len(batch.Ops) > 0 {
|
||||
start := time.Now()
|
||||
err := sink.IndexTxEvents(batch.Ops)
|
||||
if err != nil {
|
||||
is.Logger.Error("failed to index block txs", "height", height, "err", err)
|
||||
} else {
|
||||
is.metrics.TxEventsSeconds.Observe(time.Since(start).Seconds())
|
||||
is.metrics.TransactionsIndexed.Add(float64(len(batch.Ops)))
|
||||
is.Logger.Debug("indexed txs", "height", height, "sink", sink.Type())
|
||||
}
|
||||
}
|
||||
@@ -114,6 +138,14 @@ func (is *Service) OnStop() {
|
||||
}
|
||||
}
|
||||
|
||||
// ServiceArgs are arguments for constructing a new indexer service.
|
||||
type ServiceArgs struct {
|
||||
Sinks []EventSink
|
||||
EventBus *types.EventBus
|
||||
Metrics *Metrics
|
||||
Logger log.Logger
|
||||
}
|
||||
|
||||
// KVSinkEnabled returns the given eventSinks is containing KVEventSink.
|
||||
func KVSinkEnabled(sinks []EventSink) bool {
|
||||
for _, sink := range sinks {
|
||||
|
||||
@@ -0,0 +1,73 @@
|
||||
package indexer
|
||||
|
||||
import (
|
||||
"github.com/go-kit/kit/metrics"
|
||||
"github.com/go-kit/kit/metrics/discard"
|
||||
|
||||
prometheus "github.com/go-kit/kit/metrics/prometheus"
|
||||
stdprometheus "github.com/prometheus/client_golang/prometheus"
|
||||
)
|
||||
|
||||
// MetricsSubsystem is a the subsystem label for the indexer package.
|
||||
const MetricsSubsystem = "indexer"
|
||||
|
||||
// Metrics contains metrics exposed by this package.
|
||||
type Metrics struct {
|
||||
// Latency for indexing block events.
|
||||
BlockEventsSeconds metrics.Histogram
|
||||
|
||||
// Latency for indexing transaction events.
|
||||
TxEventsSeconds metrics.Histogram
|
||||
|
||||
// Number of complete blocks indexed.
|
||||
BlocksIndexed metrics.Counter
|
||||
|
||||
// Number of transactions indexed.
|
||||
TransactionsIndexed metrics.Counter
|
||||
}
|
||||
|
||||
// PrometheusMetrics returns Metrics build using Prometheus client library.
|
||||
// Optionally, labels can be provided along with their values ("foo",
|
||||
// "fooValue").
|
||||
func PrometheusMetrics(namespace string, labelsAndValues ...string) *Metrics {
|
||||
labels := []string{}
|
||||
for i := 0; i < len(labelsAndValues); i += 2 {
|
||||
labels = append(labels, labelsAndValues[i])
|
||||
}
|
||||
return &Metrics{
|
||||
BlockEventsSeconds: prometheus.NewHistogramFrom(stdprometheus.HistogramOpts{
|
||||
Namespace: namespace,
|
||||
Subsystem: MetricsSubsystem,
|
||||
Name: "block_events_seconds",
|
||||
Help: "Latency for indexing block events.",
|
||||
}, labels).With(labelsAndValues...),
|
||||
TxEventsSeconds: prometheus.NewHistogramFrom(stdprometheus.HistogramOpts{
|
||||
Namespace: namespace,
|
||||
Subsystem: MetricsSubsystem,
|
||||
Name: "tx_events_seconds",
|
||||
Help: "Latency for indexing transaction events.",
|
||||
}, labels).With(labelsAndValues...),
|
||||
BlocksIndexed: prometheus.NewCounterFrom(stdprometheus.CounterOpts{
|
||||
Namespace: namespace,
|
||||
Subsystem: MetricsSubsystem,
|
||||
Name: "blocks_indexed",
|
||||
Help: "Number of complete blocks indexed.",
|
||||
}, labels).With(labelsAndValues...),
|
||||
TransactionsIndexed: prometheus.NewCounterFrom(stdprometheus.CounterOpts{
|
||||
Namespace: namespace,
|
||||
Subsystem: MetricsSubsystem,
|
||||
Name: "transactions_indexed",
|
||||
Help: "Number of transactions indexed.",
|
||||
}, labels).With(labelsAndValues...),
|
||||
}
|
||||
}
|
||||
|
||||
// NopMetrics returns an indexer metrics stub that discards all samples.
|
||||
func NopMetrics() *Metrics {
|
||||
return &Metrics{
|
||||
BlockEventsSeconds: discard.NewHistogram(),
|
||||
TxEventsSeconds: discard.NewHistogram(),
|
||||
BlocksIndexed: discard.NewCounter(),
|
||||
TransactionsIndexed: discard.NewCounter(),
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user