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
20 changes: 19 additions & 1 deletion libs/monitoring/common.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package monitoring
import (
"context"
"fmt"
"math"

"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/metric"
Expand All @@ -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.
Expand Down Expand Up @@ -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,
}
}
Expand All @@ -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
Expand Down Expand Up @@ -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
}

Expand All @@ -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)
Expand Down
62 changes: 62 additions & 0 deletions libs/monitoring/common_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package monitoring_test

import (
"context"
"math"
"testing"

"github.com/stretchr/testify/require"
Expand Down Expand Up @@ -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()

Expand Down Expand Up @@ -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
}
Loading