diff --git a/inhibit/cache.go b/inhibit/cache.go new file mode 100644 index 0000000000..31180a2688 --- /dev/null +++ b/inhibit/cache.go @@ -0,0 +1,124 @@ +// Copyright The Prometheus Authors +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package inhibit + +import ( + "context" + "sync" + "time" + + "github.com/prometheus/common/model" + + "github.com/prometheus/alertmanager/alert" +) + +// cache contains the runtime state of the inhibit rule. +type cache struct { + equal map[model.LabelName]struct{} + + mtx sync.RWMutex + // alerts is the map of alerts that match the source matchers of the inhibit rule. + alerts map[model.Fingerprint]*alert.Alert + // index is a map of equal label fingerprint to the set of source alert fingerprints stored + // in the cache. + index map[model.Fingerprint]model.FingerprintSet +} + +func newCache(equal map[model.LabelName]struct{}) *cache { + return &cache{ + equal: equal, + alerts: make(map[model.Fingerprint]*alert.Alert), + index: make(map[model.Fingerprint]model.FingerprintSet), + } +} + +// fingerprintEquals returns the fingerprint of the equal labels of the given label set. +func (c *cache) fingerprintEquals(lset model.LabelSet) model.Fingerprint { + equalSet := make(model.LabelSet, len(c.equal)) + for n := range c.equal { + equalSet[n] = lset[n] + } + return equalSet.Fingerprint() +} + +// set adds or replaces the given source alert. +func (c *cache) set(a *alert.Alert) { + fp := a.Fingerprint() + eq := c.fingerprintEquals(a.Labels) + + c.mtx.Lock() + defer c.mtx.Unlock() + + c.alerts[fp] = a + set, ok := c.index[eq] + if !ok { + set = model.FingerprintSet{} + c.index[eq] = set + } + set[fp] = struct{}{} +} + +// find returns the fingerprint of a cached source alert that shares the equal +// labels of lset, is active at now, and satisfies match. +func (c *cache) find(lset model.LabelSet, now time.Time, match func(*alert.Alert) bool) (model.Fingerprint, bool) { + eq := c.fingerprintEquals(lset) + + c.mtx.RLock() + defer c.mtx.RUnlock() + + for fp := range c.index[eq] { + a := c.alerts[fp] + if a.ResolvedAt(now) { + continue + } + if !match(a) { + continue + } + return fp, true + } + + return model.Fingerprint(0), false +} + +func (c *cache) run(ctx context.Context, interval time.Duration) { + t := time.NewTicker(interval) + defer t.Stop() + for { + select { + case <-ctx.Done(): + return + case <-t.C: + c.gc() + } + } +} + +func (c *cache) gc() { + c.mtx.Lock() + defer c.mtx.Unlock() + + for fp, a := range c.alerts { + if !a.Resolved() { + continue + } + delete(c.alerts, fp) + + eq := c.fingerprintEquals(a.Labels) + set := c.index[eq] + delete(set, fp) + if len(set) == 0 { + delete(c.index, eq) + } + } +} diff --git a/inhibit/index.go b/inhibit/index.go deleted file mode 100644 index 7ca5c627df..0000000000 --- a/inhibit/index.go +++ /dev/null @@ -1,84 +0,0 @@ -// Copyright The Prometheus Authors -// Licensed under the Apache License, Version 2.0 (the "License"); -// you may not use this file except in compliance with the License. -// You may obtain a copy of the License at -// -// http://www.apache.org/licenses/LICENSE-2.0 -// -// Unless required by applicable law or agreed to in writing, software -// distributed under the License is distributed on an "AS IS" BASIS, -// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -// See the License for the specific language governing permissions and -// limitations under the License. - -package inhibit - -import ( - "maps" - "slices" - "sync" - - "github.com/prometheus/common/model" -) - -// index contains map of fingerprints to sets of fingerprints. -// The keys are fingerprints of the equal labels of source alerts. -// The values are the fingerprints of all source alerts sharing those equal labels. -// For more info see comments on inhibitor and InhibitRule. -type index struct { - mtx sync.RWMutex - items map[model.Fingerprint]model.FingerprintSet -} - -func newIndex() *index { - return &index{ - items: make(map[model.Fingerprint]model.FingerprintSet), - } -} - -// Get returns a copy of the fingerprints indexed under key. -func (c *index) Get(key model.Fingerprint) ([]model.Fingerprint, bool) { - c.mtx.RLock() - defer c.mtx.RUnlock() - - set, ok := c.items[key] - if !ok { - return nil, false - } - return slices.Collect(maps.Keys(set)), true -} - -func (c *index) Add(key, value model.Fingerprint) { - c.mtx.Lock() - defer c.mtx.Unlock() - - set, ok := c.items[key] - if !ok { - set = model.FingerprintSet{} - c.items[key] = set - } - set[value] = struct{}{} -} - -// Delete removes value from the set of fingerprints indexed under key, -// removing the key entirely once the set is empty. -func (c *index) Delete(key, value model.Fingerprint) { - c.mtx.Lock() - defer c.mtx.Unlock() - - set, ok := c.items[key] - if !ok { - return - } - delete(set, value) - if len(set) == 0 { - delete(c.items, key) - } -} - -func (c *index) Len() int { - c.mtx.RLock() - defer c.mtx.RUnlock() - - return len(c.items) -} diff --git a/inhibit/inhibit.go b/inhibit/inhibit.go index ea0286d040..7cb330f4ee 100644 --- a/inhibit/inhibit.go +++ b/inhibit/inhibit.go @@ -23,7 +23,6 @@ import ( "github.com/prometheus/common/model" "go.opentelemetry.io/otel" "go.opentelemetry.io/otel/attribute" - "go.opentelemetry.io/otel/codes" "go.opentelemetry.io/otel/propagation" "go.opentelemetry.io/otel/trace" @@ -33,7 +32,6 @@ import ( "github.com/prometheus/alertmanager/marker" "github.com/prometheus/alertmanager/pkg/labels" "github.com/prometheus/alertmanager/provider" - "github.com/prometheus/alertmanager/store" "github.com/prometheus/alertmanager/tracing" ) @@ -123,15 +121,8 @@ func (ih *Inhibitor) processAlert(ctx context.Context, a *alert.Alert) { if r.SourceMatchers.Matches(a.Labels) { attr := attribute.String("alerting.inhibit_rule.name", r.Name) span.AddEvent("alert matched rule source", trace.WithAttributes(attr)) - if err := r.scache.Set(a); err != nil { - message := "error on set alert" - ih.logger.Error(message, "err", err) - span.SetStatus(codes.Error, message) - span.RecordError(err) - continue - } span.SetAttributes(attr) - r.sindex.Add(r.fingerprintEquals(a.Labels), a.Fingerprint()) + r.cache.set(a) } } } @@ -153,7 +144,7 @@ func (ih *Inhibitor) Run() { runCtx, runCancel := context.WithCancel(ctx) for _, rule := range ih.rules { - go rule.scache.Run(runCtx, 15*time.Minute) + go rule.cache.run(runCtx, 15*time.Minute) } g.Add(func() error { @@ -256,12 +247,7 @@ type InhibitRule struct { Equal map[model.LabelName]struct{} // Cache of alerts matching source labels. - scache *store.Alerts - - // Index of fingerprints of source alert equal labels to fingerprints of source alerts. - // The index helps speed up source alert lookups from scache significantely in scenarios with 100s of source alerts cached. - // Every source alert with the same equal labels is indexed under the same key. - sindex *index + cache *cache } // NewInhibitRule returns a new InhibitRule based on a configuration definition. @@ -318,32 +304,12 @@ func NewInhibitRule(cr amcommoncfg.InhibitRule) *InhibitRule { equal[model.LabelName(ln)] = struct{}{} } - rule := &InhibitRule{ + return &InhibitRule{ Name: cr.Name, SourceMatchers: sourcem, TargetMatchers: targetm, Equal: equal, - scache: store.NewAlerts(), - sindex: newIndex(), - } - - rule.scache.SetGCCallback(rule.gcCallback) - - return rule -} - -// fingerprintEquals returns the fingerprint of the equal labels of the given label set. -func (r *InhibitRule) fingerprintEquals(lset model.LabelSet) model.Fingerprint { - equalSet := make(model.LabelSet, len(r.Equal)) - for n := range r.Equal { - equalSet[n] = lset[n] - } - return equalSet.Fingerprint() -} - -func (r *InhibitRule) gcCallback(alerts []*alert.Alert) { - for _, a := range alerts { - r.sindex.Delete(r.fingerprintEquals(a.Labels), a.Fingerprint()) + cache: newCache(equal), } } @@ -352,24 +318,7 @@ func (r *InhibitRule) gcCallback(alerts []*alert.Alert) { // is returned. If excludeTwoSidedMatch is true, alerts that match both the // source and the target side of the rule are disregarded. func (r *InhibitRule) hasEqual(lset model.LabelSet, excludeTwoSidedMatch bool, now time.Time) (model.Fingerprint, bool) { - sourceFPs, ok := r.sindex.Get(r.fingerprintEquals(lset)) - if !ok { - return model.Fingerprint(0), false - } - - for _, sourceFP := range sourceFPs { - a, err := r.scache.Get(sourceFP) - if err != nil { - continue - } - if a.ResolvedAt(now) { - continue - } - if excludeTwoSidedMatch && r.TargetMatchers.Matches(a.Labels) { - continue - } - return sourceFP, true - } - - return model.Fingerprint(0), false + return r.cache.find(lset, now, func(a *alert.Alert) bool { + return !excludeTwoSidedMatch || !r.TargetMatchers.Matches(a.Labels) + }) } diff --git a/inhibit/inhibit_test.go b/inhibit/inhibit_test.go index 092a4c3bf4..91e57a92f6 100644 --- a/inhibit/inhibit_test.go +++ b/inhibit/inhibit_test.go @@ -29,7 +29,6 @@ import ( "github.com/prometheus/alertmanager/marker" "github.com/prometheus/alertmanager/pkg/labels" "github.com/prometheus/alertmanager/provider" - "github.com/prometheus/alertmanager/store" ) var nopLogger = promslog.NewNopLogger() @@ -189,18 +188,17 @@ func TestInhibitRuleHasEqual(t *testing.T) { for _, c := range cases { t.Run(c.name, func(t *testing.T) { + equal := map[model.LabelName]struct{}{} + for _, ln := range c.equal { + equal[ln] = struct{}{} + } r := &InhibitRule{ - Equal: map[model.LabelName]struct{}{}, TargetMatchers: c.targetMatchers, - scache: store.NewAlerts(), - sindex: newIndex(), - } - for _, ln := range c.equal { - r.Equal[ln] = struct{}{} + Equal: equal, + cache: newCache(equal), } for _, v := range c.initial { - r.scache.Set(v) - r.sindex.Add(r.fingerprintEquals(v.Labels), v.Fingerprint()) + r.cache.set(v) } if _, have := r.hasEqual(c.input, c.excludeTwoSidedMatch, time.Now()); have != c.result { @@ -239,7 +237,7 @@ func TestInhibitRuleHasEqualKeepsSourceOnlyAlertAfterGCSameEqual(t *testing.T) { _, found := r.hasEqual(target, true, now) require.True(t, found) - r.gcCallback([]*alert.Alert{expiredSameEqual}) + r.cache.gc() _, found = r.hasEqual(target, true, now) require.True(t, found) @@ -247,7 +245,6 @@ func TestInhibitRuleHasEqualKeepsSourceOnlyAlertAfterGCSameEqual(t *testing.T) { func TestInhibitRuleGCCallbackDoesNotRemoveRefreshedSameFingerprintSourceAlert(t *testing.T) { t.Parallel() - t.Skip("gcCallback races the inhibitor and drops an index entry the cache still holds, skipped until the fix lands") now := time.Now() oldSource := &alert.Alert{ @@ -270,7 +267,7 @@ func TestInhibitRuleGCCallbackDoesNotRemoveRefreshedSameFingerprintSourceAlert(t ih := runInhibitor(t, []amcommoncfg.InhibitRule{{Equal: []string{"e"}}}, oldSource, refreshedSource) r := ih.rules[0] - r.gcCallback([]*alert.Alert{oldSource}) + r.cache.gc() _, found := r.hasEqual(model.LabelSet{"t": "1", "e": "1"}, false, now) require.True(t, found) @@ -309,15 +306,8 @@ func TestInhibitRuleMatches(t *testing.T) { }, } - ih.rules[0].scache = store.NewAlerts() - ih.rules[0].scache.Set(sourceAlert1) - ih.rules[0].sindex = newIndex() - ih.rules[0].sindex.Add(ih.rules[0].fingerprintEquals(sourceAlert1.Labels), sourceAlert1.Fingerprint()) - - ih.rules[1].scache = store.NewAlerts() - ih.rules[1].scache.Set(sourceAlert2) - ih.rules[1].sindex = newIndex() - ih.rules[1].sindex.Add(ih.rules[1].fingerprintEquals(sourceAlert2.Labels), sourceAlert2.Fingerprint()) + ih.rules[0].cache.set(sourceAlert1) + ih.rules[1].cache.set(sourceAlert2) cases := []struct { target model.LabelSet @@ -407,15 +397,8 @@ func TestInhibitRuleMatchers(t *testing.T) { }, } - ih.rules[0].scache = store.NewAlerts() - ih.rules[0].scache.Set(sourceAlert1) - ih.rules[0].sindex = newIndex() - ih.rules[0].sindex.Add(ih.rules[0].fingerprintEquals(sourceAlert1.Labels), sourceAlert1.Fingerprint()) - - ih.rules[1].scache = store.NewAlerts() - ih.rules[1].scache.Set(sourceAlert2) - ih.rules[1].sindex = newIndex() - ih.rules[1].sindex.Add(ih.rules[1].fingerprintEquals(sourceAlert2.Labels), sourceAlert2.Fingerprint()) + ih.rules[0].cache.set(sourceAlert1) + ih.rules[1].cache.set(sourceAlert2) cases := []struct { target model.LabelSet @@ -677,12 +660,10 @@ func TestInhibit(t *testing.T) { } func TestInhibitRule_fingerprintEquals(t *testing.T) { - rule := &InhibitRule{ - Equal: map[model.LabelName]struct{}{ - "cluster": {}, - "service": {}, - }, - } + c := newCache(map[model.LabelName]struct{}{ + "cluster": {}, + "service": {}, + }) lset := model.LabelSet{ "cluster": "prod", @@ -690,7 +671,7 @@ func TestInhibitRule_fingerprintEquals(t *testing.T) { "instance": "host1", } - fp := rule.fingerprintEquals(lset) + fp := c.fingerprintEquals(lset) // Same equal labels should produce same fingerprint lset2 := model.LabelSet{ @@ -698,24 +679,19 @@ func TestInhibitRule_fingerprintEquals(t *testing.T) { "service": "api", "instance": "host2", // different non-equal label } - require.Equal(t, fp, rule.fingerprintEquals(lset2)) + require.Equal(t, fp, c.fingerprintEquals(lset2)) // Different equal label value should produce different fingerprint lset3 := model.LabelSet{ "cluster": "staging", "service": "api", } - require.NotEqual(t, fp, rule.fingerprintEquals(lset3)) + require.NotEqual(t, fp, c.fingerprintEquals(lset3)) } func TestInhibitRuleIndexSurvivesGC(t *testing.T) { now := time.Now() - r := &InhibitRule{ - Equal: map[model.LabelName]struct{}{"cluster": {}}, - scache: store.NewAlerts(), - sindex: newIndex(), - } - r.scache.SetGCCallback(r.gcCallback) + r := NewInhibitRule(amcommoncfg.InhibitRule{Equal: []string{"cluster"}}) active := &alert.Alert{Alert: model.Alert{ Labels: model.LabelSet{"alertname": "S1", "cluster": "c1"}, @@ -727,44 +703,38 @@ func TestInhibitRuleIndexSurvivesGC(t *testing.T) { StartsAt: now.Add(-time.Hour), EndsAt: now.Add(-time.Minute), }} - for _, a := range []*alert.Alert{active, resolved} { - require.NoError(t, r.scache.Set(a)) - r.sindex.Add(r.fingerprintEquals(a.Labels), a.Fingerprint()) - } + r.cache.set(active) + r.cache.set(resolved) target := model.LabelSet{"alertname": "T", "cluster": "c1"} fp, ok := r.hasEqual(target, false, now) require.True(t, ok) require.Equal(t, active.Fingerprint(), fp) - deleted := r.scache.GC() - require.Len(t, deleted, 1) - require.Equal(t, resolved.Fingerprint(), deleted[0].Fingerprint()) + r.cache.gc() + require.Len(t, r.cache.alerts, 1) + require.Contains(t, r.cache.alerts, active.Fingerprint()) fp, ok = r.hasEqual(target, false, now) require.True(t, ok, "active source alert must still inhibit after GC of a sibling") require.Equal(t, active.Fingerprint(), fp) - require.Equal(t, 1, r.sindex.Len()) + require.Len(t, r.cache.index, 1) active.EndsAt = now.Add(-time.Second) - require.NoError(t, r.scache.Set(active)) - deleted = r.scache.GC() - require.Len(t, deleted, 1) + r.cache.set(active) + r.cache.gc() _, ok = r.hasEqual(target, false, now) require.False(t, ok) - require.Equal(t, 0, r.sindex.Len(), "empty index keys must be removed") + require.Empty(t, r.cache.alerts) + require.Empty(t, r.cache.index, "empty index keys must be removed") } func TestInhibitRuleTwoSidedDoesNotShadow(t *testing.T) { now := time.Now() - targetMatcher, err := labels.NewMatcher(labels.MatchEqual, "severity", "warning") - require.NoError(t, err) - r := &InhibitRule{ - Equal: map[model.LabelName]struct{}{"cluster": {}}, - TargetMatchers: labels.Matchers{targetMatcher}, - scache: store.NewAlerts(), - sindex: newIndex(), - } + r := NewInhibitRule(amcommoncfg.InhibitRule{ + TargetMatchers: amcommoncfg.Matchers{&labels.Matcher{Type: labels.MatchEqual, Name: "severity", Value: "warning"}}, + Equal: []string{"cluster"}, + }) sourceOnly := &alert.Alert{Alert: model.Alert{ Labels: model.LabelSet{"alertname": "S1", "cluster": "c1", "severity": "critical"}, @@ -776,10 +746,8 @@ func TestInhibitRuleTwoSidedDoesNotShadow(t *testing.T) { StartsAt: now.Add(-time.Hour), EndsAt: now.Add(2 * time.Hour), }} - for _, a := range []*alert.Alert{sourceOnly, twoSided} { - require.NoError(t, r.scache.Set(a)) - r.sindex.Add(r.fingerprintEquals(a.Labels), a.Fingerprint()) - } + r.cache.set(sourceOnly) + r.cache.set(twoSided) target := model.LabelSet{"alertname": "T", "cluster": "c1", "severity": "warning"} fp, ok := r.hasEqual(target, true, now) @@ -796,7 +764,7 @@ func BenchmarkFingerprintEquals(b *testing.B) { equalLabels[model.LabelName(fmt.Sprintf("label_%d", i))] = struct{}{} } - rule := &InhibitRule{Equal: equalLabels} + c := newCache(equalLabels) // Create a label set with matching values lset := make(model.LabelSet, numLabels+2) @@ -810,7 +778,7 @@ func BenchmarkFingerprintEquals(b *testing.B) { b.ReportAllocs() for b.Loop() { - _ = rule.fingerprintEquals(lset) + _ = c.fingerprintEquals(lset) } }) }