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
46 changes: 30 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
// 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
}
meteredMetrics.attachMetrics(expMetrics, 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, expMetrics, cfg.AuthPublicKeyHex)

// Logger
var loggerProvider *sdklog.LoggerProvider
Expand Down Expand Up @@ -503,12 +513,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 @@ -520,13 +535,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,31 +556,29 @@ 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...)),
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
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),
}

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

@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
183 changes: 183 additions & 0 deletions pkg/beholder/metered_exporter.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,183 @@
package beholder

import (
"context"
"sync/atomic"
"time"

"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 beholderStatsHandler to fill in.
type exportSizeKey struct{}

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

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

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

func (beholderStatsHandler) 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 (beholderStatsHandler) 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"
exportDurationMetric = "beholder.export.duration"
)

// exportMetrics holds the instruments shared by all metered exporters. They live
// on the beholder MeterProvider and are distinguished per exporter only by
// attributes, so one set covers every signal.
type exportMetrics struct {
bytes otelmetric.Int64Counter
duration otelmetric.Float64Histogram
}

// newExportMetrics creates the instruments shared by all metered exporters.
func newExportMetrics(meter otelmetric.Meter) (exportMetrics, error) {
bytes, err := meter.Int64Counter(
exportBytesMetric,
otelmetric.WithDescription("Uncompressed OTLP proto size in bytes of each export batch. Recorded once per batch, on success only; retry attempts are not summed."),
otelmetric.WithUnit("By"),
)
if err != nil {
return exportMetrics{}, err
}
duration, err := meter.Float64Histogram(
exportDurationMetric,
otelmetric.WithDescription("Wall-clock duration in seconds of each OTLP export batch, covering all retry attempts and backoff. Recorded once per batch, on both success and failure."),
otelmetric.WithUnit("s"),
// Sized for network exports: sub-10ms to a 60s deadline. The SDK defaults
// are millisecond-scaled, so nearly every export would land in bucket one.
otelmetric.WithExplicitBucketBoundaries(
0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10, 30, 60,
),
)
if err != nil {
return exportMetrics{}, err
}
return exportMetrics{bytes: bytes, duration: duration}, nil
}

// exportAttrs builds the attribute set identifying one signal's exports, plus
// any per-measurement extras.
func exportAttrs(signal, csaPublicKeyHex string, extra ...attribute.KeyValue) otelmetric.MeasurementOption {
attrs := make([]attribute.KeyValue, 0, 2+len(extra))
attrs = append(attrs,
attribute.String("otel_signal", signal),
attribute.String("csa_public_key", csaPublicKeyHex),
)
return otelmetric.WithAttributes(append(attrs, extra...)...)
}

// 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 and the duration either way. Attribute options are
// precomputed so the export path allocates nothing per call.
type meteredExporter struct {
metrics exportMetrics

byteAttrs otelmetric.MeasurementOption // otel_signal, csa_public_key
okAttrs otelmetric.MeasurementOption // + error=false
errAttrs otelmetric.MeasurementOption // + error=true
}

func newBaseExporter(metrics exportMetrics, signal, csaPublicKeyHex string) meteredExporter {
return meteredExporter{
metrics: metrics,
byteAttrs: exportAttrs(signal, csaPublicKeyHex),
okAttrs: exportAttrs(signal, csaPublicKeyHex, attribute.Bool("error", false)),
errAttrs: exportAttrs(signal, csaPublicKeyHex, attribute.Bool("error", true)),
}
}

func (m meteredExporter) record(ctx context.Context, export func(context.Context) error) error {
var size atomic.Int64
start := time.Now()
err := export(context.WithValue(ctx, exportSizeKey{}, &size))
elapsed := time.Since(start).Seconds()

// Bytes are only meaningful for a batch that landed; duration is recorded
// either way
if err == nil {
m.metrics.bytes.Add(ctx, size.Load(), m.byteAttrs)
m.metrics.duration.Record(ctx, elapsed, m.okAttrs)
} else {
m.metrics.duration.Record(ctx, elapsed, m.errAttrs)
}
return err
}

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

func newMeteredLogExporter(inner sdklog.Exporter, metrics exportMetrics, csaPublicKeyHex string) *meteredLogExporter {
return &meteredLogExporter{
meteredExporter: newBaseExporter(metrics, "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 instruments until
// the MeterProvider exists and calls attachMetrics.
type meteredMetricExporter struct {
sdkmetric.Exporter
base atomic.Pointer[meteredExporter]
}

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

// attachMetrics wires the export instruments once the MeterProvider exists.
func (e *meteredMetricExporter) attachMetrics(metrics exportMetrics, csaPublicKeyHex string) {
base := newBaseExporter(metrics, "metrics", csaPublicKeyHex)
e.base.Store(&base)
}

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