mirror of
https://github.com/tendermint/tendermint.git
synced 2026-08-29 04:07:07 +00:00
* node: hook up eventlog and eventlog metrics (#7981) * align no metrics for event log * make use of metricsgen Co-authored-by: M. J. Fromberger <fromberger@interchain.io>
This commit is contained in:
co-authored by
M. J. Fromberger
parent
d47b675cda
commit
c3cc94a0e0
@@ -24,9 +24,9 @@ import (
|
||||
// any number of readers.
|
||||
type Log struct {
|
||||
// These values do not change after construction.
|
||||
windowSize time.Duration
|
||||
maxItems int
|
||||
numItemsGauge gauge
|
||||
windowSize time.Duration
|
||||
maxItems int
|
||||
metrics *Metrics
|
||||
|
||||
// Protects access to the fields below. Lock to modify the values of these
|
||||
// fields, or to read or snapshot the values.
|
||||
@@ -45,14 +45,14 @@ func New(opts LogSettings) (*Log, error) {
|
||||
return nil, errors.New("window size must be positive")
|
||||
}
|
||||
lg := &Log{
|
||||
windowSize: opts.WindowSize,
|
||||
maxItems: opts.MaxItems,
|
||||
numItemsGauge: discard{},
|
||||
ready: make(chan struct{}),
|
||||
source: opts.Source,
|
||||
windowSize: opts.WindowSize,
|
||||
maxItems: opts.MaxItems,
|
||||
metrics: NopMetrics(),
|
||||
ready: make(chan struct{}),
|
||||
source: opts.Source,
|
||||
}
|
||||
if opts.Metrics != nil {
|
||||
lg.numItemsGauge = opts.Metrics.numItemsGauge
|
||||
lg.metrics = opts.Metrics
|
||||
}
|
||||
return lg, nil
|
||||
}
|
||||
|
||||
@@ -0,0 +1,30 @@
|
||||
// Code generated by metricsgen. DO NOT EDIT.
|
||||
|
||||
package eventlog
|
||||
|
||||
import (
|
||||
"github.com/go-kit/kit/metrics/discard"
|
||||
prometheus "github.com/go-kit/kit/metrics/prometheus"
|
||||
stdprometheus "github.com/prometheus/client_golang/prometheus"
|
||||
)
|
||||
|
||||
func PrometheusMetrics(namespace string, labelsAndValues ...string) *Metrics {
|
||||
labels := []string{}
|
||||
for i := 0; i < len(labelsAndValues); i += 2 {
|
||||
labels = append(labels, labelsAndValues[i])
|
||||
}
|
||||
return &Metrics{
|
||||
numItems: prometheus.NewGaugeFrom(stdprometheus.GaugeOpts{
|
||||
Namespace: namespace,
|
||||
Subsystem: MetricsSubsystem,
|
||||
Name: "num_items",
|
||||
Help: "Number of items currently resident in the event log.",
|
||||
}, labels).With(labelsAndValues...),
|
||||
}
|
||||
}
|
||||
|
||||
func NopMetrics() *Metrics {
|
||||
return &Metrics{
|
||||
numItems: discard.NewGauge(),
|
||||
}
|
||||
}
|
||||
@@ -1,39 +1,14 @@
|
||||
package eventlog
|
||||
|
||||
import (
|
||||
"github.com/go-kit/kit/metrics/prometheus"
|
||||
stdprometheus "github.com/prometheus/client_golang/prometheus"
|
||||
)
|
||||
import "github.com/go-kit/kit/metrics"
|
||||
|
||||
// gauge is the subset of the Prometheus gauge interface used here.
|
||||
type gauge interface {
|
||||
Set(float64)
|
||||
}
|
||||
const MetricsSubsystem = "eventlog"
|
||||
|
||||
//go:generate go run ../../scripts/metricsgen -struct=Metrics
|
||||
|
||||
// Metrics define the metrics exported by the eventlog package.
|
||||
type Metrics struct {
|
||||
numItemsGauge gauge
|
||||
}
|
||||
|
||||
// discard is a no-op implementation of the gauge interface.
|
||||
type discard struct{}
|
||||
|
||||
func (discard) Set(float64) {}
|
||||
|
||||
const eventlogSubsystem = "eventlog"
|
||||
|
||||
// PrometheusMetrics returns a collection of eventlog metrics for Prometheus.
|
||||
func PrometheusMetrics(ns string, fields ...string) *Metrics {
|
||||
var labels []string
|
||||
for i := 0; i < len(fields); i += 2 {
|
||||
labels = append(labels, fields[i])
|
||||
}
|
||||
return &Metrics{
|
||||
numItemsGauge: prometheus.NewGaugeFrom(stdprometheus.GaugeOpts{
|
||||
Namespace: ns,
|
||||
Subsystem: eventlogSubsystem,
|
||||
Name: "num_items",
|
||||
Help: "Number of items currently resident in the event log.",
|
||||
}, labels).With(fields...),
|
||||
}
|
||||
// Number of items currently resident in the event log.
|
||||
numItems metrics.Gauge
|
||||
}
|
||||
|
||||
@@ -12,7 +12,7 @@ func (lg *Log) checkPrune(head *logEntry, size int, age time.Duration) error {
|
||||
const windowSlop = 30 * time.Second
|
||||
|
||||
if age < (lg.windowSize+windowSlop) && (lg.maxItems <= 0 || size <= lg.maxItems) {
|
||||
lg.numItemsGauge.Set(float64(lg.numItems))
|
||||
lg.metrics.numItems.Set(float64(lg.numItems))
|
||||
return nil // no pruning is needed
|
||||
}
|
||||
|
||||
@@ -46,7 +46,7 @@ func (lg *Log) checkPrune(head *logEntry, size int, age time.Duration) error {
|
||||
lg.mu.Lock()
|
||||
defer lg.mu.Unlock()
|
||||
lg.numItems = newState.size
|
||||
lg.numItemsGauge.Set(float64(newState.size))
|
||||
lg.metrics.numItems.Set(float64(newState.size))
|
||||
lg.oldestCursor = newState.oldest
|
||||
lg.head = newState.head
|
||||
return err
|
||||
|
||||
+19
-4
@@ -115,20 +115,21 @@ func DefaultNewNode(config *cfg.Config, logger log.Logger) (*Node, error) {
|
||||
}
|
||||
|
||||
// MetricsProvider returns a consensus, p2p and mempool Metrics.
|
||||
type MetricsProvider func(chainID string) (*cs.Metrics, *p2p.Metrics, *mempl.Metrics, *sm.Metrics, *proxy.Metrics)
|
||||
type MetricsProvider func(chainID string) (*cs.Metrics, *eventlog.Metrics, *p2p.Metrics, *mempl.Metrics, *sm.Metrics, *proxy.Metrics)
|
||||
|
||||
// DefaultMetricsProvider returns Metrics build using Prometheus client library
|
||||
// if Prometheus is enabled. Otherwise, it returns no-op Metrics.
|
||||
func DefaultMetricsProvider(config *cfg.InstrumentationConfig) MetricsProvider {
|
||||
return func(chainID string) (*cs.Metrics, *p2p.Metrics, *mempl.Metrics, *sm.Metrics, *proxy.Metrics) {
|
||||
return func(chainID string) (*cs.Metrics, *eventlog.Metrics, *p2p.Metrics, *mempl.Metrics, *sm.Metrics, *proxy.Metrics) {
|
||||
if config.Prometheus {
|
||||
return cs.PrometheusMetrics(config.Namespace, "chain_id", chainID),
|
||||
eventlog.PrometheusMetrics(config.Namespace, "chain_id", chainID),
|
||||
p2p.PrometheusMetrics(config.Namespace, "chain_id", chainID),
|
||||
mempl.PrometheusMetrics(config.Namespace, "chain_id", chainID),
|
||||
sm.PrometheusMetrics(config.Namespace, "chain_id", chainID),
|
||||
proxy.PrometheusMetrics(config.Namespace, "chain_id", chainID)
|
||||
}
|
||||
return cs.NopMetrics(), p2p.NopMetrics(), mempl.NopMetrics(), sm.NopMetrics(), proxy.NopMetrics()
|
||||
return cs.NopMetrics(), eventlog.NopMetrics(), p2p.NopMetrics(), mempl.NopMetrics(), sm.NopMetrics(), proxy.NopMetrics()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -727,7 +728,7 @@ func NewNode(config *cfg.Config,
|
||||
return nil, err
|
||||
}
|
||||
|
||||
csMetrics, p2pMetrics, memplMetrics, smMetrics, abciMetrics := metricsProvider(genDoc.ChainID)
|
||||
csMetrics, eventLogMetrics, p2pMetrics, memplMetrics, smMetrics, abciMetrics := metricsProvider(genDoc.ChainID)
|
||||
|
||||
// Create the proxyApp and establish connections to the ABCI app (consensus, mempool, query).
|
||||
proxyApp, err := createAndStartProxyAppConns(clientCreator, logger, abciMetrics)
|
||||
@@ -744,6 +745,19 @@ func NewNode(config *cfg.Config,
|
||||
return nil, err
|
||||
}
|
||||
|
||||
var eventLog *eventlog.Log
|
||||
if w := config.RPC.EventLogWindowSize; w > 0 {
|
||||
var err error
|
||||
eventLog, err = eventlog.New(eventlog.LogSettings{
|
||||
WindowSize: w,
|
||||
MaxItems: config.RPC.EventLogMaxItems,
|
||||
Metrics: eventLogMetrics,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("initializing event log: %w", err)
|
||||
}
|
||||
}
|
||||
|
||||
indexerService, txIndexer, blockIndexer, err := createAndStartIndexerService(config,
|
||||
genDoc.ChainID, dbProvider, eventBus, logger)
|
||||
if err != nil {
|
||||
@@ -927,6 +941,7 @@ func NewNode(config *cfg.Config,
|
||||
indexerService: indexerService,
|
||||
blockIndexer: blockIndexer,
|
||||
eventBus: eventBus,
|
||||
eventLog: eventLog,
|
||||
}
|
||||
node.BaseService = *service.NewBaseService(logger, "Node", node)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user