2020-03-24 07:41:10 +02:00
|
|
|
// Copyright The OpenTelemetry Authors
|
2019-10-29 22:27:22 +02:00
|
|
|
//
|
|
|
|
// 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
|
|
|
|
//
|
|
|
|
// http://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.
|
|
|
|
|
2020-11-04 19:10:58 +02:00
|
|
|
package aggregatortest // import "go.opentelemetry.io/otel/sdk/metric/aggregator/aggregatortest"
|
2019-10-29 22:27:22 +02:00
|
|
|
|
|
|
|
import (
|
|
|
|
"context"
|
2020-12-11 04:13:08 +02:00
|
|
|
"errors"
|
2019-10-29 22:27:22 +02:00
|
|
|
"math/rand"
|
2020-01-06 20:08:40 +02:00
|
|
|
"os"
|
2019-10-29 22:27:22 +02:00
|
|
|
"sort"
|
|
|
|
"testing"
|
2020-01-06 20:08:40 +02:00
|
|
|
"unsafe"
|
2019-10-29 22:27:22 +02:00
|
|
|
|
2020-12-11 04:13:08 +02:00
|
|
|
"github.com/stretchr/testify/require"
|
|
|
|
|
2020-01-06 20:08:40 +02:00
|
|
|
ottest "go.opentelemetry.io/otel/internal/testing"
|
2020-11-12 17:28:32 +02:00
|
|
|
"go.opentelemetry.io/otel/metric"
|
2020-11-11 17:24:12 +02:00
|
|
|
"go.opentelemetry.io/otel/metric/number"
|
2019-11-05 23:08:55 +02:00
|
|
|
export "go.opentelemetry.io/otel/sdk/export/metric"
|
2020-12-11 04:13:08 +02:00
|
|
|
"go.opentelemetry.io/otel/sdk/export/metric/aggregation"
|
2020-06-10 07:53:30 +02:00
|
|
|
"go.opentelemetry.io/otel/sdk/metric/aggregator"
|
2019-10-29 22:27:22 +02:00
|
|
|
)
|
|
|
|
|
2019-11-05 00:24:01 +02:00
|
|
|
const Magnitude = 1000
|
|
|
|
|
2019-10-29 22:27:22 +02:00
|
|
|
type Profile struct {
|
2020-11-11 17:24:12 +02:00
|
|
|
NumberKind number.Kind
|
|
|
|
Random func(sign int) number.Number
|
2019-10-29 22:27:22 +02:00
|
|
|
}
|
|
|
|
|
2020-12-11 04:13:08 +02:00
|
|
|
type NoopAggregator struct{}
|
|
|
|
type NoopAggregation struct{}
|
|
|
|
|
|
|
|
var _ export.Aggregator = NoopAggregator{}
|
|
|
|
var _ aggregation.Aggregation = NoopAggregation{}
|
|
|
|
|
2019-11-05 00:24:01 +02:00
|
|
|
func newProfiles() []Profile {
|
|
|
|
rnd := rand.New(rand.NewSource(rand.Int63()))
|
|
|
|
return []Profile{
|
|
|
|
{
|
2020-11-11 17:24:12 +02:00
|
|
|
NumberKind: number.Int64Kind,
|
|
|
|
Random: func(sign int) number.Number {
|
|
|
|
return number.NewInt64Number(int64(sign) * int64(rnd.Intn(Magnitude+1)))
|
2019-11-05 00:24:01 +02:00
|
|
|
},
|
2019-10-29 22:27:22 +02:00
|
|
|
},
|
2019-11-05 00:24:01 +02:00
|
|
|
{
|
2020-11-11 17:24:12 +02:00
|
|
|
NumberKind: number.Float64Kind,
|
|
|
|
Random: func(sign int) number.Number {
|
|
|
|
return number.NewFloat64Number(float64(sign) * rnd.Float64() * Magnitude)
|
2019-11-05 00:24:01 +02:00
|
|
|
},
|
2019-10-29 22:27:22 +02:00
|
|
|
},
|
2019-11-05 00:24:01 +02:00
|
|
|
}
|
2019-10-29 22:27:22 +02:00
|
|
|
}
|
|
|
|
|
2020-11-12 17:28:32 +02:00
|
|
|
func NewAggregatorTest(mkind metric.InstrumentKind, nkind number.Kind) *metric.Descriptor {
|
|
|
|
desc := metric.NewDescriptor("test.name", mkind, nkind)
|
2020-03-19 21:02:46 +02:00
|
|
|
return &desc
|
2019-10-29 22:27:22 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
func RunProfiles(t *testing.T, f func(*testing.T, Profile)) {
|
2019-11-05 00:24:01 +02:00
|
|
|
for _, profile := range newProfiles() {
|
2019-10-29 22:27:22 +02:00
|
|
|
t.Run(profile.NumberKind.String(), func(t *testing.T) {
|
|
|
|
f(t, profile)
|
|
|
|
})
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2020-01-06 20:08:40 +02:00
|
|
|
// Ensure local struct alignment prior to running tests.
|
|
|
|
func TestMain(m *testing.M) {
|
|
|
|
fields := []ottest.FieldOffset{
|
|
|
|
{
|
|
|
|
Name: "Numbers.numbers",
|
|
|
|
Offset: unsafe.Offsetof(Numbers{}.numbers),
|
|
|
|
},
|
|
|
|
}
|
|
|
|
if !ottest.Aligned8Byte(fields, os.Stderr) {
|
|
|
|
os.Exit(1)
|
|
|
|
}
|
|
|
|
|
|
|
|
os.Exit(m.Run())
|
|
|
|
}
|
|
|
|
|
2020-05-21 20:09:10 +02:00
|
|
|
// TODO: Expose Numbers in api/metric for sorting support
|
|
|
|
|
2019-11-05 00:24:01 +02:00
|
|
|
type Numbers struct {
|
2020-01-06 20:08:40 +02:00
|
|
|
// numbers has to be aligned for 64-bit atomic operations.
|
2020-11-11 17:24:12 +02:00
|
|
|
numbers []number.Number
|
|
|
|
kind number.Kind
|
2019-11-05 00:24:01 +02:00
|
|
|
}
|
|
|
|
|
2020-11-11 17:24:12 +02:00
|
|
|
func NewNumbers(kind number.Kind) Numbers {
|
2019-11-05 00:24:01 +02:00
|
|
|
return Numbers{
|
|
|
|
kind: kind,
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2020-11-11 17:24:12 +02:00
|
|
|
func (n *Numbers) Append(v number.Number) {
|
2019-11-05 00:24:01 +02:00
|
|
|
n.numbers = append(n.numbers, v)
|
|
|
|
}
|
2019-10-29 22:27:22 +02:00
|
|
|
|
|
|
|
func (n *Numbers) Sort() {
|
2019-10-31 07:15:27 +02:00
|
|
|
sort.Sort(n)
|
|
|
|
}
|
|
|
|
|
|
|
|
func (n *Numbers) Less(i, j int) bool {
|
2019-11-05 00:24:01 +02:00
|
|
|
return n.numbers[i].CompareNumber(n.kind, n.numbers[j]) < 0
|
2019-10-31 07:15:27 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
func (n *Numbers) Len() int {
|
2019-11-05 00:24:01 +02:00
|
|
|
return len(n.numbers)
|
2019-10-31 07:15:27 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
func (n *Numbers) Swap(i, j int) {
|
2019-11-05 00:24:01 +02:00
|
|
|
n.numbers[i], n.numbers[j] = n.numbers[j], n.numbers[i]
|
2019-10-29 22:27:22 +02:00
|
|
|
}
|
|
|
|
|
2020-11-11 17:24:12 +02:00
|
|
|
func (n *Numbers) Sum() number.Number {
|
|
|
|
var sum number.Number
|
2019-11-05 00:24:01 +02:00
|
|
|
for _, num := range n.numbers {
|
|
|
|
sum.AddNumber(n.kind, num)
|
2019-10-29 22:27:22 +02:00
|
|
|
}
|
|
|
|
return sum
|
|
|
|
}
|
|
|
|
|
2021-01-06 09:17:20 +02:00
|
|
|
func (n *Numbers) Count() uint64 {
|
|
|
|
return uint64(len(n.numbers))
|
2019-10-29 22:27:22 +02:00
|
|
|
}
|
2019-10-31 07:15:27 +02:00
|
|
|
|
2020-11-11 17:24:12 +02:00
|
|
|
func (n *Numbers) Min() number.Number {
|
2019-11-05 00:24:01 +02:00
|
|
|
return n.numbers[0]
|
|
|
|
}
|
2019-10-31 07:15:27 +02:00
|
|
|
|
2020-11-11 17:24:12 +02:00
|
|
|
func (n *Numbers) Max() number.Number {
|
2019-11-05 00:24:01 +02:00
|
|
|
return n.numbers[len(n.numbers)-1]
|
|
|
|
}
|
2019-10-31 07:15:27 +02:00
|
|
|
|
2019-11-05 00:24:01 +02:00
|
|
|
// Median() is an alias for Quantile(0.5).
|
2020-11-11 17:24:12 +02:00
|
|
|
func (n *Numbers) Median() number.Number {
|
2019-11-05 00:24:01 +02:00
|
|
|
// Note that len(n.numbers) is 1 greater than the max element
|
|
|
|
// index, so dividing by two rounds up. This gives the
|
|
|
|
// intended definition for Quantile() in tests, which is to
|
|
|
|
// return the smallest element that is at or above the
|
|
|
|
// specified quantile.
|
|
|
|
return n.numbers[len(n.numbers)/2]
|
2019-10-31 07:15:27 +02:00
|
|
|
}
|
2019-11-15 23:01:20 +02:00
|
|
|
|
2020-11-11 17:24:12 +02:00
|
|
|
func (n *Numbers) Points() []number.Number {
|
2019-11-26 21:47:15 +02:00
|
|
|
return n.numbers
|
|
|
|
}
|
|
|
|
|
2019-11-15 23:01:20 +02:00
|
|
|
// Performs the same range test the SDK does on behalf of the aggregator.
|
2020-11-12 17:28:32 +02:00
|
|
|
func CheckedUpdate(t *testing.T, agg export.Aggregator, number number.Number, descriptor *metric.Descriptor) {
|
2019-11-15 23:01:20 +02:00
|
|
|
ctx := context.Background()
|
|
|
|
|
|
|
|
// Note: Aggregator tests are written assuming that the SDK
|
|
|
|
// has performed the RangeTest. Therefore we skip errors that
|
|
|
|
// would have been detected by the RangeTest.
|
|
|
|
err := aggregator.RangeTest(number, descriptor)
|
|
|
|
if err != nil {
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
|
|
|
if err := agg.Update(ctx, number, descriptor); err != nil {
|
|
|
|
t.Error("Unexpected Update failure", err)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2020-11-12 17:28:32 +02:00
|
|
|
func CheckedMerge(t *testing.T, aggInto, aggFrom export.Aggregator, descriptor *metric.Descriptor) {
|
2019-11-15 23:01:20 +02:00
|
|
|
if err := aggInto.Merge(aggFrom, descriptor); err != nil {
|
|
|
|
t.Error("Unexpected Merge failure", err)
|
|
|
|
}
|
|
|
|
}
|
2020-12-11 04:13:08 +02:00
|
|
|
|
|
|
|
func (NoopAggregation) Kind() aggregation.Kind {
|
|
|
|
return aggregation.Kind("Noop")
|
|
|
|
}
|
|
|
|
|
|
|
|
func (NoopAggregator) Aggregation() aggregation.Aggregation {
|
|
|
|
return NoopAggregation{}
|
|
|
|
}
|
|
|
|
|
|
|
|
func (NoopAggregator) Update(context.Context, number.Number, *metric.Descriptor) error {
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func (NoopAggregator) SynchronizedMove(export.Aggregator, *metric.Descriptor) error {
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func (NoopAggregator) Merge(export.Aggregator, *metric.Descriptor) error {
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func SynchronizedMoveResetTest(t *testing.T, mkind metric.InstrumentKind, nf func(*metric.Descriptor) export.Aggregator) {
|
|
|
|
t.Run("reset on nil", func(t *testing.T) {
|
|
|
|
// Ensures that SynchronizedMove(nil, descriptor) discards and
|
|
|
|
// resets the aggregator.
|
|
|
|
RunProfiles(t, func(t *testing.T, profile Profile) {
|
|
|
|
descriptor := NewAggregatorTest(
|
|
|
|
mkind,
|
|
|
|
profile.NumberKind,
|
|
|
|
)
|
|
|
|
agg := nf(descriptor)
|
|
|
|
|
|
|
|
for i := 0; i < 10; i++ {
|
|
|
|
x1 := profile.Random(+1)
|
|
|
|
CheckedUpdate(t, agg, x1, descriptor)
|
|
|
|
}
|
|
|
|
|
|
|
|
require.NoError(t, agg.SynchronizedMove(nil, descriptor))
|
|
|
|
|
|
|
|
if count, ok := agg.(aggregation.Count); ok {
|
|
|
|
c, err := count.Count()
|
2021-01-06 09:17:20 +02:00
|
|
|
require.Equal(t, uint64(0), c)
|
2020-12-11 04:13:08 +02:00
|
|
|
require.NoError(t, err)
|
|
|
|
}
|
|
|
|
|
|
|
|
if sum, ok := agg.(aggregation.Sum); ok {
|
|
|
|
s, err := sum.Sum()
|
|
|
|
require.Equal(t, number.Number(0), s)
|
|
|
|
require.NoError(t, err)
|
|
|
|
}
|
|
|
|
|
|
|
|
if lv, ok := agg.(aggregation.LastValue); ok {
|
|
|
|
v, _, err := lv.LastValue()
|
|
|
|
require.Equal(t, number.Number(0), v)
|
|
|
|
require.Error(t, err)
|
|
|
|
require.True(t, errors.Is(err, aggregation.ErrNoData))
|
|
|
|
}
|
|
|
|
})
|
|
|
|
})
|
|
|
|
|
|
|
|
t.Run("no reset on incorrect type", func(t *testing.T) {
|
|
|
|
// Ensures that SynchronizedMove(wrong_type, descriptor) does not
|
|
|
|
// reset the aggregator.
|
|
|
|
RunProfiles(t, func(t *testing.T, profile Profile) {
|
|
|
|
descriptor := NewAggregatorTest(
|
|
|
|
mkind,
|
|
|
|
profile.NumberKind,
|
|
|
|
)
|
|
|
|
agg := nf(descriptor)
|
|
|
|
|
|
|
|
var input number.Number
|
|
|
|
const inval = 100
|
|
|
|
if profile.NumberKind == number.Int64Kind {
|
|
|
|
input = number.NewInt64Number(inval)
|
|
|
|
} else {
|
|
|
|
input = number.NewFloat64Number(inval)
|
|
|
|
}
|
|
|
|
|
|
|
|
CheckedUpdate(t, agg, input, descriptor)
|
|
|
|
|
|
|
|
err := agg.SynchronizedMove(NoopAggregator{}, descriptor)
|
|
|
|
require.Error(t, err)
|
|
|
|
require.True(t, errors.Is(err, aggregation.ErrInconsistentType))
|
|
|
|
|
|
|
|
// Test that the aggregator was not reset
|
|
|
|
|
|
|
|
if count, ok := agg.(aggregation.Count); ok {
|
|
|
|
c, err := count.Count()
|
2021-01-06 09:17:20 +02:00
|
|
|
require.Equal(t, uint64(1), c)
|
2020-12-11 04:13:08 +02:00
|
|
|
require.NoError(t, err)
|
|
|
|
}
|
|
|
|
|
|
|
|
if sum, ok := agg.(aggregation.Sum); ok {
|
|
|
|
s, err := sum.Sum()
|
|
|
|
require.Equal(t, input, s)
|
|
|
|
require.NoError(t, err)
|
|
|
|
}
|
|
|
|
|
|
|
|
if lv, ok := agg.(aggregation.LastValue); ok {
|
|
|
|
v, _, err := lv.LastValue()
|
|
|
|
require.Equal(t, input, v)
|
|
|
|
require.NoError(t, err)
|
|
|
|
}
|
|
|
|
|
|
|
|
})
|
|
|
|
})
|
|
|
|
|
|
|
|
}
|