diff --git a/pkg/beholder/client.go b/pkg/beholder/client.go index 73579eccc7..acf6c79c78 100644 --- a/pkg/beholder/client.go +++ b/pkg/beholder/client.go @@ -140,21 +140,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 @@ -503,12 +514,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 "": @@ -520,13 +539,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 @@ -540,8 +560,12 @@ 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 and duration are recorded. + // The instruments are attached later once this provider's + // meter exists. + metered := newMeteredMetricExporter(exporter) readerOpts := []sdkmetric.PeriodicReaderOption{ sdkmetric.WithInterval(cfg.MetricReaderInterval), // Default is 10s @@ -549,11 +573,11 @@ func newMeterProvider(cfg Config, resource *sdkresource.Resource, auth Auth, cre 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 @@ -565,6 +589,7 @@ func newLoggerOpts(cfg Config, auth Auth, creds credentials.TransportCredentials dialOpts := []grpc.DialOption{ grpc.WithStatsHandler(otelgrpc.NewClientHandler(otelOpts...)), + grpc.WithStatsHandler(exportSizeHandler{}), } opts := []otlploggrpc.Option{ diff --git a/pkg/beholder/metered_exporter.go b/pkg/beholder/metered_exporter.go new file mode 100644 index 0000000000..2818a79c0b --- /dev/null +++ b/pkg/beholder/metered_exporter.go @@ -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 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" + 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) }) +} diff --git a/pkg/beholder/metered_exporter_test.go b/pkg/beholder/metered_exporter_test.go new file mode 100644 index 0000000000..a8db3d0892 --- /dev/null +++ b/pkg/beholder/metered_exporter_test.go @@ -0,0 +1,498 @@ +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" + "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++ { + exportSizeHandler{}.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 { + exportSizeHandler{}.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 +} + +// --- 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 +} + +// --- exportSizeHandler --------------------------------------------------- + +func TestExportSizeHandler_StoresOutPayloadLength(t *testing.T) { + var holder atomic.Int64 + ctx := context.WithValue(context.Background(), exportSizeKey{}, &holder) + + exportSizeHandler{}.HandleRPC(ctx, &stats.OutPayload{Length: 1234}) + + assert.Equal(t, int64(1234), holder.Load()) +} + +func TestExportSizeHandler_IgnoresNonOutPayload(t *testing.T) { + var holder atomic.Int64 + ctx := context.WithValue(context.Background(), exportSizeKey{}, &holder) + + exportSizeHandler{}.HandleRPC(ctx, &stats.InPayload{Length: 999}) + exportSizeHandler{}.HandleRPC(ctx, &stats.Begin{}) + exportSizeHandler{}.HandleRPC(ctx, &stats.End{}) + + assert.Equal(t, int64(0), holder.Load()) +} + +func TestExportSizeHandler_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(), exportSizeKey{}, &holder) + h := exportSizeHandler{} + + 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 TestExportSizeHandler_NoHolderInContextIsNoop(t *testing.T) { + assert.NotPanics(t, func() { + exportSizeHandler{}.HandleRPC(context.Background(), &stats.OutPayload{Length: 5}) + }) +} + +func TestExportSizeHandler_TagAndConnAreInert(t *testing.T) { + h := exportSizeHandler{} + ctx := context.WithValue(context.Background(), exportSizeKey{}, &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()) +} + +// --- 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") + + require.NoError(t, logs.Export(context.Background(), nil)) + require.NoError(t, metrics.Export(context.Background(), &metricdata.ResourceMetrics{})) + + ms := collect() + logDP, ok := dpForSignal(t, ms, "logs") + require.True(t, ok) + assert.Equal(t, int64(100), logDP.Value) + metricDP, ok := dpForSignal(t, ms, "metrics") + require.True(t, ok) + assert.Equal(t, int64(200), metricDP.Value) + + // Each signal gets its own duration datapoint on the same instrument too. + for _, signal := range []string{"logs", "metrics"} { + hdp, ok := histForSignal(t, ms, signal, false) + require.True(t, ok, "expected a %s duration datapoint", 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) +}