diff --git a/metric/distribution/distribution.go b/metric/distribution/distribution.go index 06a28f67472..97b19fffb4b 100644 --- a/metric/distribution/distribution.go +++ b/metric/distribution/distribution.go @@ -19,10 +19,10 @@ type Distribution interface { Size() int // weight is 1/samplingRate - AddEntryWithUnit(value float64, weight float64, unit string) + AddEntryWithUnit(value float64, weight float64, unit string) error // weight is 1/samplingRate - AddEntry(value float64, weight float64) + AddEntry(value float64, weight float64) error AddDistribution(distribution Distribution) diff --git a/metric/distribution/regular/regular_distribution.go b/metric/distribution/regular/regular_distribution.go index aea9d99e2fc..e8ff62069d4 100644 --- a/metric/distribution/regular/regular_distribution.go +++ b/metric/distribution/regular/regular_distribution.go @@ -4,6 +4,7 @@ package regular import ( + "errors" "log" "math" @@ -65,11 +66,10 @@ func (regularDist *RegularDistribution) Size() int { } // weight is 1/samplingRate -func (regularDist *RegularDistribution) AddEntryWithUnit(value float64, weight float64, unit string) { +func (regularDist *RegularDistribution) AddEntryWithUnit(value float64, weight float64, unit string) error { if weight > 0 { if value < 0 { - log.Printf("W! Value cannot be negative: %v", value) - return + return errors.New("negative value") } //sample count regularDist.sampleCount += weight @@ -91,16 +91,17 @@ func (regularDist *RegularDistribution) AddEntryWithUnit(value float64, weight f if regularDist.unit == "" { regularDist.unit = unit } else if regularDist.unit != unit && unit != "" { - log.Printf("D! Multiple units are dected: %s, %s", regularDist.unit, unit) + log.Printf("D! Multiple units are detected: %s, %s", regularDist.unit, unit) } } else { log.Printf("D! Weight should be larger than 0: %v", weight) } + return nil } // weight is 1/samplingRate -func (regularDist *RegularDistribution) AddEntry(value float64, weight float64) { - regularDist.AddEntryWithUnit(value, weight, "") +func (regularDist *RegularDistribution) AddEntry(value float64, weight float64) error { + return regularDist.AddEntryWithUnit(value, weight, "") } func (regularDist *RegularDistribution) AddDistribution(distribution distribution.Distribution) { diff --git a/metric/distribution/regular/regular_distribution_test.go b/metric/distribution/regular/regular_distribution_test.go index 2742d3d1fb7..72b33084b4a 100644 --- a/metric/distribution/regular/regular_distribution_test.go +++ b/metric/distribution/regular/regular_distribution_test.go @@ -13,9 +13,9 @@ func TestSEH1Distribution(t *testing.T) { //dist new and add entry dist := NewRegularDistribution() - dist.AddEntry(20, 1) - dist.AddEntry(30, 1) - dist.AddEntryWithUnit(50, 1, "Count") + assert.NoError(t, dist.AddEntry(20, 1)) + assert.NoError(t, dist.AddEntry(30, 1)) + assert.NoError(t, dist.AddEntryWithUnit(50, 1, "Count")) assert.Equal(t, 100.0, dist.Sum()) assert.Equal(t, 3.0, dist.SampleCount()) diff --git a/metric/distribution/seh1/seh1_distribution.go b/metric/distribution/seh1/seh1_distribution.go index f67254aa8c7..78e25378e78 100644 --- a/metric/distribution/seh1/seh1_distribution.go +++ b/metric/distribution/seh1/seh1_distribution.go @@ -4,6 +4,7 @@ package seh1 import ( + "errors" "log" "math" @@ -75,11 +76,10 @@ func (seh1Distribution *SEH1Distribution) Size() int { } // weight is 1/samplingRate -func (seh1Distribution *SEH1Distribution) AddEntryWithUnit(value float64, weight float64, unit string) { +func (seh1Distribution *SEH1Distribution) AddEntryWithUnit(value float64, weight float64, unit string) error { if weight > 0 { if value < 0 { - log.Printf("W! Histogram value cannot be negative: %v", value) - return + return errors.New("negative value") } //sample count seh1Distribution.sampleCount += weight @@ -102,16 +102,17 @@ func (seh1Distribution *SEH1Distribution) AddEntryWithUnit(value float64, weight if seh1Distribution.unit == "" { seh1Distribution.unit = unit } else if seh1Distribution.unit != unit && unit != "" { - log.Printf("D! Multiple units are dected: %s, %s", seh1Distribution.unit, unit) + log.Printf("D! Multiple units are detected: %s, %s", seh1Distribution.unit, unit) } } else { log.Printf("D! Weight should be larger than 0: %v", weight) } + return nil } // weight is 1/samplingRate -func (seh1Distribution *SEH1Distribution) AddEntry(value float64, weight float64) { - seh1Distribution.AddEntryWithUnit(value, weight, "") +func (seh1Distribution *SEH1Distribution) AddEntry(value float64, weight float64) error { + return seh1Distribution.AddEntryWithUnit(value, weight, "") } func (seh1Distribution *SEH1Distribution) AddDistribution(distribution distribution.Distribution) { diff --git a/metric/distribution/seh1/seh1_distribution_test.go b/metric/distribution/seh1/seh1_distribution_test.go index b708f2b3fcc..411531c879e 100644 --- a/metric/distribution/seh1/seh1_distribution_test.go +++ b/metric/distribution/seh1/seh1_distribution_test.go @@ -13,9 +13,9 @@ func TestSEH1Distribution(t *testing.T) { //dist new and add entry dist := NewSEH1Distribution() - dist.AddEntry(20, 1) - dist.AddEntry(30, 1) - dist.AddEntryWithUnit(50, 1, "Count") + assert.NoError(t, dist.AddEntry(20, 1)) + assert.NoError(t, dist.AddEntry(30, 1)) + assert.NoError(t, dist.AddEntryWithUnit(50, 1, "Count")) assert.Equal(t, 100.0, dist.Sum()) assert.Equal(t, 3.0, dist.SampleCount()) @@ -34,9 +34,9 @@ func TestSEH1Distribution(t *testing.T) { //another dist new and add entry anotherDist := NewSEH1Distribution() - anotherDist.AddEntry(21, 1) - anotherDist.AddEntry(22, 1) - anotherDist.AddEntry(23, 2) + assert.NoError(t, anotherDist.AddEntry(21, 1)) + assert.NoError(t, anotherDist.AddEntry(22, 1)) + assert.NoError(t, anotherDist.AddEntry(23, 2)) assert.Equal(t, 89.0, anotherDist.Sum()) assert.Equal(t, 4.0, anotherDist.SampleCount()) diff --git a/plugins/inputs/statsd/statsd.go b/plugins/inputs/statsd/statsd.go index b2040a29510..86d190f4052 100644 --- a/plugins/inputs/statsd/statsd.go +++ b/plugins/inputs/statsd/statsd.go @@ -540,7 +540,10 @@ func (s *Statsd) aggregate(m metric) { if m.samplerate > 0 { weight = 1.0 / m.samplerate } - field.(distribution.Distribution).AddEntry(m.floatvalue, weight) + err := field.(distribution.Distribution).AddEntry(m.floatvalue, weight) + if err != nil { + log.Printf("W! error: %s, metric: %s, value: %v", err, m.name, m.floatvalue) + } cached.fields[m.field] = field s.timings[m.hash] = cached case "c": diff --git a/plugins/inputs/statsd/statsd_test.go b/plugins/inputs/statsd/statsd_test.go index c946a92a068..3c24ed2e8a9 100644 --- a/plugins/inputs/statsd/statsd_test.go +++ b/plugins/inputs/statsd/statsd_test.go @@ -278,11 +278,11 @@ func TestParse_Timings(t *testing.T) { s.Gather(acc) dist := distribution.NewDistribution() - dist.AddEntry(1, 1) - dist.AddEntry(11, 1) - dist.AddEntry(1, 1) - dist.AddEntry(1, 1) - dist.AddEntry(1, 1) + assert.NoError(t, dist.AddEntry(1, 1)) + assert.NoError(t, dist.AddEntry(11, 1)) + assert.NoError(t, dist.AddEntry(1, 1)) + assert.NoError(t, dist.AddEntry(1, 1)) + assert.NoError(t, dist.AddEntry(1, 1)) metrics := acc.Metrics assert.Equal(t, 1, len(metrics)) @@ -1039,18 +1039,18 @@ func TestParse_Timings_MultipleFieldsWithTemplate(t *testing.T) { s.Gather(acc) dist := distribution.NewDistribution() - dist.AddEntry(1, 1) - dist.AddEntry(11, 1) - dist.AddEntry(1, 1) - dist.AddEntry(1, 1) - dist.AddEntry(1, 1) + assert.NoError(t, dist.AddEntry(1, 1)) + assert.NoError(t, dist.AddEntry(11, 1)) + assert.NoError(t, dist.AddEntry(1, 1)) + assert.NoError(t, dist.AddEntry(1, 1)) + assert.NoError(t, dist.AddEntry(1, 1)) dist2 := distribution.NewDistribution() - dist2.AddEntry(2, 1) - dist2.AddEntry(22, 1) - dist2.AddEntry(2, 1) - dist2.AddEntry(2, 1) - dist2.AddEntry(2, 1) + assert.NoError(t, dist2.AddEntry(2, 1)) + assert.NoError(t, dist2.AddEntry(22, 1)) + assert.NoError(t, dist2.AddEntry(2, 1)) + assert.NoError(t, dist2.AddEntry(2, 1)) + assert.NoError(t, dist2.AddEntry(2, 1)) metrics := acc.Metrics assert.Equal(t, 1, len(metrics)) @@ -1094,18 +1094,18 @@ func TestParse_Timings_MultipleFieldsWithoutTemplate(t *testing.T) { s.Gather(acc) dist := distribution.NewDistribution() - dist.AddEntry(1, 1) - dist.AddEntry(11, 1) - dist.AddEntry(1, 1) - dist.AddEntry(1, 1) - dist.AddEntry(1, 1) + assert.NoError(t, dist.AddEntry(1, 1)) + assert.NoError(t, dist.AddEntry(11, 1)) + assert.NoError(t, dist.AddEntry(1, 1)) + assert.NoError(t, dist.AddEntry(1, 1)) + assert.NoError(t, dist.AddEntry(1, 1)) dist2 := distribution.NewDistribution() - dist2.AddEntry(2, 1) - dist2.AddEntry(22, 1) - dist2.AddEntry(2, 1) - dist2.AddEntry(2, 1) - dist2.AddEntry(2, 1) + assert.NoError(t, dist2.AddEntry(2, 1)) + assert.NoError(t, dist2.AddEntry(22, 1)) + assert.NoError(t, dist2.AddEntry(2, 1)) + assert.NoError(t, dist2.AddEntry(2, 1)) + assert.NoError(t, dist2.AddEntry(2, 1)) metrics := acc.Metrics assert.Equal(t, 2, len(metrics)) diff --git a/plugins/outputs/cloudwatch/aggregator.go b/plugins/outputs/cloudwatch/aggregator.go index ba819fed381..b198bffd67f 100644 --- a/plugins/outputs/cloudwatch/aggregator.go +++ b/plugins/outputs/cloudwatch/aggregator.go @@ -178,7 +178,10 @@ func (durationAgg *durationAggregator) aggregating() { if dist != nil { existingDist.AddDistribution(dist) } else { - existingDist.AddEntry(value, 1) + err = existingDist.AddEntry(value, 1) + if err != nil { + log.Printf("W! error: %s, metric %s, value %v", err, m.Name(), value) + } } } case <-durationAgg.ticker.C: diff --git a/plugins/outputs/cloudwatch/cloudwatch_test.go b/plugins/outputs/cloudwatch/cloudwatch_test.go index 9bd231d2155..db1320b21af 100644 --- a/plugins/outputs/cloudwatch/cloudwatch_test.go +++ b/plugins/outputs/cloudwatch/cloudwatch_test.go @@ -96,7 +96,9 @@ func TestBuildMetricDatums(t *testing.T) { } invalidDistribution := distribution.NewDistribution() - invalidDistribution.AddEntry(-1, 1) + err := invalidDistribution.AddEntry(-1, 1) + expectedErrMsg := "negative value" + assert.EqualError(err, expectedErrMsg) invalidMetrics := []telegraf.Metric{ testutil.TestMetric("Foo"), testutil.TestMetric(invalidDistribution), diff --git a/plugins/outputs/cloudwatch/util_test.go b/plugins/outputs/cloudwatch/util_test.go index d576a0f1f45..c03ae1920e6 100644 --- a/plugins/outputs/cloudwatch/util_test.go +++ b/plugins/outputs/cloudwatch/util_test.go @@ -58,9 +58,9 @@ func TestResize(t *testing.T) { assert.Equal(t, float64(1), sampleCount) assert.Equal(t, float64(1), sum) - dist.AddEntry(2, 1) - dist.AddEntry(3, 1) - dist.AddEntry(4, 1) + assert.NoError(t, dist.AddEntry(2, 1)) + assert.NoError(t, dist.AddEntry(3, 1)) + assert.NoError(t, dist.AddEntry(4, 1)) distList = resize(dist, maxListSize) assert.Equal(t, 2, len(distList))