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",
"router-file" : "router.json",
"main": {
"enable-markers": false,
"interval": "10s",
"duration": "1s"
}
+10 -7
View File
@@ -29,15 +29,17 @@ import (
)
type CentralConfigFile struct {
Interval string `json:"interval"`
Duration string `json:"duration"`
Interval string `json:"interval"`
Duration string `json:"duration"`
EnableMarkers bool `json:"enable-markers"`
}
type RuntimeConfig struct {
Interval time.Duration
Duration time.Duration
CliArgs map[string]string
ConfigFile CentralConfigFile
Interval time.Duration
Duration time.Duration
EnableMarkers bool
CliArgs map[string]string
ConfigFile CentralConfigFile
MetricRouter mr.MetricRouter
CollectManager collectors.CollectorManager
@@ -157,6 +159,7 @@ func mainFunc() int {
cclog.Error("The interval should be greater than duration")
return 1
}
rcfg.EnableMarkers = rcfg.ConfigFile.EnableMarkers
routerConf := ccconf.GetPackageConfig("router")
if len(routerConf) == 0 {
@@ -199,7 +202,7 @@ func mainFunc() int {
rcfg.MetricRouter.AddOutput(RouterToSinksChannel)
// 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 {
cclog.Error(err.Error())
return 1
+32 -14
View File
@@ -56,21 +56,22 @@ var AvailableCollectors = map[string]MetricCollector{
// Metric collector manager data structure
type collectorManager struct {
collectors []MetricCollector // List of metric collectors to read in parallel
serial []MetricCollector // List of metric collectors to read serially
output chan lp.CCMessage // Output channels
done chan bool // channel to finish / stop metric collector manager
ticker mct.MultiChanTicker // periodically ticking once each interval
duration time.Duration // duration (for metrics that measure over a given duration)
wg *sync.WaitGroup // wait group for all goroutines in cc-metric-collector
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
parallel_run bool // Flag whether the collectors are currently read in parallel
collectors []MetricCollector // List of metric collectors to read in parallel
serial []MetricCollector // List of metric collectors to read serially
output chan lp.CCMessage // Output channels
done chan bool // channel to finish / stop metric collector manager
ticker mct.MultiChanTicker // periodically ticking once each interval
duration time.Duration // duration (for metrics that measure over a given duration)
wg *sync.WaitGroup // wait group for all goroutines in cc-metric-collector
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
parallel_run bool // Flag whether the collectors are currently read in parallel
enableMarkers bool // Send ccmc-{begin,end} metrics
}
// Metric collector manager access functions
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)
Start()
Close()
@@ -83,7 +84,7 @@ type CollectorManager interface {
// * ticker (from variable ticker)
// * configuration (read from config file in variable collectConfigFile)
// 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.serial = make([]MetricCollector, 0)
cm.output = nil
@@ -91,6 +92,7 @@ func (cm *collectorManager) Init(ticker mct.MultiChanTicker, duration time.Durat
cm.wg = wg
cm.ticker = ticker
cm.duration = duration
cm.enableMarkers = enableMarkers
d := json.NewDecoder(bytes.NewReader(collectConfig))
d.DisallowUnknownFields()
@@ -148,6 +150,14 @@ func (cm *collectorManager) Start() {
done()
return
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
for _, c := range cm.collectors {
// Wait for done signal or execute the collector
@@ -179,6 +189,14 @@ func (cm *collectorManager) Start() {
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
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)
err := cm.Init(ticker, duration, wg, collectConfig)
err := cm.Init(ticker, duration, enableMarkers, wg, collectConfig)
if err != nil {
return nil, err
}
+1 -4
View File
@@ -65,7 +65,7 @@ func (m *MegwareEurekaCollector) Init(config json.RawMessage) error {
return nil
}
m.name = "MegwareEureka"
m.name = "MegwareEurekaCollector"
if err := m.setup(); err != nil {
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())
}
fmt.Printf("string: %+v\n", stdout.String())
fmt.Printf("obj: %+v\n", u20output)
u20output.GetMpsPollValues.Timestamp = time.Now()
return &u20output.GetMpsPollValues, nil
+1
View File
@@ -4,6 +4,7 @@
"receivers-file" : "./receivers.json",
"router-file" : "./router.json",
"main" : {
"enable-markers": false,
"interval": "10s",
"duration": "1s"
}
+2 -2
View File
@@ -5,12 +5,12 @@ go 1.25.0
require (
github.com/ClusterCockpit/cc-lib/v2 v2.12.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/fsnotify/fsnotify v1.10.1
github.com/tklauser/go-sysconf v0.4.0
golang.design/x/runtime v0.3.0
golang.org/x/sys v0.45.0
golang.org/x/sys v0.47.0
)
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/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.2-0 h1:7M4cFG62wSUHw8i0XSiNU7ejKODytTS6ZrW/vgB2NSI=
github.com/NVIDIA/go-nvml v0.13.2-0/go.mod h1:ahi2psRYoa+wYUBIrZPRO+wJs9lcvMhxSSkjjvsJJNQ=
github.com/NVIDIA/go-nvml v0.13.3-1 h1:P76U2h88OZSiMtdhRsJjSF5DXyXUqHIXKeDicVAaae0=
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/go.mod h1:XRFLwvmkTEdYziLdaCeCa5ImcGVrfQbeNUbVR+C6xac=
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/net v0.53.0 h1:d+qAbo5L0orcWAr0a9JweQpjXF19LMXJE8Ey7hwOdUA=
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.45.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs=
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/go.mod h1:Y4YMaQmXwGQZoFaVFk4YpCt4FLQMYKZe9oeV/f4MSno=
golang.org/x/tools v0.14.0 h1:jvNa2pY0M4r62jkRQ6RwEZZyPcymeL9XZMLBbV7U2nc=
@@ -69,8 +69,8 @@ var metricCacheLanguage = gval.NewLanguage(
gval.Function("getNumaCpuList", getCpuListOfNumaDomainFunc),
gval.Function("getDieCpuList", getCpuListOfDieFunc),
gval.Function("getCoreCpuList", getCpuListOfCoreFunc),
gval.Function("getCpuList", getCpuListOfNode),
gval.Function("getCpuListOfType", getCpuListOfType),
gval.Function("getCpuList", getCpuListOfNodeFunc),
gval.Function("getCpuListOfType", getCpuListOfTypeFunc),
)
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
func getCpuListOfNode() (any, error) {
func getCpuListOfNodeFunc() (any, error) {
return topo.HwthreadList(), nil
}
// 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
// `getCpuListOfType()`
func getCpuListOfType(args ...any) (any, error) {
func getCpuListOfTypeFunc(args ...any) (any, error) {
cpulist := make([]int, 0)
switch typ := args[0].(type) {
case string:
+116 -94
View File
@@ -9,6 +9,8 @@ package metricRouter
import (
"fmt"
"math"
"strings"
"sync"
"time"
@@ -19,57 +21,107 @@ import (
mct "github.com/ClusterCockpit/cc-metric-collector/pkg/multiChanTicker"
)
type metricCachePeriod struct {
startstamp time.Time
stopstamp time.Time
numMetrics int
sizeMetrics int
metrics []lp.CCMessage
type ccCache struct {
periodIdx int
maxPeriods int
periods [][]lp.CCMessage
periodTimes []struct {
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
type metricCache struct {
numPeriods int
curPeriod int
lock sync.Mutex
intervals []*metricCachePeriod
cache CCCache
wg *sync.WaitGroup
ticker mct.MultiChanTicker
tickchan chan time.Time
done chan bool
output chan lp.CCMessage
aggEngine agg.MetricAggregator
numPeriods int
started bool
}
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()
Add(metric lp.CCMessage)
GetPeriod(index int) (time.Time, time.Time, []lp.CCMessage)
AddAggregation(name, function, condition string, tags, meta map[string]string) error
DeleteAggregation(name string) error
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
c.done = make(chan bool)
c.wg = wg
c.ticker = ticker
c.numPeriods = numPeriods
c.started = false
c.cache = new(ccCache)
c.output = output
c.intervals = make([]*metricCachePeriod, 0)
for i := 0; i < c.numPeriods+1; i++ {
p := new(metricCachePeriod)
p.numMetrics = 0
p.sizeMetrics = 0
p.metrics = make([]lp.CCMessage, 0)
c.intervals = append(c.intervals, p)
err = c.cache.Init(numPeriods)
if err != nil {
return fmt.Errorf("MetricCache: failed to create cache: %w", err)
}
// Create a new aggregation engine. No separate goroutine at the moment
// The code is executed by the MetricCache goroutine
c.aggEngine, err = agg.NewAggregator(c.output)
c.aggEngine, err = agg.NewAggregatorExpr(c.output)
if err != nil {
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() {
c.tickchan = make(chan time.Time)
c.ticker.AddChannel(c.tickchan)
// Router cache is done
done := func() {
cclog.ComponentDebug("MetricCache", "DONE")
close(c.done)
}
// Rotate cache interval
rotate := func(timestamp time.Time) int {
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() {
c.wg.Add(1)
go func() {
for {
select {
case <-c.done:
done()
c.wg.Done()
close(c.done)
cclog.ComponentDebug("MetricCache", "DONE")
return
case tick := <-c.tickchan:
c.lock.Lock()
old := rotate(tick)
// Get the last period and evaluate aggregation metrics
starttime, endtime, metrics := c.GetPeriod(old)
c.lock.Unlock()
if len(metrics) > 0 {
c.aggEngine.Eval(starttime, endtime, metrics)
cclog.ComponentDebug("MetricCache", "Tick", tick)
allmetrics := c.cache.GetAll()
c.cache.NewPeriod()
mintime := tick
maxtime := mintime.AddDate(-1, 0, 0)
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 {
// This message is also printed in the first interval after startup
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())
// The intervals list is used as round-robin buffer and the metric list grows dynamically and
// to avoid reallocations
func (c *metricCache) Add(metric lp.CCMessage) {
if c.curPeriod >= 0 && c.curPeriod < c.numPeriods {
c.lock.Lock()
p := c.intervals[c.curPeriod]
if p.numMetrics < p.sizeMetrics {
p.metrics[p.numMetrics] = metric
p.numMetrics++
p.stopstamp = metric.Time()
} else {
p.metrics = append(p.metrics, metric)
p.numMetrics++
p.sizeMetrics++
p.stopstamp = metric.Time()
}
c.lock.Unlock()
err := c.cache.Add(metric)
if err != nil {
s := metric.ToLineProtocol(nil)
s = strings.TrimSpace(s)
cclog.ComponentErrorf("MetricCache", "Failed to add metric %s", s)
}
}
@@ -154,40 +196,20 @@ func (c *metricCache) DeleteAggregation(name string) error {
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
func (c *metricCache) Close() {
cclog.ComponentDebug("MetricCache", "CLOSE")
c.done <- true
if c.started {
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) {
c := new(metricCache)
err := c.Init(output, ticker, wg, numPeriods)
err := c.Init(output, ticker, wg, ticker.GetDuration(), numPeriods)
if err != nil {
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))
}
}
}
}
+14 -14
View File
@@ -35,19 +35,19 @@ type metricRouterTagConfig struct {
// Metric router configuration
type metricRouterConfig struct {
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
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
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
RenameMetrics map[string]string `json:"rename_metrics"` // Map to rename metric name from key to value
IntervalStamp bool `json:"interval_timestamp"` // Update timestamp periodically by ticker each interval?
NumCacheIntervals int `json:"num_cache_intervals"` // Number of intervals of cached metrics for evaluation
MaxForward int `json:"max_forward"` // Number of maximal forwarded metrics at one select
NormalizeUnits bool `json:"normalize_units"` // Check unit meta flag and normalize it using cc-units
ChangeUnitPrefix map[string]string `json:"change_unit_prefix"` // Add prefix that should be applied to the metrics
MessageProcessor json.RawMessage `json:"process_messages,omitempty"`
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
DelTags []metricRouterTagConfig `json:"delete_tags"` // List of tags that are removed when the condition is met
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
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
IntervalStamp bool `json:"interval_timestamp"` // Update timestamp periodically by ticker each interval?
NumCacheIntervals int `json:"num_cache_intervals"` // Number of intervals of cached metrics for evaluation
MaxForward int `json:"max_forward"` // Number of maximal forwarded metrics at one select
NormalizeUnits bool `json:"normalize_units"` // Check unit meta flag and normalize it using cc-units
ChangeUnitPrefix map[string]string `json:"change_unit_prefix"` // Add prefix that should be applied to the metrics
MessageProcessor json.RawMessage `json:"process_messages,omitempty"`
}
// Metric router data structure
@@ -253,7 +253,7 @@ func (r *metricRouter) Start() {
}
// even if the metric is dropped, it is stored in the cache for
// aggregations
if r.config.NumCacheIntervals > 0 {
if r.config.NumCacheIntervals > 0 && m != nil {
r.cache.Add(m)
}
}
+7
View File
@@ -17,17 +17,20 @@ type multiChanTicker struct {
ticker *time.Ticker
channels []chan time.Time
done chan bool
duration time.Duration
}
type MultiChanTicker interface {
Init(duration time.Duration)
AddChannel(channel chan time.Time)
GetDuration() time.Duration
Close()
}
func (t *multiChanTicker) Init(duration time.Duration) {
t.ticker = time.NewTicker(duration)
t.done = make(chan bool)
t.duration = duration
go func() {
done := func() {
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) {
t.channels = append(t.channels, channel)
}