Compare commits

...
Author SHA1 Message Date
Thomas Gruber 5f9d182903 Fix linting issues 2026-07-31 16:44:43 +02:00
Thomas Gruber 80dbd41b91 Extend multiChanTicker with GetDuration 2026-07-31 16:36:26 +02:00
Thomas Gruber 4c1f810f53 Simplify cache, calculate with expr 2026-07-31 16:31:04 +02:00
Michael Panzlaff 789ba7f4a5 Revert "Prepare for FAU"
This reverts commit 9bd9a07810.

This should never have made it into our main repo.
2026-07-21 13:42:11 +02:00
Michael Panzlaff 9bd9a07810 Prepare for FAU 2026-07-21 13:40:42 +02:00
Michael PanzlaffandGitHub a0aedff202 Merge pull request #235 from ClusterCockpit/marker-commits
Add marker metrics
2026-07-21 12:51:37 +02:00
Michael Panzlaff 6944454592 Update configuration documentation for new marker metrics 2026-07-21 12:47:16 +02:00
Michael Panzlaff 9e773e2e85 Fix marker metric creation 2026-07-20 16:16:10 +02:00
Michael Panzlaff c13bb55735 Actually send the created marker metrics 2026-07-20 15:39:40 +02:00
Michael Panzlaff 6124971b2e collectorManager: Make markers opt-in 2026-07-15 18:04:56 +02:00
Michael Panzlaff cdac3c9163 Add ccmc-{begin,end} markers
These markers are sent out at the beginning and end of a collection run.
This can be used to reconstruct, which metrics belong together and were
obtained during the same tick.
2026-07-15 15:45:15 +02:00
dependabot[bot]andThomas Gruber 70364d084b Bump golang.org/x/sys from 0.45.0 to 0.47.0
Bumps [golang.org/x/sys](https://github.com/golang/sys) from 0.45.0 to 0.47.0.
- [Commits](https://github.com/golang/sys/compare/v0.45.0...v0.47.0)

---
updated-dependencies:
- dependency-name: golang.org/x/sys
  dependency-version: 0.47.0
  dependency-type: direct:production
  update-type: version-update:semver-minor
...

Signed-off-by: dependabot[bot] <support@github.com>
2026-07-13 12:06:53 +02:00
dependabot[bot]andThomas Gruber ec0868eba9 Bump github.com/NVIDIA/go-nvml from 0.13.2-0 to 0.13.3-1
Bumps [github.com/NVIDIA/go-nvml](https://github.com/NVIDIA/go-nvml) from 0.13.2-0 to 0.13.3-1.
- [Release notes](https://github.com/NVIDIA/go-nvml/releases)
- [Commits](https://github.com/NVIDIA/go-nvml/compare/v0.13.2-0...v0.13.3-1)

---
updated-dependencies:
- dependency-name: github.com/NVIDIA/go-nvml
  dependency-version: 0.13.3-1
  dependency-type: direct:production
  update-type: version-update:semver-patch
...

Signed-off-by: dependabot[bot] <support@github.com>
2026-07-13 12:06:37 +02:00
Michael Panzlaff b8e76e3bf0 megwareEureka: Fix collector name 2026-07-07 19:03:57 +02:00
Michael Panzlaff f028c83a75 megwareEureka: Remove debug output
This should never have been in the release...
2026-07-07 18:47:02 +02:00
Michael PanzlaffandGitHub 4a859af358 Merge pull request #231 from ClusterCockpit/likwid-lint-fix
likwidMetric: Replace deprecated x/thread package
2026-07-07 17:08:46 +02:00
15 changed files with 807 additions and 143 deletions
+1
View File
@@ -33,6 +33,7 @@ There is a main configuration file with basic settings that point to the other c
"receivers-file" : "receivers.json", "receivers-file" : "receivers.json",
"router-file" : "router.json", "router-file" : "router.json",
"main": { "main": {
"enable-markers": false,
"interval": "10s", "interval": "10s",
"duration": "1s" "duration": "1s"
} }
+4 -1
View File
@@ -31,11 +31,13 @@ import (
type CentralConfigFile struct { type CentralConfigFile struct {
Interval string `json:"interval"` Interval string `json:"interval"`
Duration string `json:"duration"` Duration string `json:"duration"`
EnableMarkers bool `json:"enable-markers"`
} }
type RuntimeConfig struct { type RuntimeConfig struct {
Interval time.Duration Interval time.Duration
Duration time.Duration Duration time.Duration
EnableMarkers bool
CliArgs map[string]string CliArgs map[string]string
ConfigFile CentralConfigFile ConfigFile CentralConfigFile
@@ -157,6 +159,7 @@ func mainFunc() int {
cclog.Error("The interval should be greater than duration") cclog.Error("The interval should be greater than duration")
return 1 return 1
} }
rcfg.EnableMarkers = rcfg.ConfigFile.EnableMarkers
routerConf := ccconf.GetPackageConfig("router") routerConf := ccconf.GetPackageConfig("router")
if len(routerConf) == 0 { if len(routerConf) == 0 {
@@ -199,7 +202,7 @@ func mainFunc() int {
rcfg.MetricRouter.AddOutput(RouterToSinksChannel) rcfg.MetricRouter.AddOutput(RouterToSinksChannel)
// Create new collector manager // Create new collector manager
rcfg.CollectManager, err = collectors.New(rcfg.MultiChanTicker, rcfg.Duration, &rcfg.Sync, collectorConf) rcfg.CollectManager, err = collectors.New(rcfg.MultiChanTicker, rcfg.Duration, rcfg.EnableMarkers, &rcfg.Sync, collectorConf)
if err != nil { if err != nil {
cclog.Error(err.Error()) cclog.Error(err.Error())
return 1 return 1
+22 -4
View File
@@ -66,11 +66,12 @@ type collectorManager struct {
config map[string]json.RawMessage // json encoded config for collector manager config map[string]json.RawMessage // json encoded config for collector manager
collector_wg sync.WaitGroup // internally used wait group for the parallel reading of collector collector_wg sync.WaitGroup // internally used wait group for the parallel reading of collector
parallel_run bool // Flag whether the collectors are currently read in parallel parallel_run bool // Flag whether the collectors are currently read in parallel
enableMarkers bool // Send ccmc-{begin,end} metrics
} }
// Metric collector manager access functions // Metric collector manager access functions
type CollectorManager interface { type CollectorManager interface {
Init(ticker mct.MultiChanTicker, duration time.Duration, wg *sync.WaitGroup, collectConfig json.RawMessage) error Init(ticker mct.MultiChanTicker, duration time.Duration, enableMarkers bool, wg *sync.WaitGroup, collectConfig json.RawMessage) error
AddOutput(output chan lp.CCMessage) AddOutput(output chan lp.CCMessage)
Start() Start()
Close() Close()
@@ -83,7 +84,7 @@ type CollectorManager interface {
// * ticker (from variable ticker) // * ticker (from variable ticker)
// * configuration (read from config file in variable collectConfigFile) // * configuration (read from config file in variable collectConfigFile)
// Initialization is done for all configured collectors // Initialization is done for all configured collectors
func (cm *collectorManager) Init(ticker mct.MultiChanTicker, duration time.Duration, wg *sync.WaitGroup, collectConfig json.RawMessage) error { func (cm *collectorManager) Init(ticker mct.MultiChanTicker, duration time.Duration, enableMarkers bool, wg *sync.WaitGroup, collectConfig json.RawMessage) error {
cm.collectors = make([]MetricCollector, 0) cm.collectors = make([]MetricCollector, 0)
cm.serial = make([]MetricCollector, 0) cm.serial = make([]MetricCollector, 0)
cm.output = nil cm.output = nil
@@ -91,6 +92,7 @@ func (cm *collectorManager) Init(ticker mct.MultiChanTicker, duration time.Durat
cm.wg = wg cm.wg = wg
cm.ticker = ticker cm.ticker = ticker
cm.duration = duration cm.duration = duration
cm.enableMarkers = enableMarkers
d := json.NewDecoder(bytes.NewReader(collectConfig)) d := json.NewDecoder(bytes.NewReader(collectConfig))
d.DisallowUnknownFields() d.DisallowUnknownFields()
@@ -148,6 +150,14 @@ func (cm *collectorManager) Start() {
done() done()
return return
case t := <-tick: case t := <-tick:
if cm.enableMarkers {
m, err := lp.NewMetric("ccmc-begin", map[string]string{"type": "node"}, nil, 0, time.Now())
if err != nil {
cclog.ComponentErrorf("CollectorManager", "Unable to create marker metric: %v", err)
} else {
cm.output <- m
}
}
cm.parallel_run = true cm.parallel_run = true
for _, c := range cm.collectors { for _, c := range cm.collectors {
// Wait for done signal or execute the collector // Wait for done signal or execute the collector
@@ -179,6 +189,14 @@ func (cm *collectorManager) Start() {
c.Read(cm.duration, cm.output) c.Read(cm.duration, cm.output)
} }
} }
if cm.enableMarkers {
m, err := lp.NewMetric("ccmc-end", map[string]string{"type": "node"}, nil, 0, time.Now())
if err != nil {
cclog.ComponentErrorf("CollectorManager", "Unable to create marker metric: %v", err)
} else {
cm.output <- m
}
}
} }
} }
}) })
@@ -201,9 +219,9 @@ func (cm *collectorManager) Close() {
} }
// New creates a new initialized metric collector manager // New creates a new initialized metric collector manager
func New(ticker mct.MultiChanTicker, duration time.Duration, wg *sync.WaitGroup, collectConfig json.RawMessage) (CollectorManager, error) { func New(ticker mct.MultiChanTicker, duration time.Duration, enableMarkers bool, wg *sync.WaitGroup, collectConfig json.RawMessage) (CollectorManager, error) {
cm := new(collectorManager) cm := new(collectorManager)
err := cm.Init(ticker, duration, wg, collectConfig) err := cm.Init(ticker, duration, enableMarkers, wg, collectConfig)
if err != nil { if err != nil {
return nil, err return nil, err
} }
+1 -4
View File
@@ -65,7 +65,7 @@ func (m *MegwareEurekaCollector) Init(config json.RawMessage) error {
return nil return nil
} }
m.name = "MegwareEureka" m.name = "MegwareEurekaCollector"
if err := m.setup(); err != nil { if err := m.setup(); err != nil {
return fmt.Errorf("%s Init(): setup() call failed: %w", m.name, err) return fmt.Errorf("%s Init(): setup() call failed: %w", m.name, err)
} }
@@ -127,9 +127,6 @@ func (m *MegwareEurekaCollector) readMpsData() (*mpsData, error) {
return nil, fmt.Errorf("unable to decode u20 JSON output: %w (stdout=%s)", err, stdout.String()) return nil, fmt.Errorf("unable to decode u20 JSON output: %w (stdout=%s)", err, stdout.String())
} }
fmt.Printf("string: %+v\n", stdout.String())
fmt.Printf("obj: %+v\n", u20output)
u20output.GetMpsPollValues.Timestamp = time.Now() u20output.GetMpsPollValues.Timestamp = time.Now()
return &u20output.GetMpsPollValues, nil return &u20output.GetMpsPollValues, nil
+1
View File
@@ -4,6 +4,7 @@
"receivers-file" : "./receivers.json", "receivers-file" : "./receivers.json",
"router-file" : "./router.json", "router-file" : "./router.json",
"main" : { "main" : {
"enable-markers": false,
"interval": "10s", "interval": "10s",
"duration": "1s" "duration": "1s"
} }
+2 -2
View File
@@ -5,12 +5,12 @@ go 1.25.0
require ( require (
github.com/ClusterCockpit/cc-lib/v2 v2.12.0 github.com/ClusterCockpit/cc-lib/v2 v2.12.0
github.com/ClusterCockpit/go-rocm-smi v0.4.0 github.com/ClusterCockpit/go-rocm-smi v0.4.0
github.com/NVIDIA/go-nvml v0.13.2-0 github.com/NVIDIA/go-nvml v0.13.3-1
github.com/PaesslerAG/gval v1.2.4 github.com/PaesslerAG/gval v1.2.4
github.com/fsnotify/fsnotify v1.10.1 github.com/fsnotify/fsnotify v1.10.1
github.com/tklauser/go-sysconf v0.4.0 github.com/tklauser/go-sysconf v0.4.0
golang.design/x/runtime v0.3.0 golang.design/x/runtime v0.3.0
golang.org/x/sys v0.45.0 golang.org/x/sys v0.47.0
) )
require ( require (
+4 -4
View File
@@ -13,8 +13,8 @@ github.com/Microsoft/go-winio v0.6.1/go.mod h1:LRdKpFKfdobln8UmuiYcKPot9D2v6svN5
github.com/Microsoft/hcsshim v0.11.4 h1:68vKo2VN8DE9AdN4tnkWnmdhqdbpUFM8OF3Airm7fz8= github.com/Microsoft/hcsshim v0.11.4 h1:68vKo2VN8DE9AdN4tnkWnmdhqdbpUFM8OF3Airm7fz8=
github.com/Microsoft/hcsshim v0.11.4/go.mod h1:smjE4dvqPX9Zldna+t5FG3rnoHhaB7QYxPRqGcpAD9w= github.com/Microsoft/hcsshim v0.11.4/go.mod h1:smjE4dvqPX9Zldna+t5FG3rnoHhaB7QYxPRqGcpAD9w=
github.com/NVIDIA/go-nvml v0.13.0-1/go.mod h1:+KNA7c7gIBH7SKSJ1ntlwkfN80zdx8ovl4hrK3LmPt4= github.com/NVIDIA/go-nvml v0.13.0-1/go.mod h1:+KNA7c7gIBH7SKSJ1ntlwkfN80zdx8ovl4hrK3LmPt4=
github.com/NVIDIA/go-nvml v0.13.2-0 h1:7M4cFG62wSUHw8i0XSiNU7ejKODytTS6ZrW/vgB2NSI= github.com/NVIDIA/go-nvml v0.13.3-1 h1:P76U2h88OZSiMtdhRsJjSF5DXyXUqHIXKeDicVAaae0=
github.com/NVIDIA/go-nvml v0.13.2-0/go.mod h1:ahi2psRYoa+wYUBIrZPRO+wJs9lcvMhxSSkjjvsJJNQ= github.com/NVIDIA/go-nvml v0.13.3-1/go.mod h1:ahi2psRYoa+wYUBIrZPRO+wJs9lcvMhxSSkjjvsJJNQ=
github.com/PaesslerAG/gval v1.2.4 h1:rhX7MpjJlcxYwL2eTTYIOBUyEKZ+A96T9vQySWkVUiU= github.com/PaesslerAG/gval v1.2.4 h1:rhX7MpjJlcxYwL2eTTYIOBUyEKZ+A96T9vQySWkVUiU=
github.com/PaesslerAG/gval v1.2.4/go.mod h1:XRFLwvmkTEdYziLdaCeCa5ImcGVrfQbeNUbVR+C6xac= github.com/PaesslerAG/gval v1.2.4/go.mod h1:XRFLwvmkTEdYziLdaCeCa5ImcGVrfQbeNUbVR+C6xac=
github.com/PaesslerAG/jsonpath v0.1.0 h1:gADYeifvlqK3R3i2cR5B4DGgxLXIPb3TRTH1mGi0jPI= github.com/PaesslerAG/jsonpath v0.1.0 h1:gADYeifvlqK3R3i2cR5B4DGgxLXIPb3TRTH1mGi0jPI=
@@ -184,8 +184,8 @@ golang.org/x/mod v0.13.0 h1:I/DsJXRlw/8l/0c24sM9yb0T4z9liZTduXvdAWYiysY=
golang.org/x/mod v0.13.0/go.mod h1:hTbmBsO62+eylJbnUtE2MGJUyE7QWk4xUqPFrRgJ+7c= golang.org/x/mod v0.13.0/go.mod h1:hTbmBsO62+eylJbnUtE2MGJUyE7QWk4xUqPFrRgJ+7c=
golang.org/x/net v0.53.0 h1:d+qAbo5L0orcWAr0a9JweQpjXF19LMXJE8Ey7hwOdUA= golang.org/x/net v0.53.0 h1:d+qAbo5L0orcWAr0a9JweQpjXF19LMXJE8Ey7hwOdUA=
golang.org/x/net v0.53.0/go.mod h1:JvMuJH7rrdiCfbeHoo3fCQU24Lf5JJwT9W3sJFulfgs= golang.org/x/net v0.53.0/go.mod h1:JvMuJH7rrdiCfbeHoo3fCQU24Lf5JJwT9W3sJFulfgs=
golang.org/x/sys v0.45.0 h1:dO4czNzziLiiXplLQgBCEpCvXQ3dnkn0SdaZSYdQ+FY= golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs=
golang.org/x/sys v0.45.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
golang.org/x/time v0.15.0 h1:bbrp8t3bGUeFOx08pvsMYRTCVSMk89u4tKbNOZbp88U= golang.org/x/time v0.15.0 h1:bbrp8t3bGUeFOx08pvsMYRTCVSMk89u4tKbNOZbp88U=
golang.org/x/time v0.15.0/go.mod h1:Y4YMaQmXwGQZoFaVFk4YpCt4FLQMYKZe9oeV/f4MSno= golang.org/x/time v0.15.0/go.mod h1:Y4YMaQmXwGQZoFaVFk4YpCt4FLQMYKZe9oeV/f4MSno=
golang.org/x/tools v0.14.0 h1:jvNa2pY0M4r62jkRQ6RwEZZyPcymeL9XZMLBbV7U2nc= golang.org/x/tools v0.14.0 h1:jvNa2pY0M4r62jkRQ6RwEZZyPcymeL9XZMLBbV7U2nc=
@@ -69,8 +69,8 @@ var metricCacheLanguage = gval.NewLanguage(
gval.Function("getNumaCpuList", getCpuListOfNumaDomainFunc), gval.Function("getNumaCpuList", getCpuListOfNumaDomainFunc),
gval.Function("getDieCpuList", getCpuListOfDieFunc), gval.Function("getDieCpuList", getCpuListOfDieFunc),
gval.Function("getCoreCpuList", getCpuListOfCoreFunc), gval.Function("getCoreCpuList", getCpuListOfCoreFunc),
gval.Function("getCpuList", getCpuListOfNode), gval.Function("getCpuList", getCpuListOfNodeFunc),
gval.Function("getCpuListOfType", getCpuListOfType), gval.Function("getCpuListOfType", getCpuListOfTypeFunc),
) )
var language gval.Language = gval.NewLanguage( var language gval.Language = gval.NewLanguage(
@@ -0,0 +1,449 @@
package metricAggregator
import (
"fmt"
"maps"
"math"
"slices"
"strings"
"sync"
"time"
cclog "github.com/ClusterCockpit/cc-lib/v2/ccLogger"
lp "github.com/ClusterCockpit/cc-lib/v2/ccMessage"
"github.com/expr-lang/expr"
"github.com/expr-lang/expr/vm"
)
type MetricAggregatorExprIntervalConfig struct {
Name string `json:"name"` // Metric name for the new metric
Function string `json:"function"` // Function to apply on the metric
Condition string `json:"if"` // Condition for applying function
Tags map[string]string `json:"tags,omitempty"` // Tags for the new metric
Meta map[string]string `json:"meta,omitempty"` // Meta information for the new metric
ValueType string `json:"value_type,omitempty"`
IncludeNumCacheIntervals int `json:"expr_on_num_intervals,omitempty"`
exprCond *vm.Program
exprFunc *vm.Program
}
type metricAggregatorExpr struct {
constants map[string]any
output chan lp.CCMessage
aggregations []MetricAggregatorExprIntervalConfig
}
var paramMapPool = sync.Pool{
New: func() any {
return make(map[string]any)
},
}
func sanitizeExprString(key string) string {
return strings.ReplaceAll(key, "type-id", "typeid")
}
// AddAggregation(name, function, condition string, tags, meta map[string]string) error
// DeleteAggregation(name string) error
// Init(output chan lp.CCMessage) error
// Eval(starttime time.Time, endtime time.Time, metrics []lp.CCMessage)
func (m *metricAggregatorExpr) Init(output chan lp.CCMessage) error {
m.output = output
m.constants = make(map[string]any)
m.aggregations = make([]MetricAggregatorExprIntervalConfig, 0)
return nil
}
func GetParamMap(point lp.CCMessage) map[string]any {
params := paramMapPool.Get().(map[string]any)
clear(params)
// Put metric name into params map
params["name"] = point.Name()
// Put full message into params map
params["message"] = point
params["msg"] = point
// Put timestamp into params map
params["timestamp"] = point.Time().Unix()
params["time"] = params["timestamp"]
// Put fields into params map
fields := paramMapPool.Get().(map[string]any)
clear(fields)
for key, value := range point.Fields() {
fields[key] = value
switch key {
case "value":
params["messagetype"] = "metric"
params["value"] = value
params["metric"] = value
case "event":
params["messagetype"] = "event"
params["event"] = value
case "control":
params["messagetype"] = "control"
params["control"] = value
case "log":
params["messagetype"] = "log"
params["log"] = value
default:
params["messagetype"] = "unknown"
}
}
params["msgtype"] = params["messagetype"]
params["fields"] = fields
params["field"] = fields
// Put tags into params map
tags := paramMapPool.Get().(map[string]any)
clear(tags)
for key, value := range point.Tags() {
tags[sanitizeExprString(key)] = value
}
params["tags"] = tags
params["tag"] = tags
// Put meta information into params map
meta := paramMapPool.Get().(map[string]any)
clear(meta)
for key, value := range point.Meta() {
meta[sanitizeExprString(key)] = value
}
params["meta"] = meta
return params
}
var baseenv_multi_message = map[string]any{
"name": "",
"starttime": 1234567890,
"endtime": 1234567890,
"messages": make([]lp.CCMessage, 0),
"Median": medianfunc,
"Sum": sumfunc,
"Min": minfunc,
"Max": maxfunc,
"Mean": avgfunc,
"Avg": avgfunc,
"Match": matchfunc,
"getCpuCore": getCpuCoreFunc,
"getCpuSocket": getCpuSocketFunc,
"getCpuNumaDomain": getCpuNumaDomainFunc,
"getCpuDie": getCpuDieFunc,
"getCpuListOfCore": getCpuListOfCoreFunc,
"getCpuListOfSocket": getCpuListOfSocketFunc,
"getCpuListOfNumaDomain": getCpuListOfNumaDomainFunc,
"getCpuListOfDie": getCpuListOfDieFunc,
"getCpuListOfNode": getCpuListOfNodeFunc,
"getCpuListOfType": getCpuListOfTypeFunc,
}
func (m *metricAggregatorExpr) AddAggregationWithType(name, function, condition, valueType string, tags, meta map[string]string) error {
cond, err := expr.Compile(condition, expr.Env(baseenv_multi_message), expr.AsBool(), expr.AllowUndefinedVariables())
if err != nil {
err = fmt.Errorf("failed to compile condition for aggregation %s: %s", name, err.Error())
return err
}
var exprOption = expr.AsFloat64()
switch valueType {
case "int":
case "int32":
exprOption = expr.AsInt()
case "int64":
exprOption = expr.AsInt64()
case "float32":
case "float64":
exprOption = expr.AsFloat64()
case "bool":
exprOption = expr.AsBool()
default:
err := fmt.Errorf("invalid value type '%s' for aggregation %s", valueType, name)
return err
}
f, err := expr.Compile(function, expr.Env(baseenv_multi_message), exprOption, expr.AllowUndefinedVariables())
if err != nil {
err = fmt.Errorf("failed to compile function for aggregation %s: %s", name, err.Error())
return err
}
m.aggregations = append(m.aggregations, MetricAggregatorExprIntervalConfig{
Name: name,
Condition: condition,
Function: function,
Tags: tags,
Meta: meta,
ValueType: valueType,
exprCond: cond,
exprFunc: f,
})
return nil
}
func (m *metricAggregatorExpr) AddAggregation(name, function, condition string, tags, meta map[string]string) error {
cclog.ComponentDebugf("MetricAggregator", "Adding %s", name)
err := m.AddAggregationWithType(name, function, condition, "float64", tags, meta)
cclog.ComponentDebugf("MetricAggregator", "Adding %s returned %v", name, err)
return err
}
func (m *metricAggregatorExpr) Eval(starttime, endtime time.Time, metrics []lp.CCMessage) {
cclog.ComponentDebugf("MetricAggregator", "Calculating %d expressions", len(m.aggregations))
copy_tags := func(tags map[string]string, metrics []lp.CCMessage) map[string]string {
out := make(map[string]string)
for key, value := range tags {
switch value {
case "<copy>":
for _, m := range metrics {
v, err := m.GetTag(key)
if err {
out[key] = v
}
}
default:
out[key] = value
}
}
return out
}
copy_meta := func(meta map[string]string, metrics []lp.CCMessage) map[string]string {
out := make(map[string]string)
for key, value := range meta {
switch value {
case "<copy>":
for _, m := range metrics {
v, err := m.GetMeta(key)
if err {
out[key] = v
}
}
default:
out[key] = value
}
}
return out
}
for _, aggr := range m.aggregations {
selected_metrics := make([]lp.CCMessage, 0)
values := make([]float64, 0)
aggr_vars := make(map[string]any)
maps.Copy(aggr_vars, baseenv_multi_message)
maps.Copy(aggr_vars, m.constants)
aggr_vars["starttime"] = starttime
aggr_vars["endtime"] = endtime
for _, met := range metrics {
met_vars := make(map[string]any)
maps.Copy(met_vars, aggr_vars)
met_vars["message"] = met
maps.Copy(met_vars, GetParamMap(met))
res, err := expr.Run(aggr.exprCond, met_vars)
if err != nil || res == false {
continue
}
if value, ok := met.GetField("value"); ok {
switch v := value.(type) {
case float64:
values = append(values, v)
case float32:
case int:
case int8:
case int16:
case int32:
case int64:
case uint:
case uint8:
case uint16:
case uint32:
case uint64:
values = append(values, float64(v))
case bool:
if v {
values = append(values, float64(1))
} else {
values = append(values, float64(0))
}
default:
cclog.ComponentErrorf("MetricAggregator", "Cannot convert value type for %s", met.ToLineProtocol(nil))
continue
}
selected_metrics = append(selected_metrics, met)
}
}
cclog.ComponentDebugf("MetricAggregator", "Collected %d values from %d metrics", len(values), len(selected_metrics))
aggr_vars["values"] = values
aggr_vars["messages"] = selected_metrics
if len(values) > 0 && len(selected_metrics) > 0 {
res, err := expr.Run(aggr.exprFunc, aggr_vars)
if err == nil {
tags := copy_tags(aggr.Tags, selected_metrics)
meta := copy_meta(aggr.Meta, selected_metrics)
msg, err := lp.NewMetric(aggr.Name, tags, meta, res, time.Now())
if err == nil {
cclog.ComponentDebugf("MetricAggregator", "Sending %s", msg.ToLineProtocol(nil))
select {
case m.output <- msg:
default:
}
}
} else {
cclog.ComponentErrorf("MetricAggregator", "Failed to calculate aggregation with name %s: %s", aggr.Name, err.Error())
}
}
}
}
func (c *metricAggregatorExpr) AddConstant(name string, value any) {
c.constants[name] = value
}
func (c *metricAggregatorExpr) DelConstant(name string) {
delete(c.constants, name)
}
func (c *metricAggregatorExpr) DeleteAggregation(name string) error {
i := slices.IndexFunc(
c.aggregations,
func(agg MetricAggregatorExprIntervalConfig) bool {
return agg.Name == name
})
if i == -1 {
return fmt.Errorf("no aggregation for metric name %s", name)
}
copy(c.aggregations[i:], c.aggregations[i+1:])
c.aggregations = c.aggregations[:len(c.aggregations)-1]
return nil
}
func NewAggregatorExpr(output chan lp.CCMessage) (MetricAggregator, error) {
a := new(metricAggregatorExpr)
err := a.Init(output)
if err != nil {
return nil, err
}
return a, err
}
var expr_cached map[string]*vm.Program = make(map[string]*vm.Program)
var expr_cached_lock sync.Mutex
var baseenv_message = map[string]any{
"name": "",
"messagetype": "unknown",
"msgtype": "unknown",
"tag": map[string]any{
"type": "unknown",
"typeid": "0",
"stype": "unknown",
"stypeid": "0",
"hostname": "localhost",
"cluster": "nocluster",
},
"tags": map[string]any{
"type": "unknown",
"typeid": "0",
"stype": "unknown",
"stypeid": "0",
"hostname": "localhost",
"cluster": "nocluster",
},
"meta": map[string]any{
"unit": "invalid",
"source": "unknown",
},
"fields": map[string]any{
"value": 0,
"event": "",
"control": "",
"log": "",
},
"field": map[string]any{
"value": 0,
"event": "",
"control": "",
"log": "",
},
"timestamp": 1234567890,
"msg": lp.EmptyMessage(),
"message": lp.EmptyMessage(),
}
func EvalBoolConditionExpr(condition string, msg lp.CCMessage) (bool, error) {
scond := sanitizeExprString(condition)
expr_cached_lock.Lock()
evaluable, ok := expr_cached[scond]
expr_cached_lock.Unlock()
if !ok {
newcond := strings.ReplaceAll(
strings.ReplaceAll(
scond, "'", "\""), "%", "\\")
var err error
evaluable, err = expr.Compile(newcond, expr.Env(baseenv_message), expr.AsBool())
if err != nil {
return false, err
}
expr_cached_lock.Lock()
expr_cached[scond] = evaluable
expr_cached_lock.Unlock()
}
vars := GetParamMap(msg)
res, err := expr.Run(evaluable, vars)
if err != nil {
return false, err
}
paramMapPool.Put(vars)
return res.(bool), nil
}
func EvalFloat64ConditionExpr(condition string, msg lp.CCMessage) (float64, error) {
scond := sanitizeExprString(condition)
expr_cached_lock.Lock()
evaluable, ok := expr_cached[scond]
expr_cached_lock.Unlock()
if !ok {
newcond := strings.ReplaceAll(
strings.ReplaceAll(
scond, "'", "\""), "%", "\\")
var err error
evaluable, err = expr.Compile(newcond, expr.Env(baseenv_message), expr.AsFloat64())
if err != nil {
return math.NaN(), err
}
expr_cached_lock.Lock()
expr_cached[scond] = evaluable
expr_cached_lock.Unlock()
}
vars := GetParamMap(msg)
res, err := expr.Run(evaluable, vars)
paramMapPool.Put(vars)
return res.(float64), err
}
func EvalFloat64Expression(expression string, values map[string]any) (float64, error) {
sexpr := sanitizeExprString(expression)
expr_cached_lock.Lock()
evaluable, ok := expr_cached[sexpr]
expr_cached_lock.Unlock()
if !ok {
newcond := strings.ReplaceAll(
strings.ReplaceAll(
sexpr, "'", "\""), "%", "\\")
var err error
evaluable, err = expr.Compile(newcond, expr.Env(baseenv_message), expr.AsFloat64())
if err != nil {
return math.NaN(), err
}
expr_cached_lock.Lock()
expr_cached[sexpr] = evaluable
expr_cached_lock.Unlock()
}
res, err := expr.Run(evaluable, values)
return res.(float64), err
}
@@ -0,0 +1,97 @@
// Copyright (C) NHR@FAU, University Erlangen-Nuremberg.
// All rights reserved. This file is part of cc-lib.
// Use of this source code is governed by a MIT-style
// license that can be found in the LICENSE file.
// additional authors:
// Holger Obermaier (NHR@KIT)
package metricAggregator
import (
"math"
"testing"
"time"
lp "github.com/ClusterCockpit/cc-lib/v2/ccMessage"
)
func GenMetricNoCheck(value any, tags, meta map[string]string) lp.CCMessage {
msg, _ := lp.NewMetric("test", tags, meta, value, time.Now())
return msg
}
type TestBoolConfig struct {
msg lp.CCMessage
cond string
expected_result bool
should_fail bool
}
var testBoolConfig []TestBoolConfig = []TestBoolConfig{
{
msg: GenMetricNoCheck(1.0, nil, nil),
cond: "fields.value == 1",
expected_result: true,
should_fail: false,
},
{
msg: GenMetricNoCheck(1.0, map[string]string{"hostname": "testhost"}, nil),
cond: "tags.hostname == 'testhost'",
expected_result: true,
should_fail: false,
},
}
func TestEvalBoolConditionExprSimple(t *testing.T) {
for _, test := range testBoolConfig {
res, err := EvalBoolConditionExpr(test.cond, test.msg)
if err != nil && !test.should_fail {
t.Error(err.Error())
return
}
if !test.should_fail && res != test.expected_result {
t.Errorf("Condition '%s' evaluated to %v despite expecting %v", test.cond, res, test.expected_result)
return
}
}
}
type TestFloat64Config struct {
msg lp.CCMessage
cond string
expected_result float64
should_fail bool
}
var testFloat64Config []TestFloat64Config = []TestFloat64Config{
{
msg: GenMetricNoCheck(1.0, nil, nil),
cond: "fields.value + 3.14",
expected_result: 4.14,
should_fail: false,
},
{
msg: GenMetricNoCheck(2.0, nil, nil),
cond: "fields.value * 2",
expected_result: 4.0,
should_fail: false,
},
}
const compareFloat64Max = 1e-9
func TestEvalFloat64ConditionExprSimple(t *testing.T) {
for _, test := range testFloat64Config {
res, err := EvalFloat64ConditionExpr(test.cond, test.msg)
if err != nil && !test.should_fail {
t.Error(err.Error())
return
}
if !test.should_fail && math.Abs(res-test.expected_result) > compareFloat64Max {
t.Errorf("Condition '%s' evaluated to %f despite expecting %f", test.cond, res, test.expected_result)
return
}
}
}
@@ -336,14 +336,14 @@ func getCpuListOfDieFunc(args any) (any, error) {
} }
// wrapper function to get a list of all cpuids of the node // wrapper function to get a list of all cpuids of the node
func getCpuListOfNode() (any, error) { func getCpuListOfNodeFunc() (any, error) {
return topo.HwthreadList(), nil return topo.HwthreadList(), nil
} }
// helper function to get the cpuid list for a CCMetric type tag set (type and type-id) // helper function to get the cpuid list for a CCMetric type tag set (type and type-id)
// since there is no access to the metric data in the function, is should be called like // since there is no access to the metric data in the function, is should be called like
// `getCpuListOfType()` // `getCpuListOfType()`
func getCpuListOfType(args ...any) (any, error) { func getCpuListOfTypeFunc(args ...any) (any, error) {
cpulist := make([]int, 0) cpulist := make([]int, 0)
switch typ := args[0].(type) { switch typ := args[0].(type) {
case string: case string:
+115 -93
View File
@@ -9,6 +9,8 @@ package metricRouter
import ( import (
"fmt" "fmt"
"math"
"strings"
"sync" "sync"
"time" "time"
@@ -19,57 +21,107 @@ import (
mct "github.com/ClusterCockpit/cc-metric-collector/pkg/multiChanTicker" mct "github.com/ClusterCockpit/cc-metric-collector/pkg/multiChanTicker"
) )
type metricCachePeriod struct { type ccCache struct {
startstamp time.Time periodIdx int
stopstamp time.Time maxPeriods int
numMetrics int periods [][]lp.CCMessage
sizeMetrics int periodTimes []struct {
metrics []lp.CCMessage starttime time.Time
endtime time.Time
}
}
type CCCache interface {
Init(numPeriods int) error
Add(msg lp.CCMessage) error
GetPeriod(offset int) (time.Time, time.Time, []lp.CCMessage)
GetAll() []lp.CCMessage
NewPeriod()
}
func (c *ccCache) Init(numPeriods int) error {
c.maxPeriods = numPeriods
c.periodIdx = 0
c.periods = make([][]lp.CCMessage, c.maxPeriods)
c.periodTimes = make([]struct {
starttime time.Time
endtime time.Time
}, c.maxPeriods)
return nil
}
func (c *ccCache) NewPeriod() {
c.periodTimes[c.periodIdx].endtime = time.Now()
c.periodIdx = (c.periodIdx + 1) % c.maxPeriods
fmt.Printf("New period index %d\n", c.periodIdx)
c.periods[c.periodIdx] = c.periods[c.periodIdx][:0]
c.periodTimes[c.periodIdx].starttime = time.Now()
c.periodTimes[c.periodIdx].endtime = c.periodTimes[c.periodIdx].starttime
}
func (c *ccCache) Add(msg lp.CCMessage) error {
c.periods[c.periodIdx] = append(c.periods[c.periodIdx], msg)
c.periodTimes[c.periodIdx].endtime = msg.Time()
return nil
}
func (c *ccCache) GetPeriod(offset int) (time.Time, time.Time, []lp.CCMessage) {
if offset > c.maxPeriods {
offset = offset % c.maxPeriods
}
poff := int(math.Abs(float64(c.periodIdx - offset)))
out := make([]lp.CCMessage, 0, len(c.periods[poff%c.maxPeriods]))
out = append(out, c.periods[poff%c.maxPeriods]...)
return c.periodTimes[poff%c.maxPeriods].starttime, c.periodTimes[poff%c.maxPeriods].endtime, out
}
func (c *ccCache) GetAll() []lp.CCMessage {
out := make([]lp.CCMessage, 0)
for _, data := range c.periods {
out = append(out, data...)
}
return out
} }
// Metric cache data structure // Metric cache data structure
type metricCache struct { type metricCache struct {
numPeriods int cache CCCache
curPeriod int
lock sync.Mutex
intervals []*metricCachePeriod
wg *sync.WaitGroup wg *sync.WaitGroup
ticker mct.MultiChanTicker ticker mct.MultiChanTicker
tickchan chan time.Time tickchan chan time.Time
done chan bool done chan bool
output chan lp.CCMessage output chan lp.CCMessage
aggEngine agg.MetricAggregator aggEngine agg.MetricAggregator
numPeriods int
started bool
} }
type MetricCache interface { type MetricCache interface {
Init(output chan lp.CCMessage, ticker mct.MultiChanTicker, wg *sync.WaitGroup, numPeriods int) error Init(output chan lp.CCMessage, ticker mct.MultiChanTicker, wg *sync.WaitGroup, interval time.Duration, numPeriods int) error
Start() Start()
Add(metric lp.CCMessage) Add(metric lp.CCMessage)
GetPeriod(index int) (time.Time, time.Time, []lp.CCMessage)
AddAggregation(name, function, condition string, tags, meta map[string]string) error AddAggregation(name, function, condition string, tags, meta map[string]string) error
DeleteAggregation(name string) error DeleteAggregation(name string) error
Close() Close()
} }
func (c *metricCache) Init(output chan lp.CCMessage, ticker mct.MultiChanTicker, wg *sync.WaitGroup, numPeriods int) error { func (c *metricCache) Init(output chan lp.CCMessage, ticker mct.MultiChanTicker, wg *sync.WaitGroup, interval time.Duration, numPeriods int) error {
var err error var err error
c.done = make(chan bool) c.done = make(chan bool)
c.wg = wg c.wg = wg
c.ticker = ticker c.ticker = ticker
c.numPeriods = numPeriods c.numPeriods = numPeriods
c.started = false
c.cache = new(ccCache)
c.output = output c.output = output
c.intervals = make([]*metricCachePeriod, 0)
for i := 0; i < c.numPeriods+1; i++ { err = c.cache.Init(numPeriods)
p := new(metricCachePeriod) if err != nil {
p.numMetrics = 0 return fmt.Errorf("MetricCache: failed to create cache: %w", err)
p.sizeMetrics = 0
p.metrics = make([]lp.CCMessage, 0)
c.intervals = append(c.intervals, p)
} }
// Create a new aggregation engine. No separate goroutine at the moment c.aggEngine, err = agg.NewAggregatorExpr(c.output)
// The code is executed by the MetricCache goroutine
c.aggEngine, err = agg.NewAggregator(c.output)
if err != nil { if err != nil {
return fmt.Errorf("MetricCache: failed to create aggregator: %w", err) return fmt.Errorf("MetricCache: failed to create aggregator: %w", err)
} }
@@ -81,68 +133,58 @@ func (c *metricCache) Init(output chan lp.CCMessage, ticker mct.MultiChanTicker,
func (c *metricCache) Start() { func (c *metricCache) Start() {
c.tickchan = make(chan time.Time) c.tickchan = make(chan time.Time)
c.ticker.AddChannel(c.tickchan) c.ticker.AddChannel(c.tickchan)
// Router cache is done
done := func() {
cclog.ComponentDebug("MetricCache", "DONE")
close(c.done)
}
// Rotate cache interval c.wg.Add(1)
rotate := func(timestamp time.Time) int { go func() {
oldPeriod := c.curPeriod
c.curPeriod = oldPeriod + 1
if c.curPeriod >= c.numPeriods {
c.curPeriod = 0
}
c.intervals[oldPeriod].numMetrics = 0
c.intervals[oldPeriod].stopstamp = timestamp
c.intervals[c.curPeriod].startstamp = timestamp
c.intervals[c.curPeriod].stopstamp = timestamp
return oldPeriod
}
c.wg.Go(func() {
for { for {
select { select {
case <-c.done: case <-c.done:
done() c.wg.Done()
close(c.done)
cclog.ComponentDebug("MetricCache", "DONE")
return return
case tick := <-c.tickchan: case tick := <-c.tickchan:
c.lock.Lock() cclog.ComponentDebug("MetricCache", "Tick", tick)
old := rotate(tick) allmetrics := c.cache.GetAll()
// Get the last period and evaluate aggregation metrics c.cache.NewPeriod()
starttime, endtime, metrics := c.GetPeriod(old) mintime := tick
c.lock.Unlock() maxtime := mintime.AddDate(-1, 0, 0)
if len(metrics) > 0 {
c.aggEngine.Eval(starttime, endtime, metrics) for _, metric := range allmetrics {
if metric.Time().Before(mintime) {
mintime = metric.Time()
}
if metric.Time().After(maxtime) {
maxtime = metric.Time()
}
}
if len(allmetrics) > 0 {
cclog.ComponentDebugf("MetricCache", "Evaluate %d metrics from %v to %v", len(allmetrics), mintime.UnixNano(), maxtime.UnixNano())
c.wg.Go(func() {
c.aggEngine.Eval(mintime, maxtime, allmetrics)
})
} else { } else {
// This message is also printed in the first interval after startup // This message is also printed in the first interval after startup
cclog.ComponentDebug("MetricCache", "EMPTY INTERVAL?") cclog.ComponentDebug("MetricCache", "EMPTY INTERVAL?")
} }
} }
} }
}) }()
cclog.ComponentDebug("MetricCache", "START")
cclog.ComponentDebug("MetricCache", "STARTED")
} }
// Add a metric to the cache. The interval is defined by the global timer (rotate() in Start()) // Add a metric to the cache. The interval is defined by the global timer (rotate() in Start())
// The intervals list is used as round-robin buffer and the metric list grows dynamically and // The intervals list is used as round-robin buffer and the metric list grows dynamically and
// to avoid reallocations // to avoid reallocations
func (c *metricCache) Add(metric lp.CCMessage) { func (c *metricCache) Add(metric lp.CCMessage) {
if c.curPeriod >= 0 && c.curPeriod < c.numPeriods { err := c.cache.Add(metric)
c.lock.Lock() if err != nil {
p := c.intervals[c.curPeriod] s := metric.ToLineProtocol(nil)
if p.numMetrics < p.sizeMetrics { s = strings.TrimSpace(s)
p.metrics[p.numMetrics] = metric cclog.ComponentErrorf("MetricCache", "Failed to add metric %s", s)
p.numMetrics++
p.stopstamp = metric.Time()
} else {
p.metrics = append(p.metrics, metric)
p.numMetrics++
p.sizeMetrics++
p.stopstamp = metric.Time()
}
c.lock.Unlock()
} }
} }
@@ -154,40 +196,20 @@ func (c *metricCache) DeleteAggregation(name string) error {
return c.aggEngine.DeleteAggregation(name) return c.aggEngine.DeleteAggregation(name)
} }
// Get all metrics of a interval. The index is the difference to the current interval, so index=0
// is the current one, index=1 the last interval and so on. Returns and empty array if a wrong index
// is given (negative index, index larger than configured number of total intervals, ...)
func (c *metricCache) GetPeriod(index int) (time.Time, time.Time, []lp.CCMessage) {
start := time.Now()
stop := time.Now()
var metrics []lp.CCMessage
if index >= 0 && index < c.numPeriods {
pindex := c.curPeriod - index
if pindex < 0 {
pindex = c.numPeriods - pindex
}
if pindex >= 0 && pindex < c.numPeriods {
start = c.intervals[pindex].startstamp
stop = c.intervals[pindex].stopstamp
metrics = c.intervals[pindex].metrics
} else {
metrics = make([]lp.CCMessage, 0)
}
} else {
metrics = make([]lp.CCMessage, 0)
}
return start, stop, metrics
}
// Close finishes / stops the metric cache // Close finishes / stops the metric cache
func (c *metricCache) Close() { func (c *metricCache) Close() {
cclog.ComponentDebug("MetricCache", "CLOSE") cclog.ComponentDebug("MetricCache", "CLOSE")
if c.started {
c.done <- true c.done <- true
c.wg.Wait()
}
cclog.ComponentDebug("MetricCache", "CLOSED")
} }
func NewCache(output chan lp.CCMessage, ticker mct.MultiChanTicker, wg *sync.WaitGroup, numPeriods int) (MetricCache, error) { func NewCache(output chan lp.CCMessage, ticker mct.MultiChanTicker, wg *sync.WaitGroup, numPeriods int) (MetricCache, error) {
c := new(metricCache) c := new(metricCache)
err := c.Init(output, ticker, wg, numPeriods) err := c.Init(output, ticker, wg, ticker.GetDuration(), numPeriods)
if err != nil { if err != nil {
return nil, err return nil, err
} }
+69
View File
@@ -0,0 +1,69 @@
// Copyright (C) NHR@FAU, University Erlangen-Nuremberg.
// All rights reserved. This file is part of cc-lib.
// Use of this source code is governed by a MIT-style
// license that can be found in the LICENSE file.
// additional authors:
// Holger Obermaier (NHR@KIT)
package metricRouter
import (
"fmt"
"strings"
"sync"
"testing"
"time"
lp "github.com/ClusterCockpit/cc-lib/v2/ccMessage"
mct "github.com/ClusterCockpit/cc-metric-collector/pkg/multiChanTicker"
)
func TestCache(t *testing.T) {
output := make(chan lp.CCMessage, 2000)
var wg sync.WaitGroup
tickTime := time.Second
testChan := make(chan time.Time)
ticker := mct.NewTicker(tickTime)
ticker.AddChannel(testChan)
maxIntervals := 100
c, err := NewCache(output, ticker, &wg, 4)
if err != nil {
t.Errorf("failed to create new cache: %s", err.Error())
return
}
c.Start()
defer c.Close()
err = c.AddAggregation("ps_input_power", "Avg(values)", "name == 'ps1_input_power' || name == 'ps2_input_power' || name == 'ps3_input_power'", map[string]string{"hostname": "<copy>", "type": "<copy>"}, map[string]string{"unit": "<copy>"})
if err != nil {
t.Errorf("failed to add aggregation for %s", "ps_input_power")
return
}
for range maxIntervals {
timestamp := <-testChan
raw_metrics := []string{
fmt.Sprintf("ps1_input_power,type=node,hostname=myhost,unit=W value=696.0 %d", timestamp.UnixNano()),
fmt.Sprintf("ps2_input_power,type=node,hostname=myhost,unit=W value=732.0 %d", timestamp.UnixNano()),
fmt.Sprintf("ps3_input_power,type=node,hostname=myhost,unit=W value=720.0 %d", timestamp.UnixNano()),
fmt.Sprintf("cpu_load,type=hwthread,type-id=0,hostname=myhost value=45 %d", timestamp.UnixNano()),
}
metrics, err := lp.FromBytes([]byte(strings.Join(raw_metrics, "\n")))
if err != nil {
t.Errorf("failed to generate metrics: %s", err.Error())
}
for _, m := range metrics {
c.Add(m)
}
if len(output) > 0 {
for i := 0; i < len(output); i++ {
m := <-output
t.Log(m.ToLineProtocol(nil))
}
}
}
}
+2 -2
View File
@@ -38,7 +38,7 @@ type metricRouterConfig struct {
HostnameTagName string `json:"hostname_tag"` // Key name used when adding the hostname to a metric (default 'hostname') HostnameTagName string `json:"hostname_tag"` // Key name used when adding the hostname to a metric (default 'hostname')
AddTags []metricRouterTagConfig `json:"add_tags"` // List of tags that are added when the condition is met AddTags []metricRouterTagConfig `json:"add_tags"` // List of tags that are added when the condition is met
DelTags []metricRouterTagConfig `json:"delete_tags"` // List of tags that are removed when the condition is met DelTags []metricRouterTagConfig `json:"delete_tags"` // List of tags that are removed when the condition is met
IntervalAgg []agg.MetricAggregatorIntervalConfig `json:"interval_aggregates"` // List of aggregation function processed at the end of an interval IntervalAgg []agg.MetricAggregatorExprIntervalConfig `json:"interval_aggregates"` // List of aggregation function processed at the end of an interval
DropMetrics []string `json:"drop_metrics"` // List of metric names to drop. For fine-grained dropping use drop_metrics_if DropMetrics []string `json:"drop_metrics"` // List of metric names to drop. For fine-grained dropping use drop_metrics_if
DropMetricsIf []string `json:"drop_metrics_if"` // List of evaluatable terms to drop metrics DropMetricsIf []string `json:"drop_metrics_if"` // List of evaluatable terms to drop metrics
RenameMetrics map[string]string `json:"rename_metrics"` // Map to rename metric name from key to value RenameMetrics map[string]string `json:"rename_metrics"` // Map to rename metric name from key to value
@@ -253,7 +253,7 @@ func (r *metricRouter) Start() {
} }
// even if the metric is dropped, it is stored in the cache for // even if the metric is dropped, it is stored in the cache for
// aggregations // aggregations
if r.config.NumCacheIntervals > 0 { if r.config.NumCacheIntervals > 0 && m != nil {
r.cache.Add(m) r.cache.Add(m)
} }
} }
+7
View File
@@ -17,17 +17,20 @@ type multiChanTicker struct {
ticker *time.Ticker ticker *time.Ticker
channels []chan time.Time channels []chan time.Time
done chan bool done chan bool
duration time.Duration
} }
type MultiChanTicker interface { type MultiChanTicker interface {
Init(duration time.Duration) Init(duration time.Duration)
AddChannel(channel chan time.Time) AddChannel(channel chan time.Time)
GetDuration() time.Duration
Close() Close()
} }
func (t *multiChanTicker) Init(duration time.Duration) { func (t *multiChanTicker) Init(duration time.Duration) {
t.ticker = time.NewTicker(duration) t.ticker = time.NewTicker(duration)
t.done = make(chan bool) t.done = make(chan bool)
t.duration = duration
go func() { go func() {
done := func() { done := func() {
close(t.done) close(t.done)
@@ -53,6 +56,10 @@ func (t *multiChanTicker) Init(duration time.Duration) {
}() }()
} }
func (t *multiChanTicker) GetDuration() time.Duration {
return t.duration
}
func (t *multiChanTicker) AddChannel(channel chan time.Time) { func (t *multiChanTicker) AddChannel(channel chan time.Time) {
t.channels = append(t.channels, channel) t.channels = append(t.channels, channel)
} }