From 4c1f810f5382b87d95b01296c146da24abd04bdc Mon Sep 17 00:00:00 2001 From: Thomas Gruber Date: Fri, 31 Jul 2026 16:31:04 +0200 Subject: [PATCH] Simplify cache, calculate with expr --- internal/metricAggregator/metricAggregator.go | 4 +- .../metricAggregator/metricAggregatorExpr.go | 449 ++++++++++++++++++ .../metricAggregatorExpr_test.go | 97 ++++ .../metricAggregatorFunctions.go | 4 +- internal/metricRouter/metricCache.go | 207 ++++---- internal/metricRouter/metricCache_test.go | 69 +++ internal/metricRouter/metricRouter.go | 28 +- 7 files changed, 744 insertions(+), 114 deletions(-) create mode 100644 internal/metricAggregator/metricAggregatorExpr.go create mode 100644 internal/metricAggregator/metricAggregatorExpr_test.go create mode 100644 internal/metricRouter/metricCache_test.go diff --git a/internal/metricAggregator/metricAggregator.go b/internal/metricAggregator/metricAggregator.go index 0637239..8bdf932 100644 --- a/internal/metricAggregator/metricAggregator.go +++ b/internal/metricAggregator/metricAggregator.go @@ -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( diff --git a/internal/metricAggregator/metricAggregatorExpr.go b/internal/metricAggregator/metricAggregatorExpr.go new file mode 100644 index 0000000..8dd7be6 --- /dev/null +++ b/internal/metricAggregator/metricAggregatorExpr.go @@ -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.Option = 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 "": + 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 "": + 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 +} diff --git a/internal/metricAggregator/metricAggregatorExpr_test.go b/internal/metricAggregator/metricAggregatorExpr_test.go new file mode 100644 index 0000000..234f742 --- /dev/null +++ b/internal/metricAggregator/metricAggregatorExpr_test.go @@ -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 + } + } +} diff --git a/internal/metricAggregator/metricAggregatorFunctions.go b/internal/metricAggregator/metricAggregatorFunctions.go index ed64cc8..14734b5 100644 --- a/internal/metricAggregator/metricAggregatorFunctions.go +++ b/internal/metricAggregator/metricAggregatorFunctions.go @@ -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: diff --git a/internal/metricRouter/metricCache.go b/internal/metricRouter/metricCache.go index 9f22bdf..8269754 100644 --- a/internal/metricRouter/metricCache.go +++ b/internal/metricRouter/metricCache.go @@ -9,6 +9,7 @@ package metricRouter import ( "fmt" + "math" "sync" "time" @@ -19,57 +20,104 @@ 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 + } + out := make([]lp.CCMessage, 0) + poff := int(math.Abs(float64(c.periodIdx - offset))) + 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) - } - // 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.cache.Init(numPeriods) + + c.aggEngine, err = agg.NewAggregatorExpr(c.output) if err != nil { return fmt.Errorf("MetricCache: failed to create aggregator: %w", err) } @@ -81,69 +129,56 @@ 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.Add(1) + go func() { + c.aggEngine.Eval(mintime, maxtime, allmetrics) + c.wg.Done() + }() } 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() - } + c.cache.Add(metric) } func (c *metricCache) AddAggregation(name, function, condition string, tags, meta map[string]string) error { @@ -154,40 +189,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 } diff --git a/internal/metricRouter/metricCache_test.go b/internal/metricRouter/metricCache_test.go new file mode 100644 index 0000000..e393537 --- /dev/null +++ b/internal/metricRouter/metricCache_test.go @@ -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": "", "type": ""}, map[string]string{"unit": ""}) + 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)) + } + } + + } + +} diff --git a/internal/metricRouter/metricRouter.go b/internal/metricRouter/metricRouter.go index 718b6aa..20ebaad 100644 --- a/internal/metricRouter/metricRouter.go +++ b/internal/metricRouter/metricRouter.go @@ -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) } }