2016-01-27 21:21:36 +00:00
|
|
|
package agent
|
|
|
|
|
|
|
|
import (
|
|
|
|
"fmt"
|
|
|
|
"log"
|
|
|
|
"math"
|
2016-07-25 12:09:49 +00:00
|
|
|
"sync/atomic"
|
2016-01-27 21:21:36 +00:00
|
|
|
"time"
|
|
|
|
|
2016-01-27 23:15:14 +00:00
|
|
|
"github.com/influxdata/telegraf"
|
2016-01-27 21:21:36 +00:00
|
|
|
"github.com/influxdata/telegraf/internal/models"
|
|
|
|
)
|
|
|
|
|
|
|
|
func NewAccumulator(
|
2016-07-28 11:31:11 +00:00
|
|
|
inputConfig *models.InputConfig,
|
2016-01-27 23:15:14 +00:00
|
|
|
metrics chan telegraf.Metric,
|
2016-01-27 21:21:36 +00:00
|
|
|
) *accumulator {
|
|
|
|
acc := accumulator{}
|
2016-01-27 23:15:14 +00:00
|
|
|
acc.metrics = metrics
|
2016-01-27 21:21:36 +00:00
|
|
|
acc.inputConfig = inputConfig
|
2016-06-13 14:21:11 +00:00
|
|
|
acc.precision = time.Nanosecond
|
2016-01-27 21:21:36 +00:00
|
|
|
return &acc
|
|
|
|
}
|
|
|
|
|
|
|
|
type accumulator struct {
|
2016-01-27 23:15:14 +00:00
|
|
|
metrics chan telegraf.Metric
|
2016-01-27 21:21:36 +00:00
|
|
|
|
|
|
|
defaultTags map[string]string
|
|
|
|
|
|
|
|
debug bool
|
2016-05-19 15:36:58 +00:00
|
|
|
// print every point added to the accumulator
|
|
|
|
trace bool
|
2016-01-27 21:21:36 +00:00
|
|
|
|
2016-07-28 11:31:11 +00:00
|
|
|
inputConfig *models.InputConfig
|
2016-01-27 21:21:36 +00:00
|
|
|
|
2016-06-13 14:21:11 +00:00
|
|
|
precision time.Duration
|
2016-07-25 12:09:49 +00:00
|
|
|
|
|
|
|
errCount uint64
|
2016-01-27 21:21:36 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
func (ac *accumulator) Add(
|
|
|
|
measurement string,
|
|
|
|
value interface{},
|
|
|
|
tags map[string]string,
|
|
|
|
t ...time.Time,
|
|
|
|
) {
|
|
|
|
fields := make(map[string]interface{})
|
|
|
|
fields["value"] = value
|
2016-02-20 05:35:12 +00:00
|
|
|
|
|
|
|
if !ac.inputConfig.Filter.ShouldNamePass(measurement) {
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
2016-01-27 21:21:36 +00:00
|
|
|
ac.AddFields(measurement, fields, tags, t...)
|
|
|
|
}
|
|
|
|
|
|
|
|
func (ac *accumulator) AddFields(
|
|
|
|
measurement string,
|
|
|
|
fields map[string]interface{},
|
|
|
|
tags map[string]string,
|
|
|
|
t ...time.Time,
|
|
|
|
) {
|
|
|
|
if len(fields) == 0 || len(measurement) == 0 {
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
2016-02-20 05:35:12 +00:00
|
|
|
if !ac.inputConfig.Filter.ShouldNamePass(measurement) {
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
2016-01-27 21:21:36 +00:00
|
|
|
if !ac.inputConfig.Filter.ShouldTagsPass(tags) {
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
|
|
|
// Override measurement name if set
|
|
|
|
if len(ac.inputConfig.NameOverride) != 0 {
|
|
|
|
measurement = ac.inputConfig.NameOverride
|
|
|
|
}
|
|
|
|
// Apply measurement prefix and suffix if set
|
|
|
|
if len(ac.inputConfig.MeasurementPrefix) != 0 {
|
|
|
|
measurement = ac.inputConfig.MeasurementPrefix + measurement
|
|
|
|
}
|
|
|
|
if len(ac.inputConfig.MeasurementSuffix) != 0 {
|
|
|
|
measurement = measurement + ac.inputConfig.MeasurementSuffix
|
|
|
|
}
|
|
|
|
|
|
|
|
if tags == nil {
|
|
|
|
tags = make(map[string]string)
|
|
|
|
}
|
2016-04-19 00:12:58 +00:00
|
|
|
// Apply plugin-wide tags if set
|
|
|
|
for k, v := range ac.inputConfig.Tags {
|
2016-05-19 09:01:19 +00:00
|
|
|
if _, ok := tags[k]; !ok {
|
|
|
|
tags[k] = v
|
|
|
|
}
|
|
|
|
}
|
|
|
|
// Apply daemon-wide tags if set
|
|
|
|
for k, v := range ac.defaultTags {
|
|
|
|
if _, ok := tags[k]; !ok {
|
|
|
|
tags[k] = v
|
|
|
|
}
|
2016-01-27 21:21:36 +00:00
|
|
|
}
|
2016-04-12 23:06:27 +00:00
|
|
|
ac.inputConfig.Filter.FilterTags(tags)
|
2016-01-27 21:21:36 +00:00
|
|
|
|
|
|
|
result := make(map[string]interface{})
|
|
|
|
for k, v := range fields {
|
|
|
|
// Filter out any filtered fields
|
|
|
|
if ac.inputConfig != nil {
|
2016-02-20 05:35:12 +00:00
|
|
|
if !ac.inputConfig.Filter.ShouldFieldsPass(k) {
|
2016-01-27 21:21:36 +00:00
|
|
|
continue
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
// Validate uint64 and float64 fields
|
|
|
|
switch val := v.(type) {
|
|
|
|
case uint64:
|
|
|
|
// InfluxDB does not support writing uint64
|
|
|
|
if val < uint64(9223372036854775808) {
|
|
|
|
result[k] = int64(val)
|
|
|
|
} else {
|
|
|
|
result[k] = int64(9223372036854775807)
|
|
|
|
}
|
2016-03-07 14:46:23 +00:00
|
|
|
continue
|
2016-01-27 21:21:36 +00:00
|
|
|
case float64:
|
|
|
|
// NaNs are invalid values in influxdb, skip measurement
|
|
|
|
if math.IsNaN(val) || math.IsInf(val, 0) {
|
|
|
|
if ac.debug {
|
|
|
|
log.Printf("Measurement [%s] field [%s] has a NaN or Inf "+
|
|
|
|
"field, skipping",
|
|
|
|
measurement, k)
|
|
|
|
}
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
}
|
2016-03-07 14:46:23 +00:00
|
|
|
|
|
|
|
result[k] = v
|
2016-01-27 21:21:36 +00:00
|
|
|
}
|
|
|
|
fields = nil
|
|
|
|
if len(result) == 0 {
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
|
|
|
var timestamp time.Time
|
|
|
|
if len(t) > 0 {
|
|
|
|
timestamp = t[0]
|
|
|
|
} else {
|
|
|
|
timestamp = time.Now()
|
|
|
|
}
|
2016-06-13 14:21:11 +00:00
|
|
|
timestamp = timestamp.Round(ac.precision)
|
2016-01-27 21:21:36 +00:00
|
|
|
|
2016-01-27 23:15:14 +00:00
|
|
|
m, err := telegraf.NewMetric(measurement, tags, result, timestamp)
|
2016-01-27 21:21:36 +00:00
|
|
|
if err != nil {
|
|
|
|
log.Printf("Error adding point [%s]: %s\n", measurement, err.Error())
|
|
|
|
return
|
|
|
|
}
|
2016-05-19 15:36:58 +00:00
|
|
|
if ac.trace {
|
2016-01-27 23:15:14 +00:00
|
|
|
fmt.Println("> " + m.String())
|
2016-01-27 21:21:36 +00:00
|
|
|
}
|
2016-01-27 23:15:14 +00:00
|
|
|
ac.metrics <- m
|
2016-01-27 21:21:36 +00:00
|
|
|
}
|
|
|
|
|
2016-07-25 12:09:49 +00:00
|
|
|
// AddError passes a runtime error to the accumulator.
|
|
|
|
// The error will be tagged with the plugin name and written to the log.
|
|
|
|
func (ac *accumulator) AddError(err error) {
|
|
|
|
if err == nil {
|
|
|
|
return
|
|
|
|
}
|
|
|
|
atomic.AddUint64(&ac.errCount, 1)
|
|
|
|
//TODO suppress/throttle consecutive duplicate errors?
|
|
|
|
log.Printf("ERROR in input [%s]: %s", ac.inputConfig.Name, err)
|
|
|
|
}
|
|
|
|
|
2016-01-27 21:21:36 +00:00
|
|
|
func (ac *accumulator) Debug() bool {
|
|
|
|
return ac.debug
|
|
|
|
}
|
|
|
|
|
|
|
|
func (ac *accumulator) SetDebug(debug bool) {
|
|
|
|
ac.debug = debug
|
|
|
|
}
|
|
|
|
|
2016-05-19 15:36:58 +00:00
|
|
|
func (ac *accumulator) Trace() bool {
|
|
|
|
return ac.trace
|
|
|
|
}
|
|
|
|
|
|
|
|
func (ac *accumulator) SetTrace(trace bool) {
|
|
|
|
ac.trace = trace
|
|
|
|
}
|
|
|
|
|
2016-06-13 14:21:11 +00:00
|
|
|
// SetPrecision takes two time.Duration objects. If the first is non-zero,
|
|
|
|
// it sets that as the precision. Otherwise, it takes the second argument
|
|
|
|
// as the order of time that the metrics should be rounded to, with the
|
|
|
|
// maximum being 1s.
|
|
|
|
func (ac *accumulator) SetPrecision(precision, interval time.Duration) {
|
|
|
|
if precision > 0 {
|
|
|
|
ac.precision = precision
|
|
|
|
return
|
|
|
|
}
|
|
|
|
switch {
|
|
|
|
case interval >= time.Second:
|
|
|
|
ac.precision = time.Second
|
|
|
|
case interval >= time.Millisecond:
|
|
|
|
ac.precision = time.Millisecond
|
|
|
|
case interval >= time.Microsecond:
|
|
|
|
ac.precision = time.Microsecond
|
|
|
|
default:
|
|
|
|
ac.precision = time.Nanosecond
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func (ac *accumulator) DisablePrecision() {
|
|
|
|
ac.precision = time.Nanosecond
|
|
|
|
}
|
|
|
|
|
2016-01-27 21:21:36 +00:00
|
|
|
func (ac *accumulator) setDefaultTags(tags map[string]string) {
|
|
|
|
ac.defaultTags = tags
|
|
|
|
}
|
|
|
|
|
|
|
|
func (ac *accumulator) addDefaultTag(key, value string) {
|
2016-03-07 14:46:23 +00:00
|
|
|
if ac.defaultTags == nil {
|
|
|
|
ac.defaultTags = make(map[string]string)
|
|
|
|
}
|
2016-01-27 21:21:36 +00:00
|
|
|
ac.defaultTags[key] = value
|
|
|
|
}
|