Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .mockery.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -222,6 +222,7 @@ packages:
interfaces:
HeadReporter:
PrometheusBackend:
HeadMetrics:
github.com/smartcontractkit/libocr/commontypes:
config:
dir: "common/types/mocks"
Expand Down
6 changes: 6 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,5 +1,11 @@
# Changelog Chainlink Core

## 2.58.1

### Minor Changes

- [#23250](https://github.com/smartcontractkit/chainlink/pull/23250) [`1543766`](https://github.com/smartcontractkit/chainlink/commit/15437660cdb66db2e9706f7e4b5cd179b319fb18) - #added Added head-reporter ability to expose heads in beholder metrics along with prom metrics and sending telemetry over OTI

## 2.58.0

### Minor Changes
Expand Down
11 changes: 10 additions & 1 deletion core/services/chainlink/application.go
Original file line number Diff line number Diff line change
Expand Up @@ -215,7 +215,7 @@
// the logger at the same directory and returns the Application to
// be used by the node.
// TODO: Inject more dependencies here to save booting up useless stuff in tests
func NewApplication(ctx context.Context, opts ApplicationOpts) (Application, error) {

Check warning on line 218 in core/services/chainlink/application.go

View check run for this annotation

CL-sonarqube-production / SonarQube Code Analysis

Refactor this method to reduce its Cognitive Complexity from 95 to the 30 allowed.

[S3776] Cognitive Complexity of functions should not be too high See more on https://sonarqube.main.prod.cldev.sh/project/issues?id=smartcontractkit_chainlink&pullRequest=23557&issues=6810e5c3-acc3-46ff-84a4-e94ee72ae3ca&open=6810e5c3-acc3-46ff-84a4-e94ee72ae3ca
var srvcs []services.ServiceCtx

heartbeat := NewHeartbeat(NewHeartbeatConfig(opts))
Expand Down Expand Up @@ -600,7 +600,16 @@

legacyEVMTelemReporter := headreporter.NewLegacyEVMTelemetryReporter(telemetryManager, globalLogger, evmChainIDs...)
loopTelemReporter := headreporter.NewTelemetryReporter(telemetryManager, globalLogger, relayChainInterops.GetIDToRelayerMap())
headReporter := headreporter.NewHeadReporterService(opts.DS, globalLogger, promReporter, legacyEVMTelemReporter, loopTelemReporter)
headReporters := []headreporter.HeadReporter{promReporter, legacyEVMTelemReporter, loopTelemReporter}
if headMetrics, metricsErr := headreporter.NewBeholderHeadMetrics(); metricsErr != nil {
globalLogger.Errorw("Failed to initialize head reporter Beholder metrics; skipping head metrics reporters", "err", metricsErr)
} else {
headReporters = append(headReporters, headreporter.NewEVMMetricsReporter(headMetrics, globalLogger, evmChainIDs...))
if relayerMetricsReporter := headreporter.NewRelayerMetricsReporter(headMetrics, globalLogger, relayChainInterops.GetIDToRelayerMap()); relayerMetricsReporter != nil {
headReporters = append(headReporters, relayerMetricsReporter)
}
}
headReporter := headreporter.NewHeadReporterService(opts.DS, globalLogger, headReporters...)
srvcs = append(srvcs, headReporter)
for _, chain := range legacyEVMChains.Slice() {
legacyChain, ok := chain.(legacyevm.Chain)
Expand Down
113 changes: 113 additions & 0 deletions core/services/headreporter/beholder_metrics.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,113 @@
package headreporter

import (
"context"
"fmt"
"strconv"

"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/metric"

"github.com/smartcontractkit/chainlink-common/pkg/beholder"
)

const (
metricLatestBlockNumber = "head_reporter_latest_block_number"
metricLatestBlockTimestamp = "head_reporter_latest_block_timestamp_seconds"
metricFinalizedBlockNumber = "head_reporter_finalized_block_number"
metricFinalizedBlockTimestamp = "head_reporter_finalized_block_timestamp_seconds"
metricFinalityDepth = "head_reporter_finality_depth_blocks"

attrChainID = "chain_id"
attrNetwork = "network"
attrChainSelector = "chain_selector"
)

type (
// finalizedBlock is the finalized counterpart of a headReport's latest block.
finalizedBlock struct {
number int64
ts uint64
}

// headReport is a family-neutral snapshot of chain head data, produced by both the
// in-process EVM reporter and the generic relayer reporter.
headReport struct {
chainID string
network string // chain family, e.g. "evm", "solana"
chainSelector uint64
hasSelector bool
latestNumber int64
latestTs uint64
finalized *finalizedBlock // nil if not available
}

// HeadMetrics records head-report data as Beholder (OTel) metrics.
HeadMetrics interface {
RecordHeadReport(ctx context.Context, r headReport)
}

beholderHeadMetrics struct {
latestNumber metric.Int64Gauge
latestTs metric.Int64Gauge
finalizedNumber metric.Int64Gauge
finalizedTs metric.Int64Gauge
finalityDepth metric.Int64Gauge
}
)

// NewBeholderHeadMetrics registers the head-reporter gauge instruments once via
// beholder.GetMeter(). It is safe to call unconditionally: GetMeter() returns a no-op meter
// when Beholder telemetry is disabled.
func NewBeholderHeadMetrics() (HeadMetrics, error) {
m := beholder.GetMeter()

latestNumber, err := m.Int64Gauge(metricLatestBlockNumber)
if err != nil {
return nil, fmt.Errorf("failed to register %s gauge: %w", metricLatestBlockNumber, err)

Check warning on line 67 in core/services/headreporter/beholder_metrics.go

View check run for this annotation

CL-sonarqube-production / SonarQube Code Analysis

Define a constant instead of duplicating this literal "failed to register %s gauge: %w" 5 times.

[S1192] String literals should not be duplicated See more on https://sonarqube.main.prod.cldev.sh/project/issues?id=smartcontractkit_chainlink&pullRequest=23557&issues=c95a2bce-8c9a-48ee-8ed3-eff5023c72cd&open=c95a2bce-8c9a-48ee-8ed3-eff5023c72cd
}
latestTs, err := m.Int64Gauge(metricLatestBlockTimestamp)
if err != nil {
return nil, fmt.Errorf("failed to register %s gauge: %w", metricLatestBlockTimestamp, err)
}
finalizedNumber, err := m.Int64Gauge(metricFinalizedBlockNumber)
if err != nil {
return nil, fmt.Errorf("failed to register %s gauge: %w", metricFinalizedBlockNumber, err)
}
finalizedTs, err := m.Int64Gauge(metricFinalizedBlockTimestamp)
if err != nil {
return nil, fmt.Errorf("failed to register %s gauge: %w", metricFinalizedBlockTimestamp, err)
}
finalityDepth, err := m.Int64Gauge(metricFinalityDepth)
if err != nil {
return nil, fmt.Errorf("failed to register %s gauge: %w", metricFinalityDepth, err)
}

return &beholderHeadMetrics{
latestNumber: latestNumber,
latestTs: latestTs,
finalizedNumber: finalizedNumber,
finalizedTs: finalizedTs,
finalityDepth: finalityDepth,
}, nil
}

func (b *beholderHeadMetrics) RecordHeadReport(ctx context.Context, r headReport) {
kvs := []attribute.KeyValue{
attribute.String(attrChainID, r.chainID),
attribute.String(attrNetwork, r.network),
}
if r.hasSelector {
kvs = append(kvs, attribute.String(attrChainSelector, strconv.FormatUint(r.chainSelector, 10)))
}
attrs := metric.WithAttributes(kvs...)

b.latestNumber.Record(ctx, r.latestNumber, attrs)
b.latestTs.Record(ctx, int64(r.latestTs), attrs) //nolint:gosec // unix timestamps fit int64 until year 292277026596

if r.finalized != nil {
b.finalizedNumber.Record(ctx, r.finalized.number, attrs)
b.finalizedTs.Record(ctx, int64(r.finalized.ts), attrs) //nolint:gosec // unix timestamps fit int64 until year 292277026596
b.finalityDepth.Record(ctx, r.latestNumber-r.finalized.number, attrs)
}
}
114 changes: 114 additions & 0 deletions core/services/headreporter/beholder_metrics_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,114 @@
package headreporter

import (
"context"
"testing"

"github.com/stretchr/testify/require"
sdkmetric "go.opentelemetry.io/otel/sdk/metric"
"go.opentelemetry.io/otel/sdk/metric/metricdata"
)

func newTestBeholderHeadMetrics(t *testing.T) (*beholderHeadMetrics, *sdkmetric.ManualReader) {
t.Helper()
reader := sdkmetric.NewManualReader()
mp := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader))
t.Cleanup(func() { _ = mp.Shutdown(context.Background()) })

m := mp.Meter("test")
latestNumber, err := m.Int64Gauge("head_reporter_latest_block_number")
require.NoError(t, err)
latestTs, err := m.Int64Gauge("head_reporter_latest_block_timestamp_seconds")
require.NoError(t, err)
finalizedNumber, err := m.Int64Gauge("head_reporter_finalized_block_number")
require.NoError(t, err)
finalizedTs, err := m.Int64Gauge("head_reporter_finalized_block_timestamp_seconds")
require.NoError(t, err)
finalityDepth, err := m.Int64Gauge("head_reporter_finality_depth_blocks")
require.NoError(t, err)

return &beholderHeadMetrics{
latestNumber: latestNumber,
latestTs: latestTs,
finalizedNumber: finalizedNumber,
finalizedTs: finalizedTs,
finalityDepth: finalityDepth,
}, reader
}

func collectGauge(t *testing.T, rm metricdata.ResourceMetrics, name string) metricdata.Gauge[int64] {
t.Helper()
for _, sm := range rm.ScopeMetrics {
for _, m := range sm.Metrics {
if m.Name == name {
data, ok := m.Data.(metricdata.Gauge[int64])
require.True(t, ok, "metric %s is not an int64 gauge", name)
return data
}
}
}
t.Fatalf("metric %s not found", name)
return metricdata.Gauge[int64]{}
}

func Test_BeholderHeadMetrics_RecordHeadReport_WithFinalized(t *testing.T) {
t.Parallel()
metrics, reader := newTestBeholderHeadMetrics(t)

metrics.RecordHeadReport(t.Context(), headReport{
chainID: "100",
network: "evm",
chainSelector: 465200170687744372,
hasSelector: true,
latestNumber: 42,
latestTs: 1000,
finalized: &finalizedBlock{
number: 40,
ts: 900,
},
})

var rm metricdata.ResourceMetrics
require.NoError(t, reader.Collect(t.Context(), &rm))

latest := collectGauge(t, rm, "head_reporter_latest_block_number")
require.Len(t, latest.DataPoints, 1)
require.Equal(t, int64(42), latest.DataPoints[0].Value)

depth := collectGauge(t, rm, "head_reporter_finality_depth_blocks")
require.Len(t, depth.DataPoints, 1)
require.Equal(t, int64(2), depth.DataPoints[0].Value)

finalizedNumber := collectGauge(t, rm, "head_reporter_finalized_block_number")
require.Len(t, finalizedNumber.DataPoints, 1)
require.Equal(t, int64(40), finalizedNumber.DataPoints[0].Value)
}

func Test_BeholderHeadMetrics_RecordHeadReport_NoFinalized(t *testing.T) {
t.Parallel()
metrics, reader := newTestBeholderHeadMetrics(t)

metrics.RecordHeadReport(t.Context(), headReport{
chainID: "testchain",
network: "solana",
latestNumber: 42,
latestTs: 1000,
})

var rm metricdata.ResourceMetrics
require.NoError(t, reader.Collect(t.Context(), &rm))

latest := collectGauge(t, rm, "head_reporter_latest_block_number")
require.Len(t, latest.DataPoints, 1)
require.Equal(t, int64(42), latest.DataPoints[0].Value)

for _, sm := range rm.ScopeMetrics {
for _, m := range sm.Metrics {
if m.Name == "head_reporter_finalized_block_number" || m.Name == "head_reporter_finality_depth_blocks" {
data, ok := m.Data.(metricdata.Gauge[int64])
require.True(t, ok)
require.Empty(t, data.DataPoints)
}
}
}
}
70 changes: 70 additions & 0 deletions core/services/headreporter/head_metrics_mock.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading
Loading