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
5 changes: 5 additions & 0 deletions tests/v2/e2e/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -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:"-"`
Expand Down Expand Up @@ -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 {
Expand Down
25 changes: 25 additions & 0 deletions tests/v2/e2e/crud/strategy_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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 {
Expand Down Expand Up @@ -367,6 +391,7 @@ func executeWithTimings[T interface {
var cancel context.CancelFunc
ctx, cancel = context.WithTimeout(ctx, dur)
defer cancel()

}
}

Expand Down
185 changes: 185 additions & 0 deletions tests/v2/e2e/metrics/otel.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,185 @@
//
// Copyright (C) 2019-2026 vdaas.org vald team <vald@vdaas.org>
//
// 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)
}
Loading