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
54 changes: 54 additions & 0 deletions output/duration_test.go
Original file line number Diff line number Diff line change
@@ -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")
}
2 changes: 1 addition & 1 deletion output/hec/ack.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
34 changes: 34 additions & 0 deletions output/hec/ackpoll_unit_test.go
Original file line number Diff line number Diff line change
@@ -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")
}
2 changes: 1 addition & 1 deletion output/hec/hec.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
8 changes: 4 additions & 4 deletions output/hec/metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down Expand Up @@ -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)
}
2 changes: 1 addition & 1 deletion output/hec/monitoring.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
4 changes: 2 additions & 2 deletions output/hec/monitoring.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |

Expand Down Expand Up @@ -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 |
Expand Down
2 changes: 1 addition & 1 deletion output/hec/monitoring/metric.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -54,7 +54,7 @@ groups:
stability: stable
brief: "latency of ACK polling requests"
instrument: histogram
unit: "s"
unit: "ms"
annotations:
histogram:
type: Float64Histogram
2 changes: 1 addition & 1 deletion output/monitoring.go
Original file line number Diff line number Diff line change
Expand Up @@ -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}
Expand Down
4 changes: 2 additions & 2 deletions output/monitoring.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |

Expand Down Expand Up @@ -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 |
Expand Down
2 changes: 1 addition & 1 deletion output/monitoring/metric.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -94,7 +94,7 @@ groups:
stability: stable
brief: "latency of output requests"
instrument: histogram
unit: "s"
unit: "ms"
annotations:
exported: true
histogram:
Expand Down
6 changes: 3 additions & 3 deletions output/otlp_grpc/otlp_grpc.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down Expand Up @@ -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")
Expand Down Expand Up @@ -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")
Expand Down
2 changes: 1 addition & 1 deletion output/tcp/tcp.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down
9 changes: 9 additions & 0 deletions output/util.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Loading