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
43 changes: 33 additions & 10 deletions pkg/beholder/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -140,21 +140,31 @@ 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
// Shared export-size counter (beholder.export.bytes, labelled per signal).
// It lives on this MeterProvider, so the metrics exporter can only be wired
// up once the provider and its meter exist.
exportBytes, err := newExportBytesCounter(meter)
if err != nil {
return nil, err
}
meteredMetrics.attachCounter(exportBytes, cfg.AuthPublicKeyHex)

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

// Logger
var loggerProvider *sdklog.LoggerProvider
Expand Down Expand Up @@ -503,12 +513,20 @@ 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),
}
// All gRPC dial options must go through a single WithDialOption call:
// each otlpmetricgrpc.WithDialOption call replaces (does not append) the
// previously configured dial options, so multiple calls would clobber each
// other. Capture the outbound proto size of each metric export
// (beholder.export.bytes) here.
dialOpts := []grpc.DialOption{
grpc.WithStatsHandler(exportSizeHandler{}),
}
switch compressor := cfg.MetricCompressor; compressor {
case "none":
case "":
Expand All @@ -520,13 +538,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 @@ -540,20 +559,23 @@ 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
}
// Wrap so each metric export batch's proto size is recorded. The counter is
// attached later (attachCounter) once this provider's meter exists.
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...)),
return sdkmetric.NewMeterProvider(
sdkmetric.WithReader(sdkmetric.NewPeriodicReader(metered, readerOpts...)),
sdkmetric.WithResource(resource),
)
return sdkmetric.NewMeterProvider(mpOpts...), nil
sdkmetric.WithView(cfg.MetricViews...),
), metered, nil
}

// newLoggerOpts creates options for a logger exporter
Expand All @@ -565,6 +587,7 @@ func newLoggerOpts(cfg Config, auth Auth, creds credentials.TransportCredentials

dialOpts := []grpc.DialOption{
grpc.WithStatsHandler(otelgrpc.NewClientHandler(otelOpts...)),
grpc.WithStatsHandler(exportSizeHandler{}),
}
Comment on lines 588 to 591

@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


opts := []otlploggrpc.Option{
Expand Down
127 changes: 127 additions & 0 deletions pkg/beholder/metered_exporter.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,127 @@
package beholder

import (
"context"
"sync/atomic"

"go.opentelemetry.io/otel/attribute"
otelmetric "go.opentelemetry.io/otel/metric"
sdklog "go.opentelemetry.io/otel/sdk/log"
sdkmetric "go.opentelemetry.io/otel/sdk/metric"
"go.opentelemetry.io/otel/sdk/metric/metricdata"
"google.golang.org/grpc/stats"
)

// exportSizeKey is the context key under which a metered exporter stashes a
// per-export byte holder for exportSizeHandler to fill in.
type exportSizeKey struct{}

// exportSizeHandler is a minimal, stateless gRPC stats.Handler that records the
// uncompressed proto size of each outbound message.
type exportSizeHandler struct{}

func (exportSizeHandler) TagConn(ctx context.Context, _ *stats.ConnTagInfo) context.Context {
return ctx
}

func (exportSizeHandler) HandleConn(context.Context, stats.ConnStats) {}

func (exportSizeHandler) TagRPC(ctx context.Context, _ *stats.RPCTagInfo) context.Context {
return ctx
}

// HandleRPC fires on every gRPC stats event. On OutPayload it stores the
// uncompressed message length, the same field otelgrpc used for
// rpc.client.request.size
func (exportSizeHandler) HandleRPC(ctx context.Context, rs stats.RPCStats) {
op, ok := rs.(*stats.OutPayload)
if !ok {
return
}
if holder, ok := ctx.Value(exportSizeKey{}).(*atomic.Int64); ok {
holder.Store(int64(op.Length))
}
}

const exportBytesMetric = "beholder.export.bytes"
Comment thread
kirqz23 marked this conversation as resolved.

// newExportBytesCounter creates the counter shared by all metered exporters.
func newExportBytesCounter(meter otelmetric.Meter) (otelmetric.Int64Counter, error) {
return meter.Int64Counter(
exportBytesMetric,
otelmetric.WithDescription("Uncompressed OTLP proto size in bytes of each successful export batch, by signal."),
otelmetric.WithUnit("By"),
)
}

func exportAttrs(signal, csaPublicKeyHex string) otelmetric.MeasurementOption {
return otelmetric.WithAttributes(
attribute.String("otel_signal", signal),
attribute.String("csa_public_key", csaPublicKeyHex),
)
}

// meteredExporter holds the shared metering logic: run an export with a per-call
// size holder in the context, then record the captured OutPayload size on
// success.
type meteredExporter struct {
counter otelmetric.Int64Counter
attrs otelmetric.MeasurementOption
}

func (m meteredExporter) record(ctx context.Context, export func(context.Context) error) error {
var size atomic.Int64
err := export(context.WithValue(ctx, exportSizeKey{}, &size))
if err == nil {
m.counter.Add(ctx, size.Load(), m.attrs)
}
return err
}

// meteredLogExporter wraps an sdklog.Exporter and records each export batch's
// uncompressed proto size. It sits above the otlploggrpc retry loop, so Export is
// called once per logical batch and bytes are counted only on success.
type meteredLogExporter struct {
meteredExporter
inner sdklog.Exporter
}

func newMeteredLogExporter(inner sdklog.Exporter, counter otelmetric.Int64Counter, csaPublicKeyHex string) *meteredLogExporter {
return &meteredLogExporter{
meteredExporter: meteredExporter{counter: counter, attrs: exportAttrs("logs", csaPublicKeyHex)},
inner: inner,
}
}

func (e *meteredLogExporter) Export(ctx context.Context, records []sdklog.Record) error {
return e.record(ctx, func(c context.Context) error { return e.inner.Export(c, records) })
}

func (e *meteredLogExporter) Shutdown(ctx context.Context) error { return e.inner.Shutdown(ctx) }

func (e *meteredLogExporter) ForceFlush(ctx context.Context) error { return e.inner.ForceFlush(ctx) }

// meteredMetricExporter wraps an sdkmetric.Exporter.
// It is created by the MeterProvider and has no access to the counter until
// the MeterProvider exits and calls attachCounter.
type meteredMetricExporter struct {
sdkmetric.Exporter
base atomic.Pointer[meteredExporter]
}

func newMeteredMetricExporter(inner sdkmetric.Exporter) *meteredMetricExporter {
return &meteredMetricExporter{Exporter: inner}
}

// attachCounter wires the size counter once the MeterProvider exits.
func (e *meteredMetricExporter) attachCounter(counter otelmetric.Int64Counter, csaPublicKeyHex string) {
e.base.Store(&meteredExporter{counter: counter, attrs: exportAttrs("metrics", csaPublicKeyHex)})
}

func (e *meteredMetricExporter) Export(ctx context.Context, rm *metricdata.ResourceMetrics) error {
base := e.base.Load()
if base == nil {
return e.Exporter.Export(ctx, rm)
}
return base.record(ctx, func(c context.Context) error { return e.Exporter.Export(c, rm) })
}
Loading
Loading