Simplify cache, calculate with expr

This commit is contained in:
Thomas Gruber
2026-07-31 16:31:04 +02:00
parent 789ba7f4a5
commit 4c1f810f53
7 changed files with 744 additions and 114 deletions
@@ -69,8 +69,8 @@ var metricCacheLanguage = gval.NewLanguage(
gval.Function("getNumaCpuList", getCpuListOfNumaDomainFunc), gval.Function("getNumaCpuList", getCpuListOfNumaDomainFunc),
gval.Function("getDieCpuList", getCpuListOfDieFunc), gval.Function("getDieCpuList", getCpuListOfDieFunc),
gval.Function("getCoreCpuList", getCpuListOfCoreFunc), gval.Function("getCoreCpuList", getCpuListOfCoreFunc),
gval.Function("getCpuList", getCpuListOfNode), gval.Function("getCpuList", getCpuListOfNodeFunc),
gval.Function("getCpuListOfType", getCpuListOfType), gval.Function("getCpuListOfType", getCpuListOfTypeFunc),
) )
var language gval.Language = gval.NewLanguage( var language gval.Language = gval.NewLanguage(
@@ -0,0 +1,449 @@
package metricAggregator
import (
"fmt"
"maps"
"math"
"slices"
"strings"
"sync"
"time"
cclog "github.com/ClusterCockpit/cc-lib/v2/ccLogger"
lp "github.com/ClusterCockpit/cc-lib/v2/ccMessage"
"github.com/expr-lang/expr"
"github.com/expr-lang/expr/vm"
)
type MetricAggregatorExprIntervalConfig struct {
Name string `json:"name"` // Metric name for the new metric
Function string `json:"function"` // Function to apply on the metric
Condition string `json:"if"` // Condition for applying function
Tags map[string]string `json:"tags,omitempty"` // Tags for the new metric
Meta map[string]string `json:"meta,omitempty"` // Meta information for the new metric
ValueType string `json:"value_type,omitempty"`
IncludeNumCacheIntervals int `json:"expr_on_num_intervals,omitempty"`
exprCond *vm.Program
exprFunc *vm.Program
}
type metricAggregatorExpr struct {
constants map[string]any
output chan lp.CCMessage
aggregations []MetricAggregatorExprIntervalConfig
}
var paramMapPool = sync.Pool{
New: func() any {
return make(map[string]any)
},
}
func sanitizeExprString(key string) string {
return strings.ReplaceAll(key, "type-id", "typeid")
}
// AddAggregation(name, function, condition string, tags, meta map[string]string) error
// DeleteAggregation(name string) error
// Init(output chan lp.CCMessage) error
// Eval(starttime time.Time, endtime time.Time, metrics []lp.CCMessage)
func (m *metricAggregatorExpr) Init(output chan lp.CCMessage) error {
m.output = output
m.constants = make(map[string]any)
m.aggregations = make([]MetricAggregatorExprIntervalConfig, 0)
return nil
}
func GetParamMap(point lp.CCMessage) map[string]any {
params := paramMapPool.Get().(map[string]any)
clear(params)
// Put metric name into params map
params["name"] = point.Name()
// Put full message into params map
params["message"] = point
params["msg"] = point
// Put timestamp into params map
params["timestamp"] = point.Time().Unix()
params["time"] = params["timestamp"]
// Put fields into params map
fields := paramMapPool.Get().(map[string]any)
clear(fields)
for key, value := range point.Fields() {
fields[key] = value
switch key {
case "value":
params["messagetype"] = "metric"
params["value"] = value
params["metric"] = value
case "event":
params["messagetype"] = "event"
params["event"] = value
case "control":
params["messagetype"] = "control"
params["control"] = value
case "log":
params["messagetype"] = "log"
params["log"] = value
default:
params["messagetype"] = "unknown"
}
}
params["msgtype"] = params["messagetype"]
params["fields"] = fields
params["field"] = fields
// Put tags into params map
tags := paramMapPool.Get().(map[string]any)
clear(tags)
for key, value := range point.Tags() {
tags[sanitizeExprString(key)] = value
}
params["tags"] = tags
params["tag"] = tags
// Put meta information into params map
meta := paramMapPool.Get().(map[string]any)
clear(meta)
for key, value := range point.Meta() {
meta[sanitizeExprString(key)] = value
}
params["meta"] = meta
return params
}
var baseenv_multi_message = map[string]any{
"name": "",
"starttime": 1234567890,
"endtime": 1234567890,
"messages": make([]lp.CCMessage, 0),
"Median": medianfunc,
"Sum": sumfunc,
"Min": minfunc,
"Max": maxfunc,
"Mean": avgfunc,
"Avg": avgfunc,
"Match": matchfunc,
"getCpuCore": getCpuCoreFunc,
"getCpuSocket": getCpuSocketFunc,
"getCpuNumaDomain": getCpuNumaDomainFunc,
"getCpuDie": getCpuDieFunc,
"getCpuListOfCore": getCpuListOfCoreFunc,
"getCpuListOfSocket": getCpuListOfSocketFunc,
"getCpuListOfNumaDomain": getCpuListOfNumaDomainFunc,
"getCpuListOfDie": getCpuListOfDieFunc,
"getCpuListOfNode": getCpuListOfNodeFunc,
"getCpuListOfType": getCpuListOfTypeFunc,
}
func (m *metricAggregatorExpr) AddAggregationWithType(name, function, condition, valueType string, tags, meta map[string]string) error {
cond, err := expr.Compile(condition, expr.Env(baseenv_multi_message), expr.AsBool(), expr.AllowUndefinedVariables())
if err != nil {
err = fmt.Errorf("failed to compile condition for aggregation %s: %s", name, err.Error())
return err
}
var exprOption expr.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 "<copy>":
for _, m := range metrics {
v, err := m.GetTag(key)
if err {
out[key] = v
}
}
default:
out[key] = value
}
}
return out
}
copy_meta := func(meta map[string]string, metrics []lp.CCMessage) map[string]string {
out := make(map[string]string)
for key, value := range meta {
switch value {
case "<copy>":
for _, m := range metrics {
v, err := m.GetMeta(key)
if err {
out[key] = v
}
}
default:
out[key] = value
}
}
return out
}
for _, aggr := range m.aggregations {
selected_metrics := make([]lp.CCMessage, 0)
values := make([]float64, 0)
aggr_vars := make(map[string]any)
maps.Copy(aggr_vars, baseenv_multi_message)
maps.Copy(aggr_vars, m.constants)
aggr_vars["starttime"] = starttime
aggr_vars["endtime"] = endtime
for _, met := range metrics {
met_vars := make(map[string]any)
maps.Copy(met_vars, aggr_vars)
met_vars["message"] = met
maps.Copy(met_vars, GetParamMap(met))
res, err := expr.Run(aggr.exprCond, met_vars)
if err != nil || res == false {
continue
}
if value, ok := met.GetField("value"); ok {
switch v := value.(type) {
case float64:
values = append(values, v)
case float32:
case int:
case int8:
case int16:
case int32:
case int64:
case uint:
case uint8:
case uint16:
case uint32:
case uint64:
values = append(values, float64(v))
case bool:
if v {
values = append(values, float64(1))
} else {
values = append(values, float64(0))
}
default:
cclog.ComponentErrorf("MetricAggregator", "Cannot convert value type for %s", met.ToLineProtocol(nil))
continue
}
selected_metrics = append(selected_metrics, met)
}
}
cclog.ComponentDebugf("MetricAggregator", "Collected %d values from %d metrics", len(values), len(selected_metrics))
aggr_vars["values"] = values
aggr_vars["messages"] = selected_metrics
if len(values) > 0 && len(selected_metrics) > 0 {
res, err := expr.Run(aggr.exprFunc, aggr_vars)
if err == nil {
tags := copy_tags(aggr.Tags, selected_metrics)
meta := copy_meta(aggr.Meta, selected_metrics)
msg, err := lp.NewMetric(aggr.Name, tags, meta, res, time.Now())
if err == nil {
cclog.ComponentDebugf("MetricAggregator", "Sending %s", msg.ToLineProtocol(nil))
select {
case m.output <- msg:
default:
}
}
} else {
cclog.ComponentErrorf("MetricAggregator", "Failed to calculate aggregation with name %s: %s", aggr.Name, err.Error())
}
}
}
}
func (c *metricAggregatorExpr) AddConstant(name string, value any) {
c.constants[name] = value
}
func (c *metricAggregatorExpr) DelConstant(name string) {
delete(c.constants, name)
}
func (c *metricAggregatorExpr) DeleteAggregation(name string) error {
i := slices.IndexFunc(
c.aggregations,
func(agg MetricAggregatorExprIntervalConfig) bool {
return agg.Name == name
})
if i == -1 {
return fmt.Errorf("no aggregation for metric name %s", name)
}
copy(c.aggregations[i:], c.aggregations[i+1:])
c.aggregations = c.aggregations[:len(c.aggregations)-1]
return nil
}
func NewAggregatorExpr(output chan lp.CCMessage) (MetricAggregator, error) {
a := new(metricAggregatorExpr)
err := a.Init(output)
if err != nil {
return nil, err
}
return a, err
}
var expr_cached map[string]*vm.Program = make(map[string]*vm.Program)
var expr_cached_lock sync.Mutex
var baseenv_message = map[string]any{
"name": "",
"messagetype": "unknown",
"msgtype": "unknown",
"tag": map[string]any{
"type": "unknown",
"typeid": "0",
"stype": "unknown",
"stypeid": "0",
"hostname": "localhost",
"cluster": "nocluster",
},
"tags": map[string]any{
"type": "unknown",
"typeid": "0",
"stype": "unknown",
"stypeid": "0",
"hostname": "localhost",
"cluster": "nocluster",
},
"meta": map[string]any{
"unit": "invalid",
"source": "unknown",
},
"fields": map[string]any{
"value": 0,
"event": "",
"control": "",
"log": "",
},
"field": map[string]any{
"value": 0,
"event": "",
"control": "",
"log": "",
},
"timestamp": 1234567890,
"msg": lp.EmptyMessage(),
"message": lp.EmptyMessage(),
}
func EvalBoolConditionExpr(condition string, msg lp.CCMessage) (bool, error) {
scond := sanitizeExprString(condition)
expr_cached_lock.Lock()
evaluable, ok := expr_cached[scond]
expr_cached_lock.Unlock()
if !ok {
newcond := strings.ReplaceAll(
strings.ReplaceAll(
scond, "'", "\""), "%", "\\")
var err error
evaluable, err = expr.Compile(newcond, expr.Env(baseenv_message), expr.AsBool())
if err != nil {
return false, err
}
expr_cached_lock.Lock()
expr_cached[scond] = evaluable
expr_cached_lock.Unlock()
}
vars := GetParamMap(msg)
res, err := expr.Run(evaluable, vars)
if err != nil {
return false, err
}
paramMapPool.Put(vars)
return res.(bool), nil
}
func EvalFloat64ConditionExpr(condition string, msg lp.CCMessage) (float64, error) {
scond := sanitizeExprString(condition)
expr_cached_lock.Lock()
evaluable, ok := expr_cached[scond]
expr_cached_lock.Unlock()
if !ok {
newcond := strings.ReplaceAll(
strings.ReplaceAll(
scond, "'", "\""), "%", "\\")
var err error
evaluable, err = expr.Compile(newcond, expr.Env(baseenv_message), expr.AsFloat64())
if err != nil {
return math.NaN(), err
}
expr_cached_lock.Lock()
expr_cached[scond] = evaluable
expr_cached_lock.Unlock()
}
vars := GetParamMap(msg)
res, err := expr.Run(evaluable, vars)
paramMapPool.Put(vars)
return res.(float64), err
}
func EvalFloat64Expression(expression string, values map[string]any) (float64, error) {
sexpr := sanitizeExprString(expression)
expr_cached_lock.Lock()
evaluable, ok := expr_cached[sexpr]
expr_cached_lock.Unlock()
if !ok {
newcond := strings.ReplaceAll(
strings.ReplaceAll(
sexpr, "'", "\""), "%", "\\")
var err error
evaluable, err = expr.Compile(newcond, expr.Env(baseenv_message), expr.AsFloat64())
if err != nil {
return math.NaN(), err
}
expr_cached_lock.Lock()
expr_cached[sexpr] = evaluable
expr_cached_lock.Unlock()
}
res, err := expr.Run(evaluable, values)
return res.(float64), err
}
@@ -0,0 +1,97 @@
// Copyright (C) NHR@FAU, University Erlangen-Nuremberg.
// All rights reserved. This file is part of cc-lib.
// Use of this source code is governed by a MIT-style
// license that can be found in the LICENSE file.
// additional authors:
// Holger Obermaier (NHR@KIT)
package metricAggregator
import (
"math"
"testing"
"time"
lp "github.com/ClusterCockpit/cc-lib/v2/ccMessage"
)
func GenMetricNoCheck(value any, tags, meta map[string]string) lp.CCMessage {
msg, _ := lp.NewMetric("test", tags, meta, value, time.Now())
return msg
}
type TestBoolConfig struct {
msg lp.CCMessage
cond string
expected_result bool
should_fail bool
}
var testBoolConfig []TestBoolConfig = []TestBoolConfig{
{
msg: GenMetricNoCheck(1.0, nil, nil),
cond: "fields.value == 1",
expected_result: true,
should_fail: false,
},
{
msg: GenMetricNoCheck(1.0, map[string]string{"hostname": "testhost"}, nil),
cond: "tags.hostname == 'testhost'",
expected_result: true,
should_fail: false,
},
}
func TestEvalBoolConditionExprSimple(t *testing.T) {
for _, test := range testBoolConfig {
res, err := EvalBoolConditionExpr(test.cond, test.msg)
if err != nil && !test.should_fail {
t.Error(err.Error())
return
}
if !test.should_fail && res != test.expected_result {
t.Errorf("Condition '%s' evaluated to %v despite expecting %v", test.cond, res, test.expected_result)
return
}
}
}
type TestFloat64Config struct {
msg lp.CCMessage
cond string
expected_result float64
should_fail bool
}
var testFloat64Config []TestFloat64Config = []TestFloat64Config{
{
msg: GenMetricNoCheck(1.0, nil, nil),
cond: "fields.value + 3.14",
expected_result: 4.14,
should_fail: false,
},
{
msg: GenMetricNoCheck(2.0, nil, nil),
cond: "fields.value * 2",
expected_result: 4.0,
should_fail: false,
},
}
const compareFloat64Max = 1e-9
func TestEvalFloat64ConditionExprSimple(t *testing.T) {
for _, test := range testFloat64Config {
res, err := EvalFloat64ConditionExpr(test.cond, test.msg)
if err != nil && !test.should_fail {
t.Error(err.Error())
return
}
if !test.should_fail && math.Abs(res-test.expected_result) > compareFloat64Max {
t.Errorf("Condition '%s' evaluated to %f despite expecting %f", test.cond, res, test.expected_result)
return
}
}
}
@@ -336,14 +336,14 @@ func getCpuListOfDieFunc(args any) (any, error) {
} }
// wrapper function to get a list of all cpuids of the node // wrapper function to get a list of all cpuids of the node
func getCpuListOfNode() (any, error) { func getCpuListOfNodeFunc() (any, error) {
return topo.HwthreadList(), nil return topo.HwthreadList(), nil
} }
// helper function to get the cpuid list for a CCMetric type tag set (type and type-id) // helper function to get the cpuid list for a CCMetric type tag set (type and type-id)
// since there is no access to the metric data in the function, is should be called like // since there is no access to the metric data in the function, is should be called like
// `getCpuListOfType()` // `getCpuListOfType()`
func getCpuListOfType(args ...any) (any, error) { func getCpuListOfTypeFunc(args ...any) (any, error) {
cpulist := make([]int, 0) cpulist := make([]int, 0)
switch typ := args[0].(type) { switch typ := args[0].(type) {
case string: case string:
+110 -95
View File
@@ -9,6 +9,7 @@ package metricRouter
import ( import (
"fmt" "fmt"
"math"
"sync" "sync"
"time" "time"
@@ -19,57 +20,104 @@ import (
mct "github.com/ClusterCockpit/cc-metric-collector/pkg/multiChanTicker" mct "github.com/ClusterCockpit/cc-metric-collector/pkg/multiChanTicker"
) )
type metricCachePeriod struct { type ccCache struct {
startstamp time.Time periodIdx int
stopstamp time.Time maxPeriods int
numMetrics int periods [][]lp.CCMessage
sizeMetrics int periodTimes []struct {
metrics []lp.CCMessage starttime time.Time
endtime time.Time
}
}
type CCCache interface {
Init(numPeriods int) error
Add(msg lp.CCMessage) error
GetPeriod(offset int) (time.Time, time.Time, []lp.CCMessage)
GetAll() []lp.CCMessage
NewPeriod()
}
func (c *ccCache) Init(numPeriods int) error {
c.maxPeriods = numPeriods
c.periodIdx = 0
c.periods = make([][]lp.CCMessage, c.maxPeriods)
c.periodTimes = make([]struct {
starttime time.Time
endtime time.Time
}, c.maxPeriods)
return nil
}
func (c *ccCache) NewPeriod() {
c.periodTimes[c.periodIdx].endtime = time.Now()
c.periodIdx = (c.periodIdx + 1) % c.maxPeriods
fmt.Printf("New period index %d\n", c.periodIdx)
c.periods[c.periodIdx] = c.periods[c.periodIdx][:0]
c.periodTimes[c.periodIdx].starttime = time.Now()
c.periodTimes[c.periodIdx].endtime = c.periodTimes[c.periodIdx].starttime
}
func (c *ccCache) Add(msg lp.CCMessage) error {
c.periods[c.periodIdx] = append(c.periods[c.periodIdx], msg)
c.periodTimes[c.periodIdx].endtime = msg.Time()
return nil
}
func (c *ccCache) GetPeriod(offset int) (time.Time, time.Time, []lp.CCMessage) {
if offset > c.maxPeriods {
offset = offset % c.maxPeriods
}
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 // Metric cache data structure
type metricCache struct { type metricCache struct {
numPeriods int cache CCCache
curPeriod int
lock sync.Mutex
intervals []*metricCachePeriod
wg *sync.WaitGroup wg *sync.WaitGroup
ticker mct.MultiChanTicker ticker mct.MultiChanTicker
tickchan chan time.Time tickchan chan time.Time
done chan bool done chan bool
output chan lp.CCMessage output chan lp.CCMessage
aggEngine agg.MetricAggregator aggEngine agg.MetricAggregator
numPeriods int
started bool
} }
type MetricCache interface { type MetricCache interface {
Init(output chan lp.CCMessage, ticker mct.MultiChanTicker, wg *sync.WaitGroup, numPeriods int) error Init(output chan lp.CCMessage, ticker mct.MultiChanTicker, wg *sync.WaitGroup, interval time.Duration, numPeriods int) error
Start() Start()
Add(metric lp.CCMessage) Add(metric lp.CCMessage)
GetPeriod(index int) (time.Time, time.Time, []lp.CCMessage)
AddAggregation(name, function, condition string, tags, meta map[string]string) error AddAggregation(name, function, condition string, tags, meta map[string]string) error
DeleteAggregation(name string) error DeleteAggregation(name string) error
Close() Close()
} }
func (c *metricCache) Init(output chan lp.CCMessage, ticker mct.MultiChanTicker, wg *sync.WaitGroup, numPeriods int) error { func (c *metricCache) Init(output chan lp.CCMessage, ticker mct.MultiChanTicker, wg *sync.WaitGroup, interval time.Duration, numPeriods int) error {
var err error var err error
c.done = make(chan bool) c.done = make(chan bool)
c.wg = wg c.wg = wg
c.ticker = ticker c.ticker = ticker
c.numPeriods = numPeriods c.numPeriods = numPeriods
c.started = false
c.cache = new(ccCache)
c.output = output c.output = output
c.intervals = make([]*metricCachePeriod, 0)
for i := 0; i < c.numPeriods+1; i++ {
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 c.cache.Init(numPeriods)
// 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 { if err != nil {
return fmt.Errorf("MetricCache: failed to create aggregator: %w", err) 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() { func (c *metricCache) Start() {
c.tickchan = make(chan time.Time) c.tickchan = make(chan time.Time)
c.ticker.AddChannel(c.tickchan) c.ticker.AddChannel(c.tickchan)
// Router cache is done
done := func() {
cclog.ComponentDebug("MetricCache", "DONE")
close(c.done)
}
// Rotate cache interval c.wg.Add(1)
rotate := func(timestamp time.Time) int { go func() {
oldPeriod := c.curPeriod
c.curPeriod = oldPeriod + 1
if c.curPeriod >= c.numPeriods {
c.curPeriod = 0
}
c.intervals[oldPeriod].numMetrics = 0
c.intervals[oldPeriod].stopstamp = timestamp
c.intervals[c.curPeriod].startstamp = timestamp
c.intervals[c.curPeriod].stopstamp = timestamp
return oldPeriod
}
c.wg.Go(func() {
for { for {
select { select {
case <-c.done: case <-c.done:
done() c.wg.Done()
close(c.done)
cclog.ComponentDebug("MetricCache", "DONE")
return return
case tick := <-c.tickchan: case tick := <-c.tickchan:
c.lock.Lock() cclog.ComponentDebug("MetricCache", "Tick", tick)
old := rotate(tick) allmetrics := c.cache.GetAll()
// Get the last period and evaluate aggregation metrics c.cache.NewPeriod()
starttime, endtime, metrics := c.GetPeriod(old) mintime := tick
c.lock.Unlock() maxtime := mintime.AddDate(-1, 0, 0)
if len(metrics) > 0 {
c.aggEngine.Eval(starttime, endtime, metrics) for _, metric := range allmetrics {
if metric.Time().Before(mintime) {
mintime = metric.Time()
}
if metric.Time().After(maxtime) {
maxtime = metric.Time()
}
}
if len(allmetrics) > 0 {
cclog.ComponentDebugf("MetricCache", "Evaluate %d metrics from %v to %v", len(allmetrics), mintime.UnixNano(), maxtime.UnixNano())
c.wg.Add(1)
go func() {
c.aggEngine.Eval(mintime, maxtime, allmetrics)
c.wg.Done()
}()
} else { } else {
// This message is also printed in the first interval after startup // This message is also printed in the first interval after startup
cclog.ComponentDebug("MetricCache", "EMPTY INTERVAL?") cclog.ComponentDebug("MetricCache", "EMPTY INTERVAL?")
} }
} }
} }
}) }()
cclog.ComponentDebug("MetricCache", "START")
cclog.ComponentDebug("MetricCache", "STARTED")
} }
// Add a metric to the cache. The interval is defined by the global timer (rotate() in Start()) // Add a metric to the cache. The interval is defined by the global timer (rotate() in Start())
// The intervals list is used as round-robin buffer and the metric list grows dynamically and // The intervals list is used as round-robin buffer and the metric list grows dynamically and
// to avoid reallocations // to avoid reallocations
func (c *metricCache) Add(metric lp.CCMessage) { func (c *metricCache) Add(metric lp.CCMessage) {
if c.curPeriod >= 0 && c.curPeriod < c.numPeriods { c.cache.Add(metric)
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()
}
} }
func (c *metricCache) AddAggregation(name, function, condition string, tags, meta map[string]string) error { 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) return c.aggEngine.DeleteAggregation(name)
} }
// Get all metrics of a interval. The index is the difference to the current interval, so index=0
// is the current one, index=1 the last interval and so on. Returns and empty array if a wrong index
// is given (negative index, index larger than configured number of total intervals, ...)
func (c *metricCache) GetPeriod(index int) (time.Time, time.Time, []lp.CCMessage) {
start := time.Now()
stop := time.Now()
var metrics []lp.CCMessage
if index >= 0 && index < c.numPeriods {
pindex := c.curPeriod - index
if pindex < 0 {
pindex = c.numPeriods - pindex
}
if pindex >= 0 && pindex < c.numPeriods {
start = c.intervals[pindex].startstamp
stop = c.intervals[pindex].stopstamp
metrics = c.intervals[pindex].metrics
} else {
metrics = make([]lp.CCMessage, 0)
}
} else {
metrics = make([]lp.CCMessage, 0)
}
return start, stop, metrics
}
// Close finishes / stops the metric cache // Close finishes / stops the metric cache
func (c *metricCache) Close() { func (c *metricCache) Close() {
cclog.ComponentDebug("MetricCache", "CLOSE") cclog.ComponentDebug("MetricCache", "CLOSE")
if c.started {
c.done <- true c.done <- true
c.wg.Wait()
}
cclog.ComponentDebug("MetricCache", "CLOSED")
} }
func NewCache(output chan lp.CCMessage, ticker mct.MultiChanTicker, wg *sync.WaitGroup, numPeriods int) (MetricCache, error) { func NewCache(output chan lp.CCMessage, ticker mct.MultiChanTicker, wg *sync.WaitGroup, numPeriods int) (MetricCache, error) {
c := new(metricCache) c := new(metricCache)
err := c.Init(output, ticker, wg, numPeriods) err := c.Init(output, ticker, wg, ticker.GetDuration(), numPeriods)
if err != nil { if err != nil {
return nil, err return nil, err
} }
+69
View File
@@ -0,0 +1,69 @@
// Copyright (C) NHR@FAU, University Erlangen-Nuremberg.
// All rights reserved. This file is part of cc-lib.
// Use of this source code is governed by a MIT-style
// license that can be found in the LICENSE file.
// additional authors:
// Holger Obermaier (NHR@KIT)
package metricRouter
import (
"fmt"
"strings"
"sync"
"testing"
"time"
lp "github.com/ClusterCockpit/cc-lib/v2/ccMessage"
mct "github.com/ClusterCockpit/cc-metric-collector/pkg/multiChanTicker"
)
func TestCache(t *testing.T) {
output := make(chan lp.CCMessage, 2000)
var wg sync.WaitGroup
tickTime := time.Second
testChan := make(chan time.Time)
ticker := mct.NewTicker(tickTime)
ticker.AddChannel(testChan)
maxIntervals := 100
c, err := NewCache(output, ticker, &wg, 4)
if err != nil {
t.Errorf("failed to create new cache: %s", err.Error())
return
}
c.Start()
defer c.Close()
err = c.AddAggregation("ps_input_power", "Avg(values)", "name == 'ps1_input_power' || name == 'ps2_input_power' || name == 'ps3_input_power'", map[string]string{"hostname": "<copy>", "type": "<copy>"}, map[string]string{"unit": "<copy>"})
if err != nil {
t.Errorf("failed to add aggregation for %s", "ps_input_power")
return
}
for range maxIntervals {
timestamp := <-testChan
raw_metrics := []string{
fmt.Sprintf("ps1_input_power,type=node,hostname=myhost,unit=W value=696.0 %d", timestamp.UnixNano()),
fmt.Sprintf("ps2_input_power,type=node,hostname=myhost,unit=W value=732.0 %d", timestamp.UnixNano()),
fmt.Sprintf("ps3_input_power,type=node,hostname=myhost,unit=W value=720.0 %d", timestamp.UnixNano()),
fmt.Sprintf("cpu_load,type=hwthread,type-id=0,hostname=myhost value=45 %d", timestamp.UnixNano()),
}
metrics, err := lp.FromBytes([]byte(strings.Join(raw_metrics, "\n")))
if err != nil {
t.Errorf("failed to generate metrics: %s", err.Error())
}
for _, m := range metrics {
c.Add(m)
}
if len(output) > 0 {
for i := 0; i < len(output); i++ {
m := <-output
t.Log(m.ToLineProtocol(nil))
}
}
}
}
+2 -2
View File
@@ -38,7 +38,7 @@ type metricRouterConfig struct {
HostnameTagName string `json:"hostname_tag"` // Key name used when adding the hostname to a metric (default 'hostname') HostnameTagName string `json:"hostname_tag"` // Key name used when adding the hostname to a metric (default 'hostname')
AddTags []metricRouterTagConfig `json:"add_tags"` // List of tags that are added when the condition is met AddTags []metricRouterTagConfig `json:"add_tags"` // List of tags that are added when the condition is met
DelTags []metricRouterTagConfig `json:"delete_tags"` // List of tags that are removed when the condition is met DelTags []metricRouterTagConfig `json:"delete_tags"` // List of tags that are removed when the condition is met
IntervalAgg []agg.MetricAggregatorIntervalConfig `json:"interval_aggregates"` // List of aggregation function processed at the end of an interval IntervalAgg []agg.MetricAggregatorExprIntervalConfig `json:"interval_aggregates"` // List of aggregation function processed at the end of an interval
DropMetrics []string `json:"drop_metrics"` // List of metric names to drop. For fine-grained dropping use drop_metrics_if DropMetrics []string `json:"drop_metrics"` // List of metric names to drop. For fine-grained dropping use drop_metrics_if
DropMetricsIf []string `json:"drop_metrics_if"` // List of evaluatable terms to drop metrics DropMetricsIf []string `json:"drop_metrics_if"` // List of evaluatable terms to drop metrics
RenameMetrics map[string]string `json:"rename_metrics"` // Map to rename metric name from key to value RenameMetrics map[string]string `json:"rename_metrics"` // Map to rename metric name from key to value
@@ -253,7 +253,7 @@ func (r *metricRouter) Start() {
} }
// even if the metric is dropped, it is stored in the cache for // even if the metric is dropped, it is stored in the cache for
// aggregations // aggregations
if r.config.NumCacheIntervals > 0 { if r.config.NumCacheIntervals > 0 && m != nil {
r.cache.Add(m) r.cache.Add(m)
} }
} }