Skip to content
Open
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
47 changes: 31 additions & 16 deletions pkg/beholder/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,6 @@ import (
"fmt"
"io"

"go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc"
"go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploghttp"
Expand Down Expand Up @@ -140,21 +139,32 @@ func NewGRPCClient(cfg Config, otlploggrpcNew otlploggrpcFactory) (*Client, erro
tracer := tracerProvider.Tracer(defaultPackageName)

// Meter
meterProvider, err := newMeterProvider(cfg, baseResource, auth, creds)
meterProvider, meteredMetrics, err := newMeterProvider(cfg, baseResource, auth, creds)
if err != nil {
return nil, err
}
meter := meterProvider.Meter(defaultPackageName)

// Shared log exporter for both logger and message emitter
logOpts, err := newLoggerOpts(cfg, auth, creds, meterProvider, tracerProvider)
// Shared export instruments beholder.export.bytes and
// beholder.export.duration, labelled per signal. They live on this
// MeterProvider, so the metrics exporter can only be wired up once the
// provider and its meter exist.
expMetrics, err := newExportMetrics(meter)
if err != nil {
return nil, err
}
sharedLogExporter, err := otlploggrpcNew(logOpts...)
meteredMetrics.attachMetrics(expMetrics, cfg.AuthPublicKeyHex)

// Shared log exporter for both logger and message emitter.
logOpts, err := newLoggerOpts(cfg, auth, creds)
if err != nil {
return nil, err
}
rawLogExporter, err := otlploggrpcNew(logOpts...)
if err != nil {
return nil, err
}
sharedLogExporter := newMeteredLogExporter(rawLogExporter, expMetrics, cfg.AuthPublicKeyHex)

// Logger
var loggerProvider *sdklog.LoggerProvider
Expand Down Expand Up @@ -505,12 +515,17 @@ func newTracerProvider(config Config, resource *sdkresource.Resource, auth Auth,
return sdktrace.NewTracerProvider(opts...), nil
}

func newMeterProvider(cfg Config, resource *sdkresource.Resource, auth Auth, creds credentials.TransportCredentials) (*sdkmetric.MeterProvider, error) {
func newMeterProvider(cfg Config, resource *sdkresource.Resource, auth Auth, creds credentials.TransportCredentials) (*sdkmetric.MeterProvider, *meteredMetricExporter, error) {
ctx := context.Background()
opts := []otlpmetricgrpc.Option{
otlpmetricgrpc.WithTLSCredentials(creds),
otlpmetricgrpc.WithEndpoint(cfg.OtelExporterGRPCEndpoint),
}

dialOpts := []grpc.DialOption{
grpc.WithStatsHandler(beholderStatsHandler{}),
}

switch compressor := cfg.MetricCompressor; compressor {
case "none":
case "":
Expand All @@ -522,13 +537,14 @@ func newMeterProvider(cfg Config, resource *sdkresource.Resource, auth Auth, cre
switch {
// Rotating auth
case auth != nil:
opts = append(opts, otlpmetricgrpc.WithDialOption(authDialOpt(auth)))
dialOpts = append(dialOpts, authDialOpt(auth))
// Static auth
case len(cfg.AuthHeaders) > 0:
opts = append(opts, otlpmetricgrpc.WithHeaders(cfg.AuthHeaders))
// No auth
default:
}
opts = append(opts, otlpmetricgrpc.WithDialOption(dialOpts...))

if cfg.MetricRetryConfig != nil {
// NOTE: By default, the retry is enabled in the OTel SDK
Expand All @@ -542,31 +558,30 @@ func newMeterProvider(cfg Config, resource *sdkresource.Resource, auth Auth, cre
// note: context is unused internally
exporter, err := otlpmetricgrpc.New(ctx, opts...)
if err != nil {
return nil, err
return nil, nil, err
}

metered := newMeteredMetricExporter(exporter)

readerOpts := []sdkmetric.PeriodicReaderOption{
sdkmetric.WithInterval(cfg.MetricReaderInterval), // Default is 10s
}
for _, p := range cfg.MetricProducers {
readerOpts = append(readerOpts, sdkmetric.WithProducer(p))
}

mpOpts := append(cfg.metricOptions(),
sdkmetric.WithReader(sdkmetric.NewPeriodicReader(exporter, readerOpts...)),
sdkmetric.WithReader(sdkmetric.NewPeriodicReader(metered, readerOpts...)),
sdkmetric.WithResource(resource),
)
return sdkmetric.NewMeterProvider(mpOpts...), nil
return sdkmetric.NewMeterProvider(mpOpts...), metered, nil
}

// newLoggerOpts creates options for a logger exporter
func newLoggerOpts(cfg Config, auth Auth, creds credentials.TransportCredentials, meter *sdkmetric.MeterProvider, tracer *sdktrace.TracerProvider) ([]otlploggrpc.Option, error) {
otelOpts := []otelgrpc.Option{
otelgrpc.WithMeterProvider(meter),
otelgrpc.WithTracerProvider(tracer),
}
func newLoggerOpts(cfg Config, auth Auth, creds credentials.TransportCredentials) ([]otlploggrpc.Option, error) {

dialOpts := []grpc.DialOption{
grpc.WithStatsHandler(otelgrpc.NewClientHandler(otelOpts...)),
grpc.WithStatsHandler(beholderStatsHandler{}),
}
Comment on lines 583 to 585

@kirqz23 kirqz23 Jul 20, 2026

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

gRPC builds istats.NewCombinedHandler(...) from the registered slice of multiple handlers, that is exactly the delegating handler proposed by the Copilot. This proposal is exactly what gRPC already builds for us internally. The downside of having multiple stats handlers however is that gRPC fires every handler on every event, in registration order. So for one OutPayload event there will be two HandleRPC(...) calls which might be an overhead. We might consider either keeping both statsHandlers if we need them, or dropping grpc.WithStatsHandler(otelgrpc.NewClientHandler(otelOpts...)).

This old handler:

  1. Emits rpc.client.* metrics (which part of them are going to be removed),
  2. Creates a trace span per RPC
  3. Injects trace-context headers

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

So, we have two approaches here:

  1. Keep two stats handlers as we have it now. exportSizeHandler shouldn't add much to the performance.
  2. Drop grpc.WithStatsHandler(otelgrpc.NewClientHandler(otelOpts...)) taking into account that we loose 3 points mentioned in the above comment, however rpc.client.* metrics are totally deprecated anyway expect duration, which can be added to our custom handler alongside beholder.export.bytes, e.g. sth like beholder.export.duration

cc @pkcll

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I decided to drop the otelgrpc stats handler and use only our custom beholderStatsHandler, cause it has now both size and duration metrics. We can later on expand it to produce more needed metrics.


opts := []otlploggrpc.Option{
Expand Down
94 changes: 94 additions & 0 deletions pkg/beholder/meter_provider_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,14 +2,17 @@ package beholder

import (
"context"
"sync"
"testing"
"time"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/metric"
sdkmetric "go.opentelemetry.io/otel/sdk/metric"
"go.opentelemetry.io/otel/sdk/metric/metricdata"
"google.golang.org/grpc/credentials/insecure"

"github.com/smartcontractkit/chainlink-common/pkg/beholder/metricviews"
)
Expand Down Expand Up @@ -89,3 +92,94 @@ func TestConfig_metricViews_emptyDenylistOmitsCatchAll(t *testing.T) {
require.Len(t, views, len(metricviews.Default(nil)))
require.NotEmpty(t, views)
}

// capturingMetricExporter records the ResourceMetrics handed to Export so a test
// can inspect what a MeterProvider's reader actually collected.
type capturingMetricExporter struct {
mu sync.Mutex
rm metricdata.ResourceMetrics
called bool
}

func (c *capturingMetricExporter) Temporality(k sdkmetric.InstrumentKind) metricdata.Temporality {
return sdkmetric.DefaultTemporalitySelector(k)
}

func (c *capturingMetricExporter) Aggregation(k sdkmetric.InstrumentKind) sdkmetric.Aggregation {
return sdkmetric.DefaultAggregationSelector(k)
}

func (c *capturingMetricExporter) Export(_ context.Context, rm *metricdata.ResourceMetrics) error {
c.mu.Lock()
defer c.mu.Unlock()
c.rm, c.called = *rm, true
return nil
}

func (c *capturingMetricExporter) ForceFlush(context.Context) error { return nil }
func (c *capturingMetricExporter) Shutdown(context.Context) error { return nil }

func (c *capturingMetricExporter) collected() (metricdata.ResourceMetrics, bool) {
c.mu.Lock()
defer c.mu.Unlock()
return c.rm, c.called
}

// TestNewMeterProvider_AppliesCardinalityLimit is the regression guard for the
// gRPC meter provider bypassing Config.metricOptions. Asserting on
// metricOptions alone cannot catch that, since the bug is the call site not
// using it.
func TestNewMeterProvider_AppliesCardinalityLimit(t *testing.T) {
t.Parallel()

const (
uniqueAttributes = 10
limit = 5
)

cfg := TestDefaultConfig()
cfg.MetricCardinalityLimit = limit
// Isolate the assertion from metricviews.Default, same as
// TestConfig_metricOptions_cardinalityLimit above.
cfg.metricViewsDisabled = true
// Long interval so only the explicit ForceFlush below triggers a collection.
cfg.MetricReaderInterval = time.Hour

resource, err := newOtelResource(cfg)
require.NoError(t, err)

mp, metered, err := newMeterProvider(cfg, resource, nil, insecure.NewCredentials())
require.NoError(t, err)
t.Cleanup(func() { _ = mp.Shutdown(context.Background()) })

// Swap the real OTLP exporter for one we can read. Shut the original down
// separately, since mp.Shutdown will now reach the capturing exporter.
capture := &capturingMetricExporter{}
original := metered.Exporter
metered.Exporter = capture
t.Cleanup(func() { _ = original.Shutdown(context.Background()) })

counter, err := mp.Meter("test").Int64Counter("overflow_test_total")
require.NoError(t, err)
for i := range uniqueAttributes {
counter.Add(context.Background(), 1, metric.WithAttributes(attribute.Int("key", i)))
}

require.NoError(t, mp.ForceFlush(context.Background()))

rm, called := capture.collected()
require.True(t, called, "ForceFlush should have exported a batch")

var sum metricdata.Sum[int64]
var found bool
for _, sm := range rm.ScopeMetrics {
for _, m := range sm.Metrics {
if m.Name == "overflow_test_total" {
sum, found = m.Data.(metricdata.Sum[int64]), true
}
}
}
require.True(t, found, "expected overflow_test_total in the collected batch")
assert.Len(t, sum.DataPoints, limit,
"newMeterProvider must honour MetricCardinalityLimit via Config.metricOptions")
}
Loading
Loading