diff --git a/libs/monitoring/common.go b/libs/monitoring/common.go index dd11ec0e8..3b462ec43 100644 --- a/libs/monitoring/common.go +++ b/libs/monitoring/common.go @@ -3,6 +3,7 @@ package monitoring import ( "context" "fmt" + "math" "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/metric" @@ -15,7 +16,8 @@ type MetricsInfoCapBasic struct { capTimestampStart beholder.MetricInfo capTimestampEmit beholder.MetricInfo capDuration beholder.MetricInfo // ts.emit - ts.start - capDurationBuckets []float64 // explicit histogram buckets; nil uses the SDK default + invalidTelemetry beholder.MetricInfo + capDurationBuckets []float64 // explicit histogram buckets; nil uses the SDK default } // NewMetricsInfoCapBasic creates a MetricsInfoCapBasic with default histogram buckets. @@ -51,6 +53,11 @@ func newMetricsInfoCapBasic(metricPrefix, eventRef string, buckets []float64) Me Unit: "ms", Description: fmt.Sprintf("The duration (local) since capability exec start to message: '%s' emit", eventRef), }, + invalidTelemetry: beholder.MetricInfo{ + Name: fmt.Sprintf("%s_invalid_telemetry_count", metricPrefix), + Unit: "", + Description: fmt.Sprintf("The count of message: '%s' emitted with invalid telemetry timestamps", eventRef), + }, capDurationBuckets: buckets, } } @@ -61,6 +68,7 @@ type MetricsCapBasic struct { capTimestampStart metric.Int64Gauge capTimestampEmit metric.Int64Gauge capDuration metric.Int64Histogram // ts.emit - ts.start + invalidTelemetry metric.Int64Counter } // NewMetricsCapBasic creates a new MetricsCapBasic using the provided MetricsInfoCapBasic @@ -98,6 +106,11 @@ func NewMetricsCapBasic(info MetricsInfoCapBasic) (MetricsCapBasic, error) { return set, fmt.Errorf("failed to create new histogram: %w", err) } + set.invalidTelemetry, err = info.invalidTelemetry.NewInt64Counter(meter) + if err != nil { + return set, fmt.Errorf("failed to create new counter: %w", err) + } + return set, nil } @@ -108,6 +121,11 @@ func (m *MetricsCapBasic) RecordEmit(ctx context.Context, start, emit uint64, at // Count events m.count.Add(ctx, 1, attrs) + if start > math.MaxInt64 || emit > math.MaxInt64 || emit < start { + m.invalidTelemetry.Add(ctx, 1, attrs) + return + } + // Timestamp events m.capTimestampStart.Record(ctx, int64(start), attrs) m.capTimestampEmit.Record(ctx, int64(emit), attrs) diff --git a/libs/monitoring/common_test.go b/libs/monitoring/common_test.go index dbfcac2b6..1c8f08f53 100644 --- a/libs/monitoring/common_test.go +++ b/libs/monitoring/common_test.go @@ -2,6 +2,7 @@ package monitoring_test import ( "context" + "math" "testing" "github.com/stretchr/testify/require" @@ -48,6 +49,44 @@ func TestNewMetricsCapBasic_WithoutBuckets(t *testing.T) { require.NotNil(t, findCounter(resourceMetrics, "test_metric_default_count")) } +func TestMetricsCapBasic_RecordEmitSkipsReversedTimestamps(t *testing.T) { + reader := useManualMetricReader(t) + + info := capmonitoring.NewMetricsInfoCapBasic("test_metric_reversed", "test.event.reversed") + metrics, err := capmonitoring.NewMetricsCapBasic(info) + require.NoError(t, err) + + metrics.RecordEmit(t.Context(), 200, 100) + + var resourceMetrics metricdata.ResourceMetrics + require.NoError(t, reader.Collect(context.Background(), &resourceMetrics)) + + require.EqualValues(t, 1, counterValue(t, resourceMetrics, "test_metric_reversed_count")) + require.EqualValues(t, 1, counterValue(t, resourceMetrics, "test_metric_reversed_invalid_telemetry_count")) + require.Nil(t, findHistogram(resourceMetrics, "test_metric_reversed_cap_duration")) + require.Nil(t, findGauge(resourceMetrics, "test_metric_reversed_cap_timestamp_start")) + require.Nil(t, findGauge(resourceMetrics, "test_metric_reversed_cap_timestamp_emit")) +} + +func TestMetricsCapBasic_RecordEmitSkipsTimestampsThatOverflowInt64(t *testing.T) { + reader := useManualMetricReader(t) + + info := capmonitoring.NewMetricsInfoCapBasic("test_metric_overflow", "test.event.overflow") + metrics, err := capmonitoring.NewMetricsCapBasic(info) + require.NoError(t, err) + + metrics.RecordEmit(t.Context(), uint64(math.MaxInt64)+1, uint64(math.MaxInt64)+2) + + var resourceMetrics metricdata.ResourceMetrics + require.NoError(t, reader.Collect(context.Background(), &resourceMetrics)) + + require.EqualValues(t, 1, counterValue(t, resourceMetrics, "test_metric_overflow_count")) + require.EqualValues(t, 1, counterValue(t, resourceMetrics, "test_metric_overflow_invalid_telemetry_count")) + require.Nil(t, findHistogram(resourceMetrics, "test_metric_overflow_cap_duration")) + require.Nil(t, findGauge(resourceMetrics, "test_metric_overflow_cap_timestamp_start")) + require.Nil(t, findGauge(resourceMetrics, "test_metric_overflow_cap_timestamp_emit")) +} + func useManualMetricReader(t *testing.T) *sdkmetric.ManualReader { t.Helper() @@ -94,3 +133,26 @@ func findCounter(resourceMetrics metricdata.ResourceMetrics, name string) *metri } return nil } + +func findGauge(resourceMetrics metricdata.ResourceMetrics, name string) *metricdata.Gauge[int64] { + for _, scopeMetrics := range resourceMetrics.ScopeMetrics { + for _, metric := range scopeMetrics.Metrics { + if metric.Name != name { + continue + } + if gauge, ok := metric.Data.(metricdata.Gauge[int64]); ok { + return &gauge + } + } + } + return nil +} + +func counterValue(t *testing.T, resourceMetrics metricdata.ResourceMetrics, name string) int64 { + t.Helper() + + counter := findCounter(resourceMetrics, name) + require.NotNil(t, counter) + require.Len(t, counter.DataPoints, 1) + return counter.DataPoints[0].Value +}