diff --git a/pkg/beholder/client.go b/pkg/beholder/client.go index 7f52d911bf..ffa68f4eb9 100644 --- a/pkg/beholder/client.go +++ b/pkg/beholder/client.go @@ -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" @@ -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 @@ -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 "": @@ -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 @@ -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{}), } opts := []otlploggrpc.Option{ diff --git a/pkg/beholder/meter_provider_test.go b/pkg/beholder/meter_provider_test.go index 518d935799..e491d77bfa 100644 --- a/pkg/beholder/meter_provider_test.go +++ b/pkg/beholder/meter_provider_test.go @@ -2,7 +2,9 @@ package beholder import ( "context" + "sync" "testing" + "time" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" @@ -10,6 +12,7 @@ import ( "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" ) @@ -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") +} diff --git a/pkg/beholder/metered_exporter.go b/pkg/beholder/metered_exporter.go new file mode 100644 index 0000000000..65d0d0514b --- /dev/null +++ b/pkg/beholder/metered_exporter.go @@ -0,0 +1,208 @@ +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" + sdktrace "go.opentelemetry.io/otel/sdk/trace" + "google.golang.org/grpc/stats" +) + +// exportBytesKey is the context key under which a metered exporter stashes a +// per-export byte holder for beholderStatsHandler to fill in. +type exportBytesKey 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(exportBytesKey{}).(*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, exportBytesKey{}, &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) } + +// lazyMetered defers instrument wiring for exporters constructed before +// beholder MeterProvider exist +type lazyMetered struct { + signal string + base atomic.Pointer[meteredExporter] +} + +// attachMetrics wires the export instruments once the MeterProvider exists. +func (l *lazyMetered) attachMetrics(metrics exportMetrics, csaPublicKeyHex string) { + base := newBaseExporter(metrics, l.signal, csaPublicKeyHex) + l.base.Store(&base) +} + +func (l *lazyMetered) record(ctx context.Context, export func(context.Context) error) error { + base := l.base.Load() + if base == nil { + return export(ctx) + } + return base.record(ctx, export) +} + +// 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 + lazyMetered +} + +func newMeteredMetricExporter(inner sdkmetric.Exporter) *meteredMetricExporter { + return &meteredMetricExporter{Exporter: inner, lazyMetered: lazyMetered{signal: "metrics"}} +} + +func (e *meteredMetricExporter) Export(ctx context.Context, rm *metricdata.ResourceMetrics) error { + return e.record(ctx, func(c context.Context) error { return e.Exporter.Export(c, rm) }) +} + +type meteredTraceExporter struct { + sdktrace.SpanExporter + lazyMetered +} + +func newMeteredTraceExporter(inner sdktrace.SpanExporter) *meteredTraceExporter { + return &meteredTraceExporter{SpanExporter: inner, lazyMetered: lazyMetered{signal: "traces"}} +} + +func (e *meteredTraceExporter) ExportSpans(ctx context.Context, spans []sdktrace.ReadOnlySpan) error { + return e.record(ctx, func(c context.Context) error { return e.SpanExporter.ExportSpans(c, spans) }) +} diff --git a/pkg/beholder/metered_exporter_test.go b/pkg/beholder/metered_exporter_test.go new file mode 100644 index 0000000000..6db20bd731 --- /dev/null +++ b/pkg/beholder/metered_exporter_test.go @@ -0,0 +1,616 @@ +package beholder + +import ( + "context" + "errors" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "go.opentelemetry.io/otel/attribute" + sdklog "go.opentelemetry.io/otel/sdk/log" + sdkmetric "go.opentelemetry.io/otel/sdk/metric" + "go.opentelemetry.io/otel/sdk/metric/metricdata" + sdktrace "go.opentelemetry.io/otel/sdk/trace" + "google.golang.org/grpc/stats" +) + +// --- test doubles --------------------------------------------------------- + +// fakeLogExporter stands in for the otlploggrpc exporter. On Export it mimics +// what the real gRPC stack does +type fakeLogExporter struct { + size int // OutPayload length to report + sizeFromRecords bool // if true, report len(records) instead of size + fireN int // number of OutPayload events (retry simulation); 0 => 1 + delay time.Duration // widen the store→read window for the race test + err error + + shutdownCalled atomic.Bool + forceFlushCalled atomic.Bool +} + +func (f *fakeLogExporter) Export(ctx context.Context, records []sdklog.Record) error { + sz := f.size + if f.sizeFromRecords { + sz = len(records) + } + n := f.fireN + if n == 0 { + n = 1 + } + for i := 0; i < n; i++ { + beholderStatsHandler{}.HandleRPC(ctx, &stats.OutPayload{Length: sz}) + } + if f.delay > 0 { + time.Sleep(f.delay) + } + return f.err +} + +func (f *fakeLogExporter) Shutdown(context.Context) error { f.shutdownCalled.Store(true); return nil } +func (f *fakeLogExporter) ForceFlush(context.Context) error { + f.forceFlushCalled.Store(true) + return nil +} + +// fakeMetricExporter stands in for the otlpmetricgrpc exporter. +type fakeMetricExporter struct { + size int + delay time.Duration + err error + + temporalityCalled atomic.Bool + aggregationCalled atomic.Bool + shutdownCalled atomic.Bool + forceFlushCalled atomic.Bool +} + +func (f *fakeMetricExporter) Temporality(sdkmetric.InstrumentKind) metricdata.Temporality { + f.temporalityCalled.Store(true) + return metricdata.CumulativeTemporality +} + +func (f *fakeMetricExporter) Aggregation(k sdkmetric.InstrumentKind) sdkmetric.Aggregation { + f.aggregationCalled.Store(true) + return sdkmetric.DefaultAggregationSelector(k) +} + +func (f *fakeMetricExporter) Export(ctx context.Context, _ *metricdata.ResourceMetrics) error { + beholderStatsHandler{}.HandleRPC(ctx, &stats.OutPayload{Length: f.size}) + if f.delay > 0 { + time.Sleep(f.delay) + } + return f.err +} + +func (f *fakeMetricExporter) ForceFlush(context.Context) error { + f.forceFlushCalled.Store(true) + return nil +} +func (f *fakeMetricExporter) Shutdown(context.Context) error { + f.shutdownCalled.Store(true) + return nil +} + +// fakeTraceExporter stands in for the otlptracegrpc exporter. +type fakeTraceExporter struct { + size int + fireN int // number of OutPayload events (retry simulation); 0 => 1 + delay time.Duration + err error + + shutdownCalled atomic.Bool +} + +func (f *fakeTraceExporter) ExportSpans(ctx context.Context, _ []sdktrace.ReadOnlySpan) error { + n := f.fireN + if n == 0 { + n = 1 + } + for range n { + beholderStatsHandler{}.HandleRPC(ctx, &stats.OutPayload{Length: f.size}) + } + if f.delay > 0 { + time.Sleep(f.delay) + } + return f.err +} + +func (f *fakeTraceExporter) Shutdown(context.Context) error { + f.shutdownCalled.Store(true) + return nil +} + +// --- helpers -------------------------------------------------------------- + +// newTestMetrics wires the export instruments to an in-memory ManualReader and +// returns them plus a collect func that reads back the recorded metrics. +func newTestMetrics(t *testing.T) (exportMetrics, func() []metricdata.Metrics) { + t.Helper() + reader := sdkmetric.NewManualReader() + mp := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader)) + metrics, err := newExportMetrics(mp.Meter("test")) + require.NoError(t, err) + + collect := func() []metricdata.Metrics { + var rm metricdata.ResourceMetrics + require.NoError(t, reader.Collect(context.Background(), &rm)) + var out []metricdata.Metrics + for _, sm := range rm.ScopeMetrics { + out = append(out, sm.Metrics...) + } + return out + } + return metrics, collect +} + +func dpForSignal(t *testing.T, ms []metricdata.Metrics, signal string) (metricdata.DataPoint[int64], bool) { + t.Helper() + for _, m := range ms { + if m.Name != exportBytesMetric { + continue + } + sum, ok := m.Data.(metricdata.Sum[int64]) + require.True(t, ok, "expected beholder.export.bytes to be Sum[int64]") + for _, dp := range sum.DataPoints { + if v, ok := dp.Attributes.Value(attribute.Key("otel_signal")); ok && v.AsString() == signal { + return dp, true + } + } + } + return metricdata.DataPoint[int64]{}, false +} + +// histForSignal finds the beholder.export.duration datapoint for one signal and +// error outcome. +func histForSignal(t *testing.T, ms []metricdata.Metrics, signal string, wantErr bool) (metricdata.HistogramDataPoint[float64], bool) { + t.Helper() + for _, m := range ms { + if m.Name != exportDurationMetric { + continue + } + h, ok := m.Data.(metricdata.Histogram[float64]) + require.True(t, ok, "expected beholder.export.duration to be Histogram[float64]") + for _, dp := range h.DataPoints { + sig, ok := dp.Attributes.Value(attribute.Key("otel_signal")) + if !ok || sig.AsString() != signal { + continue + } + isErr, ok := dp.Attributes.Value(attribute.Key("error")) + require.True(t, ok, "duration datapoint must carry an error attribute") + if isErr.AsBool() == wantErr { + return dp, true + } + } + } + return metricdata.HistogramDataPoint[float64]{}, false +} + +// --- beholderStatsHandler ----------------------------------------------- + +func TestBeholderStatsHandler_StoresOutPayloadLength(t *testing.T) { + var holder atomic.Int64 + ctx := context.WithValue(context.Background(), exportBytesKey{}, &holder) + + beholderStatsHandler{}.HandleRPC(ctx, &stats.OutPayload{Length: 1234}) + + assert.Equal(t, int64(1234), holder.Load()) +} + +func TestBeholderStatsHandler_IgnoresNonOutPayload(t *testing.T) { + var holder atomic.Int64 + ctx := context.WithValue(context.Background(), exportBytesKey{}, &holder) + + beholderStatsHandler{}.HandleRPC(ctx, &stats.InPayload{Length: 999}) + beholderStatsHandler{}.HandleRPC(ctx, &stats.Begin{}) + beholderStatsHandler{}.HandleRPC(ctx, &stats.End{}) + + assert.Equal(t, int64(0), holder.Load()) +} + +func TestBeholderStatsHandler_LastWriteWins(t *testing.T) { + // Retries resend the same proto, so each OutPayload reports the same length. + // Store means the holder ends at that length, not a multiple of it. + var holder atomic.Int64 + ctx := context.WithValue(context.Background(), exportBytesKey{}, &holder) + h := beholderStatsHandler{} + + h.HandleRPC(ctx, &stats.OutPayload{Length: 700}) + h.HandleRPC(ctx, &stats.OutPayload{Length: 700}) + h.HandleRPC(ctx, &stats.OutPayload{Length: 700}) + + assert.Equal(t, int64(700), holder.Load()) +} + +func TestBeholderStatsHandler_NoHolderInContextIsNoop(t *testing.T) { + assert.NotPanics(t, func() { + beholderStatsHandler{}.HandleRPC(context.Background(), &stats.OutPayload{Length: 5}) + }) +} + +func TestBeholderStatsHandler_TagAndConnAreInert(t *testing.T) { + h := beholderStatsHandler{} + ctx := context.WithValue(context.Background(), exportBytesKey{}, &atomic.Int64{}) + + assert.Equal(t, ctx, h.TagRPC(ctx, &stats.RPCTagInfo{})) + assert.Equal(t, ctx, h.TagConn(ctx, &stats.ConnTagInfo{})) + assert.NotPanics(t, func() { h.HandleConn(ctx, &stats.ConnBegin{}) }) +} + +// --- meteredLogExporter --------------------------------------------------- + +func TestMeteredLogExporter_RecordsBytesOnSuccess(t *testing.T) { + metrics, collect := newTestMetrics(t) + inner := &fakeLogExporter{size: 4096} + exp := newMeteredLogExporter(inner, metrics, "csa-pub-hex") + + require.NoError(t, exp.Export(context.Background(), nil)) + + dp, ok := dpForSignal(t, collect(), "logs") + require.True(t, ok, "expected a logs datapoint") + assert.Equal(t, int64(4096), dp.Value) + + csa, ok := dp.Attributes.Value(attribute.Key("csa_public_key")) + require.True(t, ok) + assert.Equal(t, "csa-pub-hex", csa.AsString()) +} + +func TestMeteredLogExporter_NoRecordOnError(t *testing.T) { + metrics, collect := newTestMetrics(t) + // Handler still fires, but Export returns an error. + inner := &fakeLogExporter{size: 4096, err: errors.New("boom")} + exp := newMeteredLogExporter(inner, metrics, "csa") + + require.Error(t, exp.Export(context.Background(), nil)) + + ms := collect() + _, ok := dpForSignal(t, ms, "logs") + assert.False(t, ok, "no bytes should be recorded when export fails") + // Duration is still recorded, labelled error=true. + hdp, ok := histForSignal(t, ms, "logs", true) + require.True(t, ok, "expected an error-labelled logs duration datapoint") + assert.Equal(t, uint64(1), hdp.Count) + _, ok = histForSignal(t, ms, "logs", false) + assert.False(t, ok, "a failed export must not be labelled error=false") +} + +func TestMeteredLogExporter_RetriesCountedOnce(t *testing.T) { + metrics, collect := newTestMetrics(t) + // Three OutPayload events for one batch, all the same size. + inner := &fakeLogExporter{size: 500, fireN: 3} + exp := newMeteredLogExporter(inner, metrics, "csa") + + require.NoError(t, exp.Export(context.Background(), nil)) + + dp, ok := dpForSignal(t, collect(), "logs") + require.True(t, ok) + assert.Equal(t, int64(500), dp.Value, "retries of the same batch must be counted once") +} + +func TestMeteredLogExporter_Passthrough(t *testing.T) { + inner := &fakeLogExporter{} + exp := newMeteredLogExporter(inner, exportMetrics{}, "csa") + + require.NoError(t, exp.Shutdown(context.Background())) + require.NoError(t, exp.ForceFlush(context.Background())) + + assert.True(t, inner.shutdownCalled.Load()) + assert.True(t, inner.forceFlushCalled.Load()) +} + +// TestMeteredExporter_ConcurrentExportsIsolated is the regression guard for the +// per-call context holder +func TestMeteredExporter_ConcurrentExportsIsolated(t *testing.T) { + metrics, collect := newTestMetrics(t) + // delay widens the window between the handler storing and record reading, + // so a broken implementation would reliably mis-attribute. + inner := &fakeLogExporter{sizeFromRecords: true, delay: 200 * time.Microsecond} + exp := newMeteredLogExporter(inner, metrics, "csa") + + const n = 50 + var wg sync.WaitGroup + var want int64 + for i := 1; i <= n; i++ { + size := i + want += int64(size) + wg.Add(1) + go func() { + defer wg.Done() + assert.NoError(t, exp.Export(context.Background(), make([]sdklog.Record, size))) + }() + } + wg.Wait() + + dp, ok := dpForSignal(t, collect(), "logs") + require.True(t, ok) + assert.Equal(t, want, dp.Value) +} + +// --- meteredMetricExporter ------------------------------------------------ + +func TestMeteredMetricExporter_RecordsBytesOnSuccess(t *testing.T) { + metrics, collect := newTestMetrics(t) + inner := &fakeMetricExporter{size: 8192} + exp := newMeteredMetricExporter(inner) + exp.attachMetrics(metrics, "csa-pub-hex") + + require.NoError(t, exp.Export(context.Background(), &metricdata.ResourceMetrics{})) + + dp, ok := dpForSignal(t, collect(), "metrics") + require.True(t, ok, "expected a metrics datapoint") + assert.Equal(t, int64(8192), dp.Value) + + csa, ok := dp.Attributes.Value(attribute.Key("csa_public_key")) + require.True(t, ok) + assert.Equal(t, "csa-pub-hex", csa.AsString()) +} + +func TestMeteredMetricExporter_UnmeteredBeforeAttach(t *testing.T) { + _, collect := newTestMetrics(t) + inner := &fakeMetricExporter{size: 8192} + exp := newMeteredMetricExporter(inner) // no attachMetrics + + require.NoError(t, exp.Export(context.Background(), &metricdata.ResourceMetrics{})) + + ms := collect() + _, ok := dpForSignal(t, ms, "metrics") + assert.False(t, ok, "no bytes should be recorded before the instruments are attached") + _, ok = histForSignal(t, ms, "metrics", false) + assert.False(t, ok, "no duration should be recorded before the instruments are attached") +} + +func TestMeteredMetricExporter_NoRecordOnError(t *testing.T) { + metrics, collect := newTestMetrics(t) + inner := &fakeMetricExporter{size: 8192, err: errors.New("boom")} + exp := newMeteredMetricExporter(inner) + exp.attachMetrics(metrics, "csa") + + require.Error(t, exp.Export(context.Background(), &metricdata.ResourceMetrics{})) + + ms := collect() + _, ok := dpForSignal(t, ms, "metrics") + assert.False(t, ok) + hdp, ok := histForSignal(t, ms, "metrics", true) + require.True(t, ok, "expected an error-labelled metrics duration datapoint") + assert.Equal(t, uint64(1), hdp.Count) +} + +func TestMeteredMetricExporter_Passthrough(t *testing.T) { + inner := &fakeMetricExporter{} + exp := newMeteredMetricExporter(inner) + + assert.Equal(t, metricdata.CumulativeTemporality, exp.Temporality(sdkmetric.InstrumentKindCounter)) + assert.NotNil(t, exp.Aggregation(sdkmetric.InstrumentKindCounter)) + require.NoError(t, exp.ForceFlush(context.Background())) + require.NoError(t, exp.Shutdown(context.Background())) + + assert.True(t, inner.temporalityCalled.Load()) + assert.True(t, inner.aggregationCalled.Load()) + assert.True(t, inner.forceFlushCalled.Load()) + assert.True(t, inner.shutdownCalled.Load()) +} + +// --- meteredTraceExporter -------------------------------------------------- +// +// The trace exporter is not wired into NewGRPCClient yet; these lock in the +// behaviour so enabling it later is a wiring change only. + +func TestMeteredTraceExporter_RecordsBytesAndDurationOnSuccess(t *testing.T) { + inst, collect := newTestMetrics(t) + const delay = 20 * time.Millisecond + exp := newMeteredTraceExporter(&fakeTraceExporter{size: 2048, delay: delay}) + exp.attachMetrics(inst, "csa-pub-hex") + + require.NoError(t, exp.ExportSpans(context.Background(), nil)) + + ms := collect() + dp, ok := dpForSignal(t, ms, "traces") + require.True(t, ok, "expected a traces bytes datapoint") + assert.Equal(t, int64(2048), dp.Value) + + csa, ok := dp.Attributes.Value(attribute.Key("csa_public_key")) + require.True(t, ok) + assert.Equal(t, "csa-pub-hex", csa.AsString()) + + hdp, ok := histForSignal(t, ms, "traces", false) + require.True(t, ok, "expected a traces duration datapoint") + assert.Equal(t, uint64(1), hdp.Count) + assert.GreaterOrEqual(t, hdp.Sum, delay.Seconds()) +} + +func TestMeteredTraceExporter_NoBytesOnErrorButDurationRecorded(t *testing.T) { + inst, collect := newTestMetrics(t) + exp := newMeteredTraceExporter(&fakeTraceExporter{size: 2048, err: errors.New("boom")}) + exp.attachMetrics(inst, "csa") + + require.Error(t, exp.ExportSpans(context.Background(), nil)) + + ms := collect() + _, ok := dpForSignal(t, ms, "traces") + assert.False(t, ok, "no bytes should be recorded when export fails") + hdp, ok := histForSignal(t, ms, "traces", true) + require.True(t, ok, "expected an error-labelled traces duration datapoint") + assert.Equal(t, uint64(1), hdp.Count) +} + +func TestMeteredTraceExporter_RetriesCountedOnce(t *testing.T) { + inst, collect := newTestMetrics(t) + exp := newMeteredTraceExporter(&fakeTraceExporter{size: 700, fireN: 3}) + exp.attachMetrics(inst, "csa") + + require.NoError(t, exp.ExportSpans(context.Background(), nil)) + + ms := collect() + dp, ok := dpForSignal(t, ms, "traces") + require.True(t, ok) + assert.Equal(t, int64(700), dp.Value, "retries of the same batch must be counted once") + hdp, ok := histForSignal(t, ms, "traces", false) + require.True(t, ok) + assert.Equal(t, uint64(1), hdp.Count, "one batch must be one duration observation") +} + +// TestMeteredTraceExporter_UnmeteredBeforeAttach is the guard for the deferred +// attach: newTracerProvider runs before the MeterProvider exists, so an +// unattached exporter must still export rather than panic. +func TestMeteredTraceExporter_UnmeteredBeforeAttach(t *testing.T) { + _, collect := newTestMetrics(t) + inner := &fakeTraceExporter{size: 2048} + exp := newMeteredTraceExporter(inner) // no attachMetrics + + require.NoError(t, exp.ExportSpans(context.Background(), nil)) + + ms := collect() + _, ok := dpForSignal(t, ms, "traces") + assert.False(t, ok, "no bytes before the instruments are attached") + _, ok = histForSignal(t, ms, "traces", false) + assert.False(t, ok, "no duration before the instruments are attached") +} + +func TestMeteredTraceExporter_Passthrough(t *testing.T) { + inner := &fakeTraceExporter{} + exp := newMeteredTraceExporter(inner) + + require.NoError(t, exp.Shutdown(context.Background())) + + assert.True(t, inner.shutdownCalled.Load()) +} + +// --- shared naming -------------------------------------------------------- + +func TestMeteredExporters_ShareOneMetricBySignal(t *testing.T) { + inst, collect := newTestMetrics(t) + logs := newMeteredLogExporter(&fakeLogExporter{size: 100}, inst, "csa") + metrics := newMeteredMetricExporter(&fakeMetricExporter{size: 200}) + metrics.attachMetrics(inst, "csa") + traces := newMeteredTraceExporter(&fakeTraceExporter{size: 300}) + traces.attachMetrics(inst, "csa") + + require.NoError(t, logs.Export(context.Background(), nil)) + require.NoError(t, metrics.Export(context.Background(), &metricdata.ResourceMetrics{})) + require.NoError(t, traces.ExportSpans(context.Background(), nil)) + + ms := collect() + for _, tc := range []struct { + signal string + want int64 + }{{"logs", 100}, {"metrics", 200}, {"traces", 300}} { + dp, ok := dpForSignal(t, ms, tc.signal) + require.True(t, ok, "expected a %s bytes datapoint", tc.signal) + assert.Equal(t, tc.want, dp.Value) + + // Each signal gets its own duration datapoint on the same instrument too. + hdp, ok := histForSignal(t, ms, tc.signal, false) + require.True(t, ok, "expected a %s duration datapoint", tc.signal) + assert.Equal(t, uint64(1), hdp.Count) + } + + // Both are datapoints of the same instrument, distinguished only by otel_signal. + for _, m := range ms { + switch m.Name { + case exportBytesMetric: + assert.Equal(t, "By", m.Unit) + case exportDurationMetric: + assert.Equal(t, "s", m.Unit) + } + } +} + +// --- beholder.export.duration --------------------------------------------- + +func TestMeteredLogExporter_RecordsDurationOnSuccess(t *testing.T) { + inst, collect := newTestMetrics(t) + const delay = 20 * time.Millisecond + inner := &fakeLogExporter{size: 4096, delay: delay} + exp := newMeteredLogExporter(inner, inst, "csa-pub-hex") + + require.NoError(t, exp.Export(context.Background(), nil)) + + hdp, ok := histForSignal(t, collect(), "logs", false) + require.True(t, ok, "expected a logs duration datapoint") + assert.Equal(t, uint64(1), hdp.Count) + assert.GreaterOrEqual(t, hdp.Sum, delay.Seconds(), + "recorded duration must cover the inner export") + assert.Less(t, hdp.Sum, 10.0, "recorded duration should be in seconds, not another unit") + + csa, ok := hdp.Attributes.Value(attribute.Key("csa_public_key")) + require.True(t, ok) + assert.Equal(t, "csa-pub-hex", csa.AsString()) +} + +func TestMeteredMetricExporter_RecordsDurationOnSuccess(t *testing.T) { + inst, collect := newTestMetrics(t) + const delay = 20 * time.Millisecond + exp := newMeteredMetricExporter(&fakeMetricExporter{size: 8192, delay: delay}) + exp.attachMetrics(inst, "csa-pub-hex") + + require.NoError(t, exp.Export(context.Background(), &metricdata.ResourceMetrics{})) + + hdp, ok := histForSignal(t, collect(), "metrics", false) + require.True(t, ok, "expected a metrics duration datapoint") + assert.Equal(t, uint64(1), hdp.Count) + assert.GreaterOrEqual(t, hdp.Sum, delay.Seconds()) + + csa, ok := hdp.Attributes.Value(attribute.Key("csa_public_key")) + require.True(t, ok) + assert.Equal(t, "csa-pub-hex", csa.AsString()) +} + +// TestMeteredExporter_DurationRetriesCountedOnce mirrors the bytes behaviour: +// the wrapper sits above the otlp retry loop, so one logical batch is one +// observation covering the whole retry sequence. +func TestMeteredExporter_DurationRetriesCountedOnce(t *testing.T) { + inst, collect := newTestMetrics(t) + inner := &fakeLogExporter{size: 500, fireN: 3} + exp := newMeteredLogExporter(inner, inst, "csa") + + require.NoError(t, exp.Export(context.Background(), nil)) + + hdp, ok := histForSignal(t, collect(), "logs", false) + require.True(t, ok) + assert.Equal(t, uint64(1), hdp.Count, "one batch must be one duration observation") +} + +func TestExportDuration_UsesSecondScaledBuckets(t *testing.T) { + inst, collect := newTestMetrics(t) + exp := newMeteredLogExporter(&fakeLogExporter{}, inst, "csa") + + require.NoError(t, exp.Export(context.Background(), nil)) + + for _, m := range collect() { + if m.Name != exportDurationMetric { + continue + } + h, ok := m.Data.(metricdata.Histogram[float64]) + require.True(t, ok) + require.NotEmpty(t, h.DataPoints) + assert.Equal(t, + []float64{0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10, 30, 60}, + h.DataPoints[0].Bounds, + "explicit second-scaled boundaries must survive to the reader") + return + } + t.Fatalf("no %s metric collected", exportDurationMetric) +} + +func TestMeteredExporter_DurationRecordedPerOutcome(t *testing.T) { + inst, collect := newTestMetrics(t) + ok1 := newMeteredLogExporter(&fakeLogExporter{size: 10}, inst, "csa") + bad := newMeteredLogExporter(&fakeLogExporter{size: 10, err: errors.New("boom")}, inst, "csa") + + require.NoError(t, ok1.Export(context.Background(), nil)) + require.NoError(t, ok1.Export(context.Background(), nil)) + require.Error(t, bad.Export(context.Background(), nil)) + + ms := collect() + okDP, found := histForSignal(t, ms, "logs", false) + require.True(t, found) + assert.Equal(t, uint64(2), okDP.Count) + + errDP, found := histForSignal(t, ms, "logs", true) + require.True(t, found) + assert.Equal(t, uint64(1), errDP.Count) +}