From 667fea6623d733d0a7be7a4ce34e96d9d5124b09 Mon Sep 17 00:00:00 2001 From: Alexander Akhmetov Date: Tue, 11 Jun 2024 18:34:58 +0200 Subject: [PATCH] Alerting: use hash of labels instead of labels string as the alert state cache key (#88956) * Alerting: use hash instead of labels as the cache key * Use data.Labels.Fingerprint to calculate the cache key --- pkg/services/ngalert/backtesting/engine.go | 2 +- .../ngalert/backtesting/engine_test.go | 17 ++++--- .../ngalert/models/instance_labels.go | 4 ++ .../ngalert/models/instance_labels_test.go | 14 ++++++ .../ngalert/schedule/alert_rule_test.go | 2 +- pkg/services/ngalert/state/cache.go | 16 +++---- pkg/services/ngalert/state/manager.go | 10 ++-- .../ngalert/state/manager_private_test.go | 38 ++++++--------- pkg/services/ngalert/state/manager_test.go | 46 +++++++------------ pkg/services/ngalert/state/persister_sync.go | 2 +- pkg/services/ngalert/state/state.go | 2 +- 11 files changed, 72 insertions(+), 81 deletions(-) diff --git a/pkg/services/ngalert/backtesting/engine.go b/pkg/services/ngalert/backtesting/engine.go index 30d4c8135b9..2a53ad61582 100644 --- a/pkg/services/ngalert/backtesting/engine.go +++ b/pkg/services/ngalert/backtesting/engine.go @@ -90,7 +90,7 @@ func (e *Engine) Test(ctx context.Context, user identity.Requester, rule *models start := time.Now() tsField := data.NewField("Time", nil, make([]time.Time, length)) - valueFields := make(map[string]*data.Field) + valueFields := make(map[data.Fingerprint]*data.Field) err = evaluator.Eval(ruleCtx, from, time.Duration(rule.IntervalSeconds)*time.Second, length, func(idx int, currentTime time.Time, results eval.Results) error { if idx >= length { diff --git a/pkg/services/ngalert/backtesting/engine_test.go b/pkg/services/ngalert/backtesting/engine_test.go index 26bb317aaad..9b0bf8304f5 100644 --- a/pkg/services/ngalert/backtesting/engine_test.go +++ b/pkg/services/ngalert/backtesting/engine_test.go @@ -160,9 +160,10 @@ func TestNewBacktestingEvaluator(t *testing.T) { func TestEvaluatorTest(t *testing.T) { states := []eval.State{eval.Normal, eval.Alerting, eval.Pending} generateState := func(prefix string) *state.State { + labels := models.GenerateAlertLabels(rand.Intn(5)+1, prefix+"-") return &state.State{ - CacheID: "state-" + prefix, - Labels: models.GenerateAlertLabels(rand.Intn(5)+1, prefix+"-"), + CacheID: labels.Fingerprint(), + Labels: labels, State: states[rand.Intn(len(states))], } } @@ -201,10 +202,11 @@ func TestEvaluatorTest(t *testing.T) { var states []state.StateTransition for _, s := range allStates { + labels := models.GenerateAlertLabels(rand.Intn(5)+1, s.String()+"-") states = append(states, state.StateTransition{ State: &state.State{ - CacheID: "state-" + s.String(), - Labels: models.GenerateAlertLabels(rand.Intn(5)+1, s.String()+"-"), + CacheID: labels.Fingerprint(), + Labels: labels, State: s, StateReason: util.GenerateShortUID(), }, @@ -226,7 +228,7 @@ func TestEvaluatorTest(t *testing.T) { require.Equal(t, data.FieldTypeTime, timestampField.Type()) }) - fieldByState := make(map[string]*data.Field, len(states)) + fieldByState := make(map[data.Fingerprint]*data.Field, len(states)) t.Run("should contain a field per state", func(t *testing.T) { for _, s := range states { @@ -269,11 +271,12 @@ func TestEvaluatorTest(t *testing.T) { from := time.Unix(0, 0) to := from.Add(5 * ruleInterval) + labels := models.GenerateAlertLabels(rand.Intn(5)+1, "test-") states := []state.StateTransition{ { State: &state.State{ - CacheID: "state-1", - Labels: models.GenerateAlertLabels(rand.Intn(5)+1, "test-"), + CacheID: labels.Fingerprint(), + Labels: labels, State: eval.Normal, StateReason: util.GenerateShortUID(), }, diff --git a/pkg/services/ngalert/models/instance_labels.go b/pkg/services/ngalert/models/instance_labels.go index 468ea3da032..1904515cbd9 100644 --- a/pkg/services/ngalert/models/instance_labels.go +++ b/pkg/services/ngalert/models/instance_labels.go @@ -65,6 +65,10 @@ func (il *InstanceLabels) StringAndHash() (string, string, error) { return string(b), fmt.Sprintf("%x", h.Sum(nil)), nil } +func (il *InstanceLabels) Fingerprint() data.Fingerprint { + return data.Labels(*il).Fingerprint() +} + // The following is based on SDK code, copied for now // tupleLables is an alternative representation of Labels (map[string]string) that can be sorted diff --git a/pkg/services/ngalert/models/instance_labels_test.go b/pkg/services/ngalert/models/instance_labels_test.go index 513841798de..e6783438f6d 100644 --- a/pkg/services/ngalert/models/instance_labels_test.go +++ b/pkg/services/ngalert/models/instance_labels_test.go @@ -4,6 +4,7 @@ import ( "fmt" "testing" + "github.com/grafana/grafana-plugin-sdk-go/data" "github.com/stretchr/testify/require" ) @@ -39,6 +40,19 @@ func TestTupleLabelsToLabels(t *testing.T) { }) } +func TestInstanceLabelsFingerprint(t *testing.T) { + t.Run("returns labels fingerprint", func(t *testing.T) { + labels := InstanceLabels{ + "foo": "bar", + "baz": "qux", + } + + fingerprint := labels.Fingerprint() + expectedFingerprint := data.Labels(labels).Fingerprint() + require.Equal(t, expectedFingerprint, fingerprint) + }) +} + func BenchmarkTupleLabelsToLabels(b *testing.B) { b.Run("10 labels", func(b *testing.B) { in := make(tupleLabels, 0, 10) diff --git a/pkg/services/ngalert/schedule/alert_rule_test.go b/pkg/services/ngalert/schedule/alert_rule_test.go index 49b118e0b75..12eb90ec661 100644 --- a/pkg/services/ngalert/schedule/alert_rule_test.go +++ b/pkg/services/ngalert/schedule/alert_rule_test.go @@ -531,7 +531,7 @@ func TestRuleRoutine(t *testing.T) { for i := 0; i < 2; i++ { states = append(states, &state.State{ AlertRuleUID: rule.UID, - CacheID: util.GenerateShortUID(), + CacheID: data.Labels(rule.Labels).Fingerprint(), OrgID: rule.OrgID, State: s, StartsAt: sch.clock.Now(), diff --git a/pkg/services/ngalert/state/cache.go b/pkg/services/ngalert/state/cache.go index ae82381310e..25e8a083b67 100644 --- a/pkg/services/ngalert/state/cache.go +++ b/pkg/services/ngalert/state/cache.go @@ -20,7 +20,7 @@ import ( ) type ruleStates struct { - states map[string]*State + states map[data.Fingerprint]*State } type cache struct { @@ -89,7 +89,7 @@ func (c *cache) getOrCreate(ctx context.Context, log log.Logger, alertRule *ngMo } var states *ruleStates if states, ok = orgStates[stateCandidate.AlertRuleUID]; !ok { - states = &ruleStates{states: make(map[string]*State)} + states = &ruleStates{states: make(map[data.Fingerprint]*State)} c.states[stateCandidate.OrgID][stateCandidate.AlertRuleUID] = states } return states.getOrAdd(stateCandidate) @@ -199,18 +199,14 @@ func calculateState(ctx context.Context, log log.Logger, alertRule *ngModels.Ale log.Warn("Evaluation result contains either reserved labels or labels declared in the rules. Those labels from the result will be ignored", "labels", dupes) } - il := ngModels.InstanceLabels(lbs) - id, err := il.StringKey() - if err != nil { - log.Error("Error getting cacheId for entry", "error", err) - } + cacheID := lbs.Fingerprint() // For new states, we set StartsAt & EndsAt to EvaluatedAt as this is the // expected value for a Normal state during state transition. newState := State{ AlertRuleUID: alertRule.UID, OrgID: alertRule.OrgID, - CacheID: id, + CacheID: cacheID, Labels: lbs, Annotations: annotations, EvaluationDuration: result.EvaluationDuration, @@ -279,12 +275,12 @@ func (c *cache) set(entry *State) { c.states[entry.OrgID] = make(map[string]*ruleStates) } if _, ok := c.states[entry.OrgID][entry.AlertRuleUID]; !ok { - c.states[entry.OrgID][entry.AlertRuleUID] = &ruleStates{states: make(map[string]*State)} + c.states[entry.OrgID][entry.AlertRuleUID] = &ruleStates{states: make(map[data.Fingerprint]*State)} } c.states[entry.OrgID][entry.AlertRuleUID].states[entry.CacheID] = entry } -func (c *cache) get(orgID int64, alertRuleUID, stateId string) *State { +func (c *cache) get(orgID int64, alertRuleUID string, stateId data.Fingerprint) *State { c.mtxStates.RLock() defer c.mtxStates.RUnlock() ruleStates, ok := c.states[orgID][alertRuleUID] diff --git a/pkg/services/ngalert/state/manager.go b/pkg/services/ngalert/state/manager.go index fd44ef1a27b..10050b60037 100644 --- a/pkg/services/ngalert/state/manager.go +++ b/pkg/services/ngalert/state/manager.go @@ -180,15 +180,12 @@ func (st *Manager) Warm(ctx context.Context, rulesReader RuleReader) { rulesStates, ok := orgStates[entry.RuleUID] if !ok { - rulesStates = &ruleStates{states: make(map[string]*State)} + rulesStates = &ruleStates{states: make(map[data.Fingerprint]*State)} orgStates[entry.RuleUID] = rulesStates } lbs := map[string]string(entry.Labels) - cacheID, err := entry.Labels.StringKey() - if err != nil { - st.log.Error("Error getting cacheId for entry", "error", err) - } + cacheID := entry.Labels.Fingerprint() var resultFp data.Fingerprint if entry.ResultFingerprint != "" { fp, err := strconv.ParseUint(entry.ResultFingerprint, 16, 64) @@ -214,11 +211,12 @@ func (st *Manager) Warm(ctx context.Context, rulesReader RuleReader) { statesCount++ } } + st.cache.setAllStates(states) st.log.Info("State cache has been initialized", "states", statesCount, "duration", time.Since(startTime)) } -func (st *Manager) Get(orgID int64, alertRuleUID, stateId string) *State { +func (st *Manager) Get(orgID int64, alertRuleUID string, stateId data.Fingerprint) *State { return st.cache.get(orgID, alertRuleUID, stateId) } diff --git a/pkg/services/ngalert/state/manager_private_test.go b/pkg/services/ngalert/state/manager_private_test.go index dd284c18391..8740186deac 100644 --- a/pkg/services/ngalert/state/manager_private_test.go +++ b/pkg/services/ngalert/state/manager_private_test.go @@ -276,7 +276,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { if !ok { // skip if nothing to assert continue } - expectedTransitionsMap := make(map[string]StateTransition, len(expectedTransitions)) + expectedTransitionsMap := make(map[data.Fingerprint]StateTransition, len(expectedTransitions)) for i := range expectedTransitions { patchState(alertRule, expectedTransitions[i].State) expectedTransitionsMap[expectedTransitions[i].CacheID] = expectedTransitions[i] @@ -294,7 +294,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { } } if len(expectedTransitionsMap) > 0 { - vals := make([]string, 0, len(expectedTransitionsMap)) + vals := make([]data.Fingerprint, 0, len(expectedTransitionsMap)) for _, s := range expectedTransitionsMap { vals = append(vals, s.CacheID) } @@ -2524,15 +2524,6 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { ngmodels.KeepLastErrState: baseRuleWith(ngmodels.RuleMuts.WithErrorExecAs(ngmodels.KeepLastErrState)), } - cacheID := func(lbls data.Labels) string { - l := ngmodels.InstanceLabels(lbls) - r, err := l.StringKey() - if err != nil { - panic(err) - } - return r - } - type errorTestCase struct { desc string ruleMutators []ngmodels.AlertRuleMutator @@ -2587,7 +2578,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { { PreviousState: eval.Normal, State: &State{ - CacheID: cacheID(labels["system + rule"]), + CacheID: labels["system + rule"].Fingerprint(), Labels: labels["system + rule + datasource-error"], State: eval.Error, Error: datasourceError, @@ -2748,7 +2739,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { { PreviousState: eval.Normal, State: &State{ - CacheID: cacheID(labels["system + rule"]), + CacheID: labels["system + rule"].Fingerprint(), Labels: labels["system + rule + datasource-error"], State: eval.Error, Error: datasourceError, @@ -2881,7 +2872,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { { PreviousState: eval.Normal, State: &State{ - CacheID: cacheID(labels["system + rule"]), + CacheID: labels["system + rule"].Fingerprint(), Labels: labels["system + rule + datasource-error"], State: eval.Error, Error: datasourceError, @@ -3028,7 +3019,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { { PreviousState: eval.Error, State: &State{ - CacheID: cacheID(labels["system + rule"]), + CacheID: labels["system + rule"].Fingerprint(), Labels: labels["system + rule + datasource-error"], Error: datasourceError, State: eval.Normal, @@ -3148,7 +3139,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { { PreviousState: eval.Normal, State: &State{ - CacheID: cacheID(labels["system + rule"]), + CacheID: labels["system + rule"].Fingerprint(), Labels: labels["system + rule + datasource-error"], State: eval.Error, Error: datasourceError, @@ -3244,7 +3235,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { { PreviousState: eval.Pending, State: &State{ - CacheID: cacheID(labels["system + rule"]), + CacheID: labels["system + rule"].Fingerprint(), Labels: labels["system + rule + datasource-error"], State: eval.Error, Error: datasourceError, @@ -3330,7 +3321,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { { PreviousState: eval.Pending, State: &State{ - CacheID: cacheID(labels["system + rule"]), + CacheID: labels["system + rule"].Fingerprint(), Labels: labels["system + rule + datasource-error"], State: eval.Error, Error: datasourceError, @@ -3539,14 +3530,11 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { } func setCacheID(s *State) *State { - if s.CacheID != "" { + if s.CacheID != 0 { return s } - il := ngmodels.InstanceLabels(s.Labels) - id, err := il.StringKey() - if err != nil { - panic(err) - } - s.CacheID = id + + s.CacheID = s.Labels.Fingerprint() + return s } diff --git a/pkg/services/ngalert/state/manager_test.go b/pkg/services/ngalert/state/manager_test.go index cc116a8b8e9..edb836b4dac 100644 --- a/pkg/services/ngalert/state/manager_test.go +++ b/pkg/services/ngalert/state/manager_test.go @@ -973,13 +973,9 @@ func TestProcessEvalResults(t *testing.T) { expectedAnnotations: 1, expectedStates: []*state.State{ { - CacheID: func() string { + CacheID: func() data.Fingerprint { lbls := models.InstanceLabels(labels["system + rule + labels1"]) - r, err := lbls.StringKey() - if err != nil { - panic(err) - } - return r + return lbls.Fingerprint() }(), Labels: mergeLabels(labels["system + rule + labels1"], data.Labels{ "datasource_uid": "datasource_uid_1", @@ -1317,7 +1313,7 @@ func TestProcessEvalResults(t *testing.T) { states := st.GetStatesForRuleUID(tc.alertRule.OrgID, tc.alertRule.UID) assert.Len(t, states, len(tc.expectedStates)) - expectedStates := make(map[string]*state.State, len(tc.expectedStates)) + expectedStates := make(map[data.Fingerprint]*state.State, len(tc.expectedStates)) for _, s := range tc.expectedStates { // patch all optional fields of the expected state setCacheID(s) @@ -1404,12 +1400,11 @@ func TestProcessEvalResults(t *testing.T) { require.NotEmpty(t, states) - savedStates := make(map[string]models.AlertInstance) + savedStates := make(map[data.Fingerprint]models.AlertInstance) for _, op := range instanceStore.RecordedOps() { switch q := op.(type) { case models.AlertInstance: - cacheId, err := q.Labels.StringKey() - require.NoError(t, err) + cacheId := q.Labels.Fingerprint() savedStates[cacheId] = q } } @@ -1578,7 +1573,7 @@ func TestStaleResultsHandler(t *testing.T) { } func TestStaleResults(t *testing.T) { - getCacheID := func(t *testing.T, rule *models.AlertRule, result eval.Result) string { + getCacheID := func(t *testing.T, rule *models.AlertRule, result eval.Result) data.Fingerprint { t.Helper() labels := data.Labels{} for key, value := range rule.Labels { @@ -1588,14 +1583,12 @@ func TestStaleResults(t *testing.T) { labels[key] = value } lbls := models.InstanceLabels(labels) - key, err := lbls.StringKey() - require.NoError(t, err) - return key + return lbls.Fingerprint() } - checkExpectedStates := func(t *testing.T, actual []*state.State, expected map[string]struct{}) map[string]*state.State { + checkExpectedStates := func(t *testing.T, actual []*state.State, expected map[data.Fingerprint]struct{}) map[data.Fingerprint]*state.State { t.Helper() - result := make(map[string]*state.State) + result := make(map[data.Fingerprint]*state.State) require.Len(t, actual, len(expected)) for _, currentState := range actual { _, ok := expected[currentState.CacheID] @@ -1604,7 +1597,7 @@ func TestStaleResults(t *testing.T) { } return result } - checkExpectedStateTransitions := func(t *testing.T, actual []state.StateTransition, expected map[string]struct{}) { + checkExpectedStateTransitions := func(t *testing.T, actual []state.StateTransition, expected map[data.Fingerprint]struct{}) { t.Helper() require.Len(t, actual, len(expected)) for _, currentState := range actual { @@ -1643,7 +1636,7 @@ func TestStaleResults(t *testing.T) { state2 := getCacheID(t, rule, initResults[1]) state3 := getCacheID(t, rule, initResults[2]) - initStates := map[string]struct{}{ + initStates := map[data.Fingerprint]struct{}{ state1: {}, state2: {}, state3: {}, @@ -1687,7 +1680,7 @@ func TestStaleResults(t *testing.T) { t.Run("should remove stale states from cache", func(t *testing.T) { currentStates = st.GetStatesForRuleUID(rule.OrgID, rule.UID) - checkExpectedStates(t, currentStates, map[string]struct{}{ + checkExpectedStates(t, currentStates, map[data.Fingerprint]struct{}{ getCacheID(t, rule, results[0]): {}, }) }) @@ -1782,7 +1775,7 @@ func TestDeleteStateByRuleUID(t *testing.T) { } for _, tc := range testCases { - expectedStatesMap := make(map[string]*state.State, len(tc.expectedStates)) + expectedStatesMap := make(map[data.Fingerprint]*state.State, len(tc.expectedStates)) for _, expectedState := range tc.expectedStates { s := setCacheID(expectedState) expectedStatesMap[s.CacheID] = s @@ -1992,20 +1985,15 @@ func TestResetStateByRuleUID(t *testing.T) { } func setCacheID(s *state.State) *state.State { - if s.CacheID != "" { + if s.CacheID != 0 { return s } - il := models.InstanceLabels(s.Labels) - id, err := il.StringKey() - if err != nil { - panic(err) - } - s.CacheID = id + s.CacheID = s.Labels.Fingerprint() return s } -func stateSliceToMap(states []*state.State) map[string]*state.State { - result := make(map[string]*state.State, len(states)) +func stateSliceToMap(states []*state.State) map[data.Fingerprint]*state.State { + result := make(map[data.Fingerprint]*state.State, len(states)) for _, s := range states { setCacheID(s) result[s.CacheID] = s diff --git a/pkg/services/ngalert/state/persister_sync.go b/pkg/services/ngalert/state/persister_sync.go index bc8c5cf1673..90ddd350498 100644 --- a/pkg/services/ngalert/state/persister_sync.go +++ b/pkg/services/ngalert/state/persister_sync.go @@ -56,7 +56,7 @@ func (a *SyncStatePersister) deleteAlertStates(ctx context.Context, states []Sta for _, s := range states { key, err := s.GetAlertInstanceKey() if err != nil { - a.log.Error("Failed to delete alert instance with invalid labels", "cacheID", s.CacheID, "error", err) + a.log.Error("Failed to delete alert instance with invalid labels", "cacheID", s.CacheID, "labels", s.Labels.String(), "error", err) continue } toDelete = append(toDelete, key) diff --git a/pkg/services/ngalert/state/state.go b/pkg/services/ngalert/state/state.go index cb6cf0767a9..cc349d7ff43 100644 --- a/pkg/services/ngalert/state/state.go +++ b/pkg/services/ngalert/state/state.go @@ -27,7 +27,7 @@ type State struct { // CacheID is a unique, opaque identifier for the state, and is used to find the state // in the state cache. It tends to be derived from the state's labels. - CacheID string + CacheID data.Fingerprint // State represents the current state. State eval.State