diff --git a/tests/v2/e2e/config/config.go b/tests/v2/e2e/config/config.go index d1e2ba28b4..e95c0d27b4 100644 --- a/tests/v2/e2e/config/config.go +++ b/tests/v2/e2e/config/config.go @@ -58,6 +58,7 @@ type Data struct { Dataset *Dataset `json:"dataset,omitempty" yaml:"dataset,omitempty"` Kubernetes *Kubernetes `json:"kubernetes,omitempty" yaml:"kubernetes,omitempty"` Metrics *Metrics `json:"metrics,omitempty" yaml:"metrics,omitempty"` + Observability *config.Observability `json:"observability,omitempty" yaml:"observability,omitempty"` Metadata map[string]string `json:"metadata,omitempty" yaml:"metadata,omitempty"` MetaString string `json:"metadata_string,omitempty" yaml:"metadata_string,omitempty"` FilePath string `json:"-" yaml:"-"` @@ -276,6 +277,10 @@ func (d *Data) Bind() (bound *Data, err error) { } } } + // Bind Observability. + if d.Observability != nil { + d.Observability.Bind() + } // Bind Dataset. if d.Dataset != nil { if ds, err := d.Dataset.Bind(); err != nil { diff --git a/tests/v2/e2e/crud/strategy_test.go b/tests/v2/e2e/crud/strategy_test.go index 1fe91fd703..6e5beb2e0d 100644 --- a/tests/v2/e2e/crud/strategy_test.go +++ b/tests/v2/e2e/crud/strategy_test.go @@ -30,6 +30,8 @@ import ( "github.com/vdaas/vald/internal/errors" "github.com/vdaas/vald/internal/log" "github.com/vdaas/vald/internal/net/grpc" + "github.com/vdaas/vald/internal/observability" + obsmetrics "github.com/vdaas/vald/internal/observability/metrics" "github.com/vdaas/vald/internal/sync/errgroup" "github.com/vdaas/vald/tests/v2/e2e/config" k8s "github.com/vdaas/vald/tests/v2/e2e/kubernetes" @@ -53,6 +55,28 @@ func TestE2EStrategy(t *testing.T) { ctx, cancel := context.WithCancel(t.Context()) defer cancel() + if cfg.Observability != nil && cfg.Observability.Enabled { + var extraMetrics []obsmetrics.Metric + if cfg.Collector != nil { + extraMetrics = append(extraMetrics, metrics.NewOTELMetrics(cfg.Collector)) + } else { + t.Log("observability enabled but collector is nil, skipping otel metrics registration") + } + + obs, err := observability.NewWithConfig(cfg.Observability, extraMetrics...) + if err != nil { + t.Fatalf("failed to create observability: %v", err) + } + if err := obs.PreStart(ctx); err != nil { + t.Fatalf("failed to start observability: %v", err) + } + defer func() { + if err := obs.Stop(ctx); err != nil { + t.Logf("failed to stop observability: %v", err) + } + }() + } + var err error r := new(runner) if cfg.Kubernetes != nil { @@ -367,6 +391,7 @@ func executeWithTimings[T interface { var cancel context.CancelFunc ctx, cancel = context.WithTimeout(ctx, dur) defer cancel() + } } diff --git a/tests/v2/e2e/metrics/otel.go b/tests/v2/e2e/metrics/otel.go new file mode 100644 index 0000000000..7eb0a4ea41 --- /dev/null +++ b/tests/v2/e2e/metrics/otel.go @@ -0,0 +1,185 @@ +// +// Copyright (C) 2019-2026 vdaas.org vald team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// You may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// https://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. +// + +package metrics + +import ( + "context" + "math" + "time" + + "go.opentelemetry.io/otel/metric" + + obsmetrics "github.com/vdaas/vald/internal/observability/metrics" +) + +type otelMetrics struct { + collector Collector +} + +func NewOTELMetrics(c Collector) obsmetrics.Metric { + return &otelMetrics{ + collector: c, + } +} + +func (o *otelMetrics) View() ([]obsmetrics.View, error) { + return nil, nil +} + +func (o *otelMetrics) Register(m obsmetrics.Meter) error { + totalRequests, err := m.Int64ObservableGauge( + "e2e_total_requests", + obsmetrics.WithDescription("Total number of requests"), + obsmetrics.WithUnit(obsmetrics.Dimensionless), + ) + if err != nil { + return err + } + + totalErrors, err := m.Int64ObservableGauge( + "e2e_total_errors", + obsmetrics.WithDescription("Total number of errors"), + obsmetrics.WithUnit(obsmetrics.Dimensionless), + ) + if err != nil { + return err + } + + // Latency metrics + latP50, err := m.Float64ObservableGauge( + "e2e_latency_p50", + obsmetrics.WithDescription("Latency P50"), + obsmetrics.WithUnit(obsmetrics.Milliseconds), + ) + if err != nil { + return err + } + + latP90, err := m.Float64ObservableGauge( + "e2e_latency_p90", + obsmetrics.WithDescription("Latency P90"), + obsmetrics.WithUnit(obsmetrics.Milliseconds), + ) + if err != nil { + return err + } + + latP99, err := m.Float64ObservableGauge( + "e2e_latency_p99", + obsmetrics.WithDescription("Latency P99"), + obsmetrics.WithUnit(obsmetrics.Milliseconds), + ) + if err != nil { + return err + } + + latMax, err := m.Float64ObservableGauge( + "e2e_latency_max", + obsmetrics.WithDescription("Latency Max"), + obsmetrics.WithUnit(obsmetrics.Milliseconds), + ) + if err != nil { + return err + } + + latMin, err := m.Float64ObservableGauge( + "e2e_latency_min", + obsmetrics.WithDescription("Latency Min"), + obsmetrics.WithUnit(obsmetrics.Milliseconds), + ) + if err != nil { + return err + } + + // Queue Wait metrics + qwP50, err := m.Float64ObservableGauge( + "e2e_queue_wait_p50", + obsmetrics.WithDescription("Queue Wait P50"), + obsmetrics.WithUnit(obsmetrics.Milliseconds), + ) + if err != nil { + return err + } + + qwP90, err := m.Float64ObservableGauge( + "e2e_queue_wait_p90", + obsmetrics.WithDescription("Queue Wait P90"), + obsmetrics.WithUnit(obsmetrics.Milliseconds), + ) + if err != nil { + return err + } + + qwP99, err := m.Float64ObservableGauge( + "e2e_queue_wait_p99", + obsmetrics.WithDescription("Queue Wait P99"), + obsmetrics.WithUnit(obsmetrics.Milliseconds), + ) + if err != nil { + return err + } + + qwMax, err := m.Float64ObservableGauge( + "e2e_queue_wait_max", + obsmetrics.WithDescription("Queue Wait Max"), + obsmetrics.WithUnit(obsmetrics.Milliseconds), + ) + if err != nil { + return err + } + + _, err = m.RegisterCallback(func(_ context.Context, obs metric.Observer) error { + snap := o.collector.GlobalSnapshot() + if snap == nil { + return nil + } + + safeInt64 := func(v uint64) int64 { + if v > math.MaxInt64 { + return math.MaxInt64 + } + return int64(v) + } + + obs.ObserveInt64(totalRequests, safeInt64(snap.Total)) + obs.ObserveInt64(totalErrors, safeInt64(snap.Errors)) + + // Latency + if snap.LatPercentiles != nil { + obs.ObserveFloat64(latP50, nsToMs(snap.LatPercentiles.Quantile(0.5))) + obs.ObserveFloat64(latP90, nsToMs(snap.LatPercentiles.Quantile(0.9))) + obs.ObserveFloat64(latP99, nsToMs(snap.LatPercentiles.Quantile(0.99))) + obs.ObserveFloat64(latMax, nsToMs(snap.LatPercentiles.Quantile(1.0))) + obs.ObserveFloat64(latMin, nsToMs(snap.LatPercentiles.Quantile(0.0))) + } + + // Queue Wait + if snap.QWPercentiles != nil { + obs.ObserveFloat64(qwP50, nsToMs(snap.QWPercentiles.Quantile(0.5))) + obs.ObserveFloat64(qwP90, nsToMs(snap.QWPercentiles.Quantile(0.9))) + obs.ObserveFloat64(qwP99, nsToMs(snap.QWPercentiles.Quantile(0.99))) + obs.ObserveFloat64(qwMax, nsToMs(snap.QWPercentiles.Quantile(1.0))) + } + return nil + }, totalRequests, totalErrors, latP50, latP90, latP99, latMax, latMin, qwP50, qwP90, qwP99, qwMax) + + return err +} + +func nsToMs(ns float64) float64 { + return ns / float64(time.Millisecond) +}