From 792856ef954ed6d84bfcfdfaf810f72f1c1b627f Mon Sep 17 00:00:00 2001 From: Dylan Myers Date: Tue, 11 Aug 2026 16:08:05 -0400 Subject: [PATCH] fix(o11y): record output request + HEC ACK-poll latency in ms (PIPE-1404) Assisted-by: Claude Opus 4.8 --- output/duration_test.go | 54 +++++++++++++++++++++++++++++++ output/hec/ack.go | 2 +- output/hec/ackpoll_unit_test.go | 34 +++++++++++++++++++ output/hec/hec.go | 2 +- output/hec/metrics.go | 8 ++--- output/hec/monitoring.go | 2 +- output/hec/monitoring.md | 4 +-- output/hec/monitoring/metric.yaml | 2 +- output/monitoring.go | 2 +- output/monitoring.md | 4 +-- output/monitoring/metric.yaml | 2 +- output/otlp_grpc/otlp_grpc.go | 6 ++-- output/tcp/tcp.go | 2 +- output/util.go | 9 ++++++ 14 files changed, 115 insertions(+), 18 deletions(-) create mode 100644 output/duration_test.go create mode 100644 output/hec/ackpoll_unit_test.go diff --git a/output/duration_test.go b/output/duration_test.go new file mode 100644 index 0000000..bd12a3e --- /dev/null +++ b/output/duration_test.go @@ -0,0 +1,54 @@ +package output_test + +import ( + "context" + "testing" + "time" + + "github.com/observiq/blitz/output" + "github.com/stretchr/testify/require" + sdkmetric "go.opentelemetry.io/otel/sdk/metric" + "go.opentelemetry.io/otel/sdk/metric/metricdata" +) + +func TestDurationMillis(t *testing.T) { + cases := []struct { + name string + in time.Duration + want float64 + }{ + {"whole milliseconds", 250 * time.Millisecond, 250}, + {"seconds scale", 2 * time.Second, 2000}, + {"sub-millisecond preserved", 500 * time.Microsecond, 0.5}, + {"zero", 0, 0}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + require.InDelta(t, tc.want, output.DurationMillis(tc.in), 1e-9) + }) + } +} + +// TestRequestLatencyUnitIsMillis confirms the request-latency histogram is +// registered in milliseconds so recorded values land across the default +// millisecond-scale buckets rather than collapsing into the first one. +func TestRequestLatencyUnitIsMillis(t *testing.T) { + reader := sdkmetric.NewManualReader() + mp := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader)) + m, err := output.NewMetrics(mp) + require.NoError(t, err) + + m.BlitzOutputRequestLatencyHistogram.Record(context.Background(), 1, "test", "logs") + + var rm metricdata.ResourceMetrics + require.NoError(t, reader.Collect(context.Background(), &rm)) + unit := "" + for _, sm := range rm.ScopeMetrics { + for _, mm := range sm.Metrics { + if mm.Name == "blitz.output.request_latency" { + unit = mm.Unit + } + } + } + require.Equal(t, "ms", unit, "request_latency unit") +} diff --git a/output/hec/ack.go b/output/hec/ack.go index 8993524..86e5d4c 100644 --- a/output/hec/ack.go +++ b/output/hec/ack.go @@ -208,7 +208,7 @@ func (p *ackPoller) poll() { // Query ACK status startTime := time.Now() confirmed, err := p.queryACK(ids) - p.metrics.recordACKPollLatency(ctx, time.Since(startTime).Seconds()) + p.metrics.recordACKPollLatency(ctx, output.DurationMillis(time.Since(startTime))) if err != nil { span.RecordError(err) diff --git a/output/hec/ackpoll_unit_test.go b/output/hec/ackpoll_unit_test.go new file mode 100644 index 0000000..19f0a2d --- /dev/null +++ b/output/hec/ackpoll_unit_test.go @@ -0,0 +1,34 @@ +package hec + +import ( + "context" + "testing" + + "github.com/stretchr/testify/require" + sdkmetric "go.opentelemetry.io/otel/sdk/metric" + "go.opentelemetry.io/otel/sdk/metric/metricdata" +) + +// TestACKPollLatencyUnitIsMillis confirms the HEC ACK-poll latency histogram is +// registered in milliseconds so recorded values spread across the default +// millisecond-scale buckets rather than collapsing into the first one. +func TestACKPollLatencyUnitIsMillis(t *testing.T) { + reader := sdkmetric.NewManualReader() + mp := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader)) + m, err := NewMetrics(mp) + require.NoError(t, err) + + m.blitzOutputHecAckPollLatencyHistogram.Record(context.Background(), 1) + + var rm metricdata.ResourceMetrics + require.NoError(t, reader.Collect(context.Background(), &rm)) + unit := "" + for _, sm := range rm.ScopeMetrics { + for _, mm := range sm.Metrics { + if mm.Name == "blitz.output.hec.ack_poll_latency" { + unit = mm.Unit + } + } + } + require.Equal(t, "ms", unit, "ack_poll_latency unit") +} diff --git a/output/hec/hec.go b/output/hec/hec.go index cbed7e1..84b7c3e 100644 --- a/output/hec/hec.go +++ b/output/hec/hec.go @@ -374,7 +374,7 @@ func (w *worker) sendBatch(batch []output.LogRecord) { startTime := time.Now() resp, err := w.postEvents(payload) - latency := time.Since(startTime).Seconds() + latency := output.DurationMillis(time.Since(startTime)) if err != nil { span.RecordError(err) diff --git a/output/hec/metrics.go b/output/hec/metrics.go index 54505d1..b8b5286 100644 --- a/output/hec/metrics.go +++ b/output/hec/metrics.go @@ -47,8 +47,8 @@ func (m *hecMetrics) recordRequestSize(ctx context.Context, bytes int64) { m.out.BlitzOutputRequestSizeHistogram.Record(ctx, bytes, outputType, "logs") } -func (m *hecMetrics) recordRequestLatency(ctx context.Context, seconds float64) { - m.out.BlitzOutputRequestLatencyHistogram.Record(ctx, seconds, outputType, "logs") +func (m *hecMetrics) recordRequestLatency(ctx context.Context, millis float64) { + m.out.BlitzOutputRequestLatencyHistogram.Record(ctx, millis, outputType, "logs") } func (m *hecMetrics) recordSendError(ctx context.Context, _ string) { @@ -79,6 +79,6 @@ func (m *hecMetrics) recordACKDropped(ctx context.Context, count int64) { m.hec.blitzOutputHecAckDroppedCounter.Add(ctx, count) } -func (m *hecMetrics) recordACKPollLatency(ctx context.Context, seconds float64) { - m.hec.blitzOutputHecAckPollLatencyHistogram.Record(ctx, seconds) +func (m *hecMetrics) recordACKPollLatency(ctx context.Context, millis float64) { + m.hec.blitzOutputHecAckPollLatencyHistogram.Record(ctx, millis) } diff --git a/output/hec/monitoring.go b/output/hec/monitoring.go index dc55350..d0d2b5d 100644 --- a/output/hec/monitoring.go +++ b/output/hec/monitoring.go @@ -83,7 +83,7 @@ func NewMetrics(mp metric.MeterProvider) (*Metrics, error) { blitzOutputHecAckPollLatencyHistogramRaw, err := m.hecMeter.Float64Histogram( "blitz.output.hec.ack_poll_latency", metric.WithDescription("latency of ACK polling requests"), - metric.WithUnit("s"), + metric.WithUnit("ms"), ) errs = errors.Join(errs, err) m.blitzOutputHecAckPollLatencyHistogram = blitzOutputHecAckPollLatencyHistogramRaw diff --git a/output/hec/monitoring.md b/output/hec/monitoring.md index eb6060a..a602c20 100644 --- a/output/hec/monitoring.md +++ b/output/hec/monitoring.md @@ -8,7 +8,7 @@ | [`blitz.output.hec.ack_dropped`](#blitzoutputhecack-dropped) | Counter | `{ack}` | total number of batches dropped after max retries | | [`blitz.output.hec.ack_expired`](#blitzoutputhecack-expired) | Counter | `{ack}` | total number of ACKs that expired without confirmation | | [`blitz.output.hec.ack_pending`](#blitzoutputhecack-pending) | Gauge | `{ack}` | number of ACKs currently pending confirmation | -| [`blitz.output.hec.ack_poll_latency`](#blitzoutputhecack-poll-latency) | Histogram | `s` | latency of ACK polling requests | +| [`blitz.output.hec.ack_poll_latency`](#blitzoutputhecack-poll-latency) | Histogram | `ms` | latency of ACK polling requests | | [`blitz.output.hec.ack_retried`](#blitzoutputhecack-retried) | Counter | `{ack}` | total number of batches retried due to ACK failure | | [`blitz.output.hec.batch_size`](#blitzoutputhecbatch-size) | Histogram | `{entry}` | number of entries per HEC batch | @@ -89,7 +89,7 @@ blitzOutputHecAckPendingGauge.Add(ctx, 1) | Property | Value | |----------|-------| | **Type** | Histogram | -| **Unit** | `s` | +| **Unit** | `ms` | | **Meter** | `hec` | | **Stability** | Stable | | **Description** | latency of ACK polling requests | diff --git a/output/hec/monitoring/metric.yaml b/output/hec/monitoring/metric.yaml index faf138f..f6ea69c 100644 --- a/output/hec/monitoring/metric.yaml +++ b/output/hec/monitoring/metric.yaml @@ -54,7 +54,7 @@ groups: stability: stable brief: "latency of ACK polling requests" instrument: histogram - unit: "s" + unit: "ms" annotations: histogram: type: Float64Histogram diff --git a/output/monitoring.go b/output/monitoring.go index 712d11c..1b741a3 100644 --- a/output/monitoring.go +++ b/output/monitoring.go @@ -205,7 +205,7 @@ func NewMetrics(mp metric.MeterProvider) (*Metrics, error) { BlitzOutputRequestLatencyHistogramRaw, err := m.outputMeter.Float64Histogram( "blitz.output.request_latency", metric.WithDescription("latency of output requests"), - metric.WithUnit("s"), + metric.WithUnit("ms"), ) errs = errors.Join(errs, err) m.BlitzOutputRequestLatencyHistogram = blitzOutputRequestLatencyHistogramType{histogram: BlitzOutputRequestLatencyHistogramRaw} diff --git a/output/monitoring.md b/output/monitoring.md index ccf740b..6c0360f 100644 --- a/output/monitoring.md +++ b/output/monitoring.md @@ -8,7 +8,7 @@ | [`blitz.output.entries_received`](#blitzoutputentries-received) | Counter | `{entry}` | total number of telemetry entries received by the output | | [`blitz.output.entry_rate`](#blitzoutputentry-rate) | Counter | `{entry}/s` | rate of telemetry entries processed per second | | [`blitz.output.queue_size`](#blitzoutputqueue-size) | Gauge | `{entry}` | current number of entries in the output queue | -| [`blitz.output.request_latency`](#blitzoutputrequest-latency) | Histogram | `s` | latency of output requests | +| [`blitz.output.request_latency`](#blitzoutputrequest-latency) | Histogram | `ms` | latency of output requests | | [`blitz.output.request_size`](#blitzoutputrequest-size) | Histogram | `By` | size of output requests in bytes | | [`blitz.output.send_errors`](#blitzoutputsend-errors) | Counter | `{error}` | total number of send errors | @@ -117,7 +117,7 @@ blitzOutputEntryRateCounter.Add(ctx, 1, outputTypeValue, telemetryTypeValue) | Property | Value | |----------|-------| | **Type** | Histogram | -| **Unit** | `s` | +| **Unit** | `ms` | | **Meter** | `output` | | **Stability** | Stable | | **Description** | latency of output requests | diff --git a/output/monitoring/metric.yaml b/output/monitoring/metric.yaml index 7701c06..c0f754f 100644 --- a/output/monitoring/metric.yaml +++ b/output/monitoring/metric.yaml @@ -94,7 +94,7 @@ groups: stability: stable brief: "latency of output requests" instrument: histogram - unit: "s" + unit: "ms" annotations: exported: true histogram: diff --git a/output/otlp_grpc/otlp_grpc.go b/output/otlp_grpc/otlp_grpc.go index 8f277c2..4562def 100644 --- a/output/otlp_grpc/otlp_grpc.go +++ b/output/otlp_grpc/otlp_grpc.go @@ -559,7 +559,7 @@ func (o *OTLPGrpc) sendMetricBatch(client collectormetrics.MetricsServiceClient, return fmt.Errorf("failed to export metrics: %w", err) } - latency := time.Since(startTime).Seconds() + latency := output.DurationMillis(time.Since(startTime)) requestSize := int64(proto.Size(request)) o.metrics.BlitzOutputEntryRateCounter.Add(context.Background(), float64(len(metrics)), outputType, "metrics") o.metrics.BlitzOutputRequestSizeHistogram.Record(context.Background(), requestSize, outputType, "metrics") @@ -594,7 +594,7 @@ func (o *OTLPGrpc) sendTraceBatch(client collectortrace.TraceServiceClient, batc return fmt.Errorf("failed to export traces: %w", err) } - latency := time.Since(startTime).Seconds() + latency := output.DurationMillis(time.Since(startTime)) requestSize := int64(proto.Size(request)) o.metrics.BlitzOutputEntryRateCounter.Add(context.Background(), float64(len(spans)), outputType, "traces") o.metrics.BlitzOutputRequestSizeHistogram.Record(context.Background(), requestSize, outputType, "traces") @@ -708,7 +708,7 @@ func (o *OTLPGrpc) sendBatch(client collectorlogs.LogsServiceClient, batch *logB } // Record successful send metrics - latency := time.Since(startTime).Seconds() + latency := output.DurationMillis(time.Since(startTime)) requestSize := int64(proto.Size(request)) o.metrics.BlitzOutputEntryRateCounter.Add(context.Background(), float64(len(logs)), outputType, "logs") o.metrics.BlitzOutputRequestSizeHistogram.Record(context.Background(), requestSize, outputType, "logs") diff --git a/output/tcp/tcp.go b/output/tcp/tcp.go index 28e0153..3cfcc11 100644 --- a/output/tcp/tcp.go +++ b/output/tcp/tcp.go @@ -263,7 +263,7 @@ func (t *TCP) sendData(conn net.Conn, data string) error { } // Record successful send metrics - latency := time.Since(startTime).Seconds() + latency := output.DurationMillis(time.Since(startTime)) t.metrics.BlitzOutputEntryRateCounter.Add(context.Background(), 1.0, outputType, "logs") t.metrics.BlitzOutputRequestSizeHistogram.Record(context.Background(), int64(bytesWritten), outputType, "logs") t.metrics.BlitzOutputRequestLatencyHistogram.Record(context.Background(), latency, outputType, "logs") diff --git a/output/util.go b/output/util.go index 149c3ab..a06794b 100644 --- a/output/util.go +++ b/output/util.go @@ -19,3 +19,12 @@ func Int64ToUint64(nanos int64) uint64 { func TimeToUnixNanoUint64(t time.Time) uint64 { return Int64ToUint64(t.UnixNano()) } + +// DurationMillis converts a duration to fractional milliseconds for recording +// on latency histograms. It divides the nanosecond count as a float so a +// sub-millisecond value keeps its fractional part. time.Duration's Milliseconds() +// truncates to an integer, which would drop sub-millisecond samples to zero and +// understate the histogram sum. +func DurationMillis(d time.Duration) float64 { + return float64(d.Nanoseconds()) / 1e6 +}