Compare commits

..
Author SHA1 Message Date
Michael Panzlaff 39cae9c79a Implement warmup delay for LIKWID collector 2026-08-25 12:28:44 +02:00
10 changed files with 151 additions and 757 deletions
+38 -1
View File
@@ -46,6 +46,7 @@ const (
LIKWID_LIB_DL_FLAGS = dl.RTLD_LAZY | dl.RTLD_GLOBAL LIKWID_LIB_DL_FLAGS = dl.RTLD_LAZY | dl.RTLD_GLOBAL
LIKWID_DEF_ACCESSMODE = "direct" LIKWID_DEF_ACCESSMODE = "direct"
LIKWID_DEF_LOCKFILE = "/var/run/likwid.lock" LIKWID_DEF_LOCKFILE = "/var/run/likwid.lock"
LIKWID_DEF_WARMUP_DEL = "0s"
) )
type LikwidCollectorMetricConfig struct { type LikwidCollectorMetricConfig struct {
@@ -83,6 +84,7 @@ type LikwidCollectorConfig struct {
DaemonPath string `json:"accessdaemon_path,omitempty"` DaemonPath string `json:"accessdaemon_path,omitempty"`
LibraryPath string `json:"liblikwid_path,omitempty"` LibraryPath string `json:"liblikwid_path,omitempty"`
LockfilePath string `json:"lockfile_path,omitempty"` LockfilePath string `json:"lockfile_path,omitempty"`
WarmupDelay string `json:"warmup_delay"`
} }
type LikwidCollector struct { type LikwidCollector struct {
@@ -104,6 +106,7 @@ type LikwidCollector struct {
likwidGroups map[C.int]LikwidEventsetConfig likwidGroups map[C.int]LikwidEventsetConfig
lock sync.Mutex lock sync.Mutex
measureThread thread.Thread measureThread thread.Thread
warmupDelay time.Duration
} }
type LikwidMetric struct { type LikwidMetric struct {
@@ -206,6 +209,7 @@ func (m *LikwidCollector) Init(config json.RawMessage) error {
m.config.AccessMode = LIKWID_DEF_ACCESSMODE m.config.AccessMode = LIKWID_DEF_ACCESSMODE
m.config.LibraryPath = LIKWID_LIB_NAME m.config.LibraryPath = LIKWID_LIB_NAME
m.config.LockfilePath = LIKWID_DEF_LOCKFILE m.config.LockfilePath = LIKWID_DEF_LOCKFILE
m.config.WarmupDelay = LIKWID_DEF_WARMUP_DEL
if len(config) > 0 { if len(config) > 0 {
d := json.NewDecoder(bytes.NewReader(config)) d := json.NewDecoder(bytes.NewReader(config))
d.DisallowUnknownFields() d.DisallowUnknownFields()
@@ -242,6 +246,17 @@ func (m *LikwidCollector) Init(config json.RawMessage) error {
m.cpu2tid[c] = i m.cpu2tid[c] = i
} }
if len(m.config.WarmupDelay) > 0 {
t, err := time.ParseDuration(m.config.WarmupDelay)
if err != nil {
return fmt.Errorf("%s Init(): Cannot parse WarmupDelay: %w", m.name, err)
}
if t < 0 {
return fmt.Errorf("%s Init(): WarmupDelay must not be negative", m.name)
}
m.warmupDelay = t
}
m.likwidGroups = make(map[C.int]LikwidEventsetConfig) m.likwidGroups = make(map[C.int]LikwidEventsetConfig)
// This is for the global metrics computation test // This is for the global metrics computation test
@@ -495,6 +510,28 @@ func (m *LikwidCollector) takeMeasurement(evidx int, evset LikwidEventsetConfig,
if ret != 0 { if ret != 0 {
return true, fmt.Errorf("failed to start events '%s', error %d", evset.go_estr, ret) return true, fmt.Errorf("failed to start events '%s', error %d", evset.go_estr, ret)
} }
// warmup measuring
if m.warmupDelay > 0 {
fmt.Printf("Performing warmup measurement for %s\n", m.warmupDelay)
select {
case <-sigchan:
ret = -1
case e := <-watcher.Events:
if e.Op != fsnotify.Chmod {
ret = C.perfmon_readCounters()
}
default:
ret = C.perfmon_readCounters()
}
if ret != 0 {
return true, fmt.Errorf("failed to read events '%s', error %d", evset.go_estr, ret)
}
time.Sleep(m.warmupDelay)
}
// begin measuring
select { select {
case <-sigchan: case <-sigchan:
ret = -1 ret = -1
@@ -512,7 +549,7 @@ func (m *LikwidCollector) takeMeasurement(evidx int, evset LikwidEventsetConfig,
// Wait // Wait
time.Sleep(interval) time.Sleep(interval)
// Read counters // end measuring
select { select {
case <-sigchan: case <-sigchan:
ret = -1 ret = -1
+1
View File
@@ -63,6 +63,7 @@ Additional options:
- `accessdaemon_path`: Folder of the accessDaemon `likwid-accessD` (like `/usr/local/sbin`) - `accessdaemon_path`: Folder of the accessDaemon `likwid-accessD` (like `/usr/local/sbin`)
- `liblikwid_path`: Location of `liblikwid.so` including file name like `/usr/local/lib/liblikwid.so` - `liblikwid_path`: Location of `liblikwid.so` including file name like `/usr/local/lib/liblikwid.so`
- `lockfile_path`: Location of LIKWID's lock file if multiple tools should access the hardware counters. Default `/var/run/likwid.lock` - `lockfile_path`: Location of LIKWID's lock file if multiple tools should access the hardware counters. Default `/var/run/likwid.lock`
- `warmup_delay`: Run an additional measurement of the specified length before the actual measurement. This can be used as a workaround for CPU starvation (during high load) in the measurement thread. Default is disabled (i.e. `0s`).
### Available metric types ### Available metric types
@@ -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", getCpuListOfNodeFunc), gval.Function("getCpuList", getCpuListOfNode),
gval.Function("getCpuListOfType", getCpuListOfTypeFunc), gval.Function("getCpuListOfType", getCpuListOfType),
) )
var language gval.Language = gval.NewLanguage( var language gval.Language = gval.NewLanguage(
@@ -1,449 +0,0 @@
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
}
@@ -1,97 +0,0 @@
// 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 getCpuListOfNodeFunc() (any, error) { func getCpuListOfNode() (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 getCpuListOfTypeFunc(args ...any) (any, error) { func getCpuListOfType(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:
+94 -116
View File
@@ -9,8 +9,6 @@ package metricRouter
import ( import (
"fmt" "fmt"
"math"
"strings"
"sync" "sync"
"time" "time"
@@ -21,107 +19,57 @@ import (
mct "github.com/ClusterCockpit/cc-metric-collector/pkg/multiChanTicker" mct "github.com/ClusterCockpit/cc-metric-collector/pkg/multiChanTicker"
) )
type ccCache struct { type metricCachePeriod struct {
periodIdx int startstamp time.Time
maxPeriods int stopstamp time.Time
periods [][]lp.CCMessage numMetrics int
periodTimes []struct { sizeMetrics int
starttime time.Time metrics []lp.CCMessage
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 {
cache CCCache numPeriods int
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, interval time.Duration, numPeriods int) error Init(output chan lp.CCMessage, ticker mct.MultiChanTicker, wg *sync.WaitGroup, 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, interval time.Duration, numPeriods int) error { func (c *metricCache) Init(output chan lp.CCMessage, ticker mct.MultiChanTicker, wg *sync.WaitGroup, 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)
err = c.cache.Init(numPeriods) for i := 0; i < c.numPeriods+1; i++ {
if err != nil { p := new(metricCachePeriod)
return fmt.Errorf("MetricCache: failed to create cache: %w", err) p.numMetrics = 0
p.sizeMetrics = 0
p.metrics = make([]lp.CCMessage, 0)
c.intervals = append(c.intervals, p)
} }
c.aggEngine, err = agg.NewAggregatorExpr(c.output) // 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)
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)
} }
@@ -133,58 +81,68 @@ 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)
}
c.wg.Add(1) // Rotate cache interval
go func() { 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() {
for { for {
select { select {
case <-c.done: case <-c.done:
c.wg.Done() done()
close(c.done)
cclog.ComponentDebug("MetricCache", "DONE")
return return
case tick := <-c.tickchan: case tick := <-c.tickchan:
cclog.ComponentDebug("MetricCache", "Tick", tick) c.lock.Lock()
allmetrics := c.cache.GetAll() old := rotate(tick)
c.cache.NewPeriod() // Get the last period and evaluate aggregation metrics
mintime := tick starttime, endtime, metrics := c.GetPeriod(old)
maxtime := mintime.AddDate(-1, 0, 0) c.lock.Unlock()
if len(metrics) > 0 {
for _, metric := range allmetrics { c.aggEngine.Eval(starttime, endtime, metrics)
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) {
err := c.cache.Add(metric) if c.curPeriod >= 0 && c.curPeriod < c.numPeriods {
if err != nil { c.lock.Lock()
s := metric.ToLineProtocol(nil) p := c.intervals[c.curPeriod]
s = strings.TrimSpace(s) if p.numMetrics < p.sizeMetrics {
cclog.ComponentErrorf("MetricCache", "Failed to add metric %s", s) 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()
} }
} }
@@ -196,20 +154,40 @@ 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, ticker.GetDuration(), numPeriods) err := c.Init(output, ticker, wg, numPeriods)
if err != nil { if err != nil {
return nil, err return nil, err
} }
-69
View File
@@ -1,69 +0,0 @@
// 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 // Metric router configuration
type metricRouterConfig struct { 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.MetricAggregatorExprIntervalConfig `json:"interval_aggregates"` // List of aggregation function processed at the end of an interval 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 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
IntervalStamp bool `json:"interval_timestamp"` // Update timestamp periodically by ticker each interval? 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 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 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 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 ChangeUnitPrefix map[string]string `json:"change_unit_prefix"` // Add prefix that should be applied to the metrics
MessageProcessor json.RawMessage `json:"process_messages,omitempty"` MessageProcessor json.RawMessage `json:"process_messages,omitempty"`
} }
// Metric router data structure // 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 // even if the metric is dropped, it is stored in the cache for
// aggregations // aggregations
if r.config.NumCacheIntervals > 0 && m != nil { if r.config.NumCacheIntervals > 0 {
r.cache.Add(m) r.cache.Add(m)
} }
} }
-7
View File
@@ -17,20 +17,17 @@ 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)
@@ -56,10 +53,6 @@ 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)
} }