diff --git a/conf/defaults.ini b/conf/defaults.ini index 898b0261fa2..224555260f2 100644 --- a/conf/defaults.ini +++ b/conf/defaults.ini @@ -1304,6 +1304,9 @@ disable_jitter = false # Retention period for Alertmanager notification log entries. notification_log_retention = 5d +# Duration for which a resolved alert state transition will continue to be sent to the Alertmanager. +resolved_alert_retention = 15m + [unified_alerting.screenshots] # Enable screenshots in notifications. You must have either installed the Grafana image rendering # plugin, or set up Grafana to use a remote rendering service. diff --git a/conf/sample.ini b/conf/sample.ini index b1d61c44187..edd415de28b 100644 --- a/conf/sample.ini +++ b/conf/sample.ini @@ -792,7 +792,7 @@ ;role_attribute_strict = false ;groups_attribute_path = ;id_token_attribute_name = -;team_ids_attribute_path +;team_ids_attribute_path ;auth_url = https://foo.bar/login/oauth/authorize ;token_url = https://foo.bar/login/oauth/access_token ;api_url = https://foo.bar/user @@ -1290,6 +1290,9 @@ # Retention period for Alertmanager notification log entries. ;notification_log_retention = 5d +# Duration for which a resolved alert state transition will continue to be sent to the Alertmanager. +;resolved_alert_retention = 15m + [unified_alerting.screenshots] # Enable screenshots in notifications. You must have either installed the Grafana image rendering # plugin, or set up Grafana to use a remote rendering service. @@ -1837,4 +1840,4 @@ timeout = 30s #################################### Public Dashboards ##################################### [public_dashboards] # Set to false to disable public dashboards -;enabled = true \ No newline at end of file +;enabled = true diff --git a/pkg/services/ngalert/ngalert.go b/pkg/services/ngalert/ngalert.go index ec86f1ed5f0..952e1d86814 100644 --- a/pkg/services/ngalert/ngalert.go +++ b/pkg/services/ngalert/ngalert.go @@ -377,6 +377,7 @@ func (ng *AlertNG) init() error { RulesPerRuleGroupLimit: ng.Cfg.UnifiedAlerting.RulesPerRuleGroupLimit, Tracer: ng.tracer, Log: log.New("ngalert.state.manager"), + ResolvedRetention: ng.Cfg.UnifiedAlerting.ResolvedAlertRetention, } logger := log.New("ngalert.state.manager.persist") statePersister := state.NewSyncStatePersisiter(logger, cfg) diff --git a/pkg/services/ngalert/schedule/alert_rule.go b/pkg/services/ngalert/schedule/alert_rule.go index 6c3d58ec3e8..f033448f10f 100644 --- a/pkg/services/ngalert/schedule/alert_rule.go +++ b/pkg/services/ngalert/schedule/alert_rule.go @@ -8,6 +8,10 @@ import ( "time" "github.com/benbjohnson/clock" + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/codes" + "go.opentelemetry.io/otel/trace" + "github.com/grafana/grafana/pkg/infra/log" "github.com/grafana/grafana/pkg/infra/tracing" "github.com/grafana/grafana/pkg/services/datasources" @@ -19,9 +23,6 @@ import ( "github.com/grafana/grafana/pkg/services/org" "github.com/grafana/grafana/pkg/services/user" "github.com/grafana/grafana/pkg/util" - "go.opentelemetry.io/otel/attribute" - "go.opentelemetry.io/otel/codes" - "go.opentelemetry.io/otel/trace" ) // Rule represents a single piece of work that is executed periodically by the ruler. @@ -418,7 +419,7 @@ func (a *alertRule) evaluate(ctx context.Context, key ngmodels.AlertRuleKey, f f processDuration.Observe(a.clock.Now().Sub(start).Seconds()) start = a.clock.Now() - alerts := state.FromStateTransitionToPostableAlerts(processedStates, a.stateManager, a.appURL) + alerts := state.FromStateTransitionToPostableAlerts(e.scheduledAt, processedStates, a.stateManager, a.appURL) span.AddEvent("results processed", trace.WithAttributes( attribute.Int64("state_transitions", int64(len(processedStates))), attribute.Int64("alerts_to_send", int64(len(alerts.PostableAlerts))), diff --git a/pkg/services/ngalert/schedule/alert_rule_test.go b/pkg/services/ngalert/schedule/alert_rule_test.go index 12eb90ec661..875c098302e 100644 --- a/pkg/services/ngalert/schedule/alert_rule_test.go +++ b/pkg/services/ngalert/schedule/alert_rule_test.go @@ -2,7 +2,7 @@ package schedule import ( "bytes" - context "context" + "context" "fmt" "math" "math/rand" @@ -11,19 +11,21 @@ import ( "testing" "time" - alertingModels "github.com/grafana/alerting/models" "github.com/grafana/grafana-plugin-sdk-go/data" "github.com/prometheus/client_golang/prometheus" "github.com/prometheus/client_golang/prometheus/testutil" prometheusModel "github.com/prometheus/common/model" "github.com/stretchr/testify/assert" - mock "github.com/stretchr/testify/mock" + "github.com/stretchr/testify/mock" "github.com/stretchr/testify/require" + alertingModels "github.com/grafana/alerting/models" + "github.com/grafana/grafana/pkg/infra/log" - definitions "github.com/grafana/grafana/pkg/services/ngalert/api/tooling/definitions" + "github.com/grafana/grafana/pkg/infra/log/logtest" + "github.com/grafana/grafana/pkg/services/ngalert/api/tooling/definitions" "github.com/grafana/grafana/pkg/services/ngalert/eval" - models "github.com/grafana/grafana/pkg/services/ngalert/models" + "github.com/grafana/grafana/pkg/services/ngalert/models" "github.com/grafana/grafana/pkg/services/ngalert/state" "github.com/grafana/grafana/pkg/util" ) @@ -762,8 +764,94 @@ func TestRuleRoutine(t *testing.T) { require.NotEmpty(t, sch.stateManager.GetStatesForRuleUID(rule.OrgID, rule.UID)) }) + + t.Run("when there are resolved alerts they should keep sending until retention period is over", func(t *testing.T) { + rule := gen.With(withQueryForState(t, eval.Normal), models.RuleMuts.WithInterval(time.Second)).GenerateRef() + + evalAppliedChan := make(chan time.Time) + + sender := NewSyncAlertsSenderMock() + sender.EXPECT().Send(mock.Anything, rule.GetKey(), mock.Anything).Return() + + sch, ruleStore, _, _ := createSchedule(evalAppliedChan, sender) + sch.stateManager.ResolvedRetention = 4 * time.Second + sch.stateManager.ResendDelay = 2 * time.Second + sch.stateManager.Put([]*state.State{ + stateForRule(rule, sch.clock.Now(), eval.Alerting), // Add existing Alerting state so evals will resolve. + }) + + ruleStore.PutRule(context.Background(), rule) + factory := ruleFactoryFromScheduler(sch) + ctx, cancel := context.WithCancel(context.Background()) + t.Cleanup(cancel) + ruleInfo := factory.new(ctx, rule) + + go func() { + _ = ruleInfo.Run(rule.GetKey()) + }() + + // Evaluate 10 times: + // 1. Send resolve #1. + // 2. 2s resend delay. + // 3. Send resolve #2. + // 4. 2s resend delay. + // 5. Send resolve #3. + // 6. No more sends, 4s retention period is over. + expectedResolves := map[time.Time]struct{}{ + sch.clock.Now().Add(1 * time.Second): {}, + sch.clock.Now().Add(3 * time.Second): {}, + sch.clock.Now().Add(5 * time.Second): {}, + } + calls := 0 + for i := 1; i < 10; i++ { + ts := sch.clock.Now().Add(time.Duration(int64(i)*rule.IntervalSeconds) * time.Second) + ruleInfo.Eval(&Evaluation{ + scheduledAt: ts, + rule: rule, + }) + waitForTimeChannel(t, evalAppliedChan) + + if _, ok := expectedResolves[ts]; ok { + calls++ + prevCallAlerts, ok := sender.Calls()[calls-1].Arguments[2].(definitions.PostableAlerts) + assert.Truef(t, ok, fmt.Sprintf("expected argument of function was supposed to be 'definitions.PostableAlerts' but got %T", sender.Calls()[calls-1].Arguments[2])) + assert.Len(t, prevCallAlerts.PostableAlerts, 1) + } + sender.AssertNumberOfCalls(t, "Send", calls) + } + }) } func ruleFactoryFromScheduler(sch *schedule) ruleFactory { return newRuleFactory(sch.appURL, sch.disableGrafanaFolder, sch.maxAttempts, sch.alertsSender, sch.stateManager, sch.evaluatorFactory, &sch.schedulableAlertRules, sch.clock, sch.featureToggles, sch.metrics, sch.log, sch.tracer, sch.recordingWriter, sch.evalAppliedFunc, sch.stopAppliedFunc) } + +func stateForRule(rule *models.AlertRule, ts time.Time, evalState eval.State) *state.State { + s := &state.State{ + OrgID: rule.OrgID, + AlertRuleUID: rule.UID, + CacheID: 0, + State: evalState, + Annotations: make(map[string]string), + Labels: make(map[string]string), + StartsAt: ts, + EndsAt: ts, + ResolvedAt: &ts, + LastSentAt: &ts, + LastEvaluationTime: ts, + } + for k, v := range rule.Labels { + s.Labels[k] = v + } + for k, v := range state.GetRuleExtraLabels(&logtest.Fake{}, rule, "", true) { + if _, ok := s.Labels[k]; !ok { + s.Labels[k] = v + } + } + il := models.InstanceLabels(s.Labels) + s.Labels = data.Labels(il) + id := il.Fingerprint() + s.CacheID = id + + return s +} diff --git a/pkg/services/ngalert/schedule/testing.go b/pkg/services/ngalert/schedule/testing.go index 2f00c213ac3..d182329cf16 100644 --- a/pkg/services/ngalert/schedule/testing.go +++ b/pkg/services/ngalert/schedule/testing.go @@ -7,15 +7,17 @@ import ( "testing" "time" - definitions "github.com/grafana/grafana/pkg/services/ngalert/api/tooling/definitions" + "github.com/stretchr/testify/mock" + + "github.com/grafana/grafana/pkg/services/ngalert/api/tooling/definitions" "github.com/grafana/grafana/pkg/services/ngalert/models" - mock "github.com/stretchr/testify/mock" ) // waitForTimeChannel blocks the execution until either the channel ch has some data or a timeout of 10 second expires. // Timeout will cause the test to fail. // Returns the data from the channel. func waitForTimeChannel(t *testing.T, ch chan time.Time) time.Time { + t.Helper() select { case result := <-ch: return result diff --git a/pkg/services/ngalert/state/compat.go b/pkg/services/ngalert/state/compat.go index 673e9e56810..1a070cfd5ac 100644 --- a/pkg/services/ngalert/state/compat.go +++ b/pkg/services/ngalert/state/compat.go @@ -10,11 +10,12 @@ import ( "github.com/benbjohnson/clock" "github.com/go-openapi/strfmt" - alertingModels "github.com/grafana/alerting/models" "github.com/grafana/grafana-plugin-sdk-go/data" "github.com/prometheus/alertmanager/api/v2/models" "github.com/prometheus/common/model" + alertingModels "github.com/grafana/alerting/models" + apimodels "github.com/grafana/grafana/pkg/services/ngalert/api/tooling/definitions" "github.com/grafana/grafana/pkg/services/ngalert/eval" ngModels "github.com/grafana/grafana/pkg/services/ngalert/models" @@ -73,7 +74,7 @@ func StateToPostableAlert(transition StateTransition, appURL *url.URL) *models.P } state := alertState.State - if alertState.Resolved { + if alertState.ResolvedAt != nil { // If this is a resolved alert, we need to send an alert with the correct labels such that they will expire the previous alert. // In most cases the labels on the state will be correct, however when the previous alert was a NoData or Error alert, we need to // ensure to modify it appropriately. @@ -139,13 +140,12 @@ func errorAlert(labels, annotations data.Labels, alertState *State, urlStr strin } } -func FromStateTransitionToPostableAlerts(firingStates []StateTransition, stateManager *Manager, appURL *url.URL) apimodels.PostableAlerts { +func FromStateTransitionToPostableAlerts(evaluatedAt time.Time, firingStates []StateTransition, stateManager *Manager, appURL *url.URL) apimodels.PostableAlerts { alerts := apimodels.PostableAlerts{PostableAlerts: make([]models.PostableAlert, 0, len(firingStates))} - ts := time.Now() sentAlerts := make([]*State, 0, len(firingStates)) for _, alertState := range firingStates { - if !alertState.NeedsSending(stateManager.ResendDelay) { + if !alertState.NeedsSending(stateManager.ResendDelay, stateManager.ResolvedRetention) { continue } alert := StateToPostableAlert(alertState, appURL) @@ -153,7 +153,7 @@ func FromStateTransitionToPostableAlerts(firingStates []StateTransition, stateMa if alertState.StateReason == ngModels.StateReasonMissingSeries { // do not put stale state back to state manager continue } - alertState.LastSentAt = ts + alertState.LastSentAt = &evaluatedAt sentAlerts = append(sentAlerts, alertState.State) } stateManager.Put(sentAlerts) diff --git a/pkg/services/ngalert/state/compat_test.go b/pkg/services/ngalert/state/compat_test.go index e0f889a96ea..bd0189badb8 100644 --- a/pkg/services/ngalert/state/compat_test.go +++ b/pkg/services/ngalert/state/compat_test.go @@ -9,12 +9,13 @@ import ( "github.com/benbjohnson/clock" "github.com/go-openapi/strfmt" - alertingModels "github.com/grafana/alerting/models" "github.com/grafana/grafana-plugin-sdk-go/data" "github.com/prometheus/alertmanager/api/v2/models" "github.com/prometheus/common/model" "github.com/stretchr/testify/require" + alertingModels "github.com/grafana/alerting/models" + "github.com/grafana/grafana/pkg/services/ngalert/eval" ngModels "github.com/grafana/grafana/pkg/services/ngalert/models" "github.com/grafana/grafana/pkg/util" @@ -267,7 +268,9 @@ func TestStateToPostableAlertFromNodataError(t *testing.T) { for _, tc := range testCases { t.Run(tc.name, func(t *testing.T) { alertState := randomTransition(tc.from, tc.to) - alertState.Resolved = tc.resolved + if tc.resolved { + alertState.ResolvedAt = &alertState.LastEvaluationTime + } alertState.Labels = data.Labels(standardLabels) result := StateToPostableAlert(alertState, appURL) require.Equal(t, tc.expectedLabels, result.Labels) @@ -339,7 +342,7 @@ func randomTransition(from, to eval.State) StateTransition { EndsAt: randomTimeInFuture(), LastEvaluationTime: randomTimeInPast(), EvaluationDuration: randomDuration(), - LastSentAt: randomTimeInPast(), + LastSentAt: util.Pointer(randomTimeInPast()), Annotations: make(map[string]string), Labels: make(map[string]string), Values: make(map[string]float64), diff --git a/pkg/services/ngalert/state/manager.go b/pkg/services/ngalert/state/manager.go index 10050b60037..b4a64d4cb32 100644 --- a/pkg/services/ngalert/state/manager.go +++ b/pkg/services/ngalert/state/manager.go @@ -39,9 +39,10 @@ type Manager struct { metrics *metrics.State tracer tracing.Tracer - clock clock.Clock - cache *cache - ResendDelay time.Duration + clock clock.Clock + cache *cache + ResendDelay time.Duration + ResolvedRetention time.Duration instanceStore InstanceStore images ImageCapturer @@ -73,6 +74,9 @@ type ManagerCfg struct { DisableExecution bool + // Duration for which a resolved alert state transition will continue to be sent to the Alertmanager. + ResolvedRetention time.Duration + Tracer tracing.Tracer Log log.Logger } @@ -88,6 +92,7 @@ func NewManager(cfg ManagerCfg, statePersister StatePersister) *Manager { m := &Manager{ cache: c, ResendDelay: ResendDelay, // TODO: make this configurable + ResolvedRetention: cfg.ResolvedRetention, log: cfg.Log, metrics: cfg.Metrics, instanceStore: cfg.InstanceStore, @@ -245,7 +250,11 @@ func (st *Manager) DeleteStateByRuleUID(ctx context.Context, ruleKey ngModels.Al s.SetNormal(reason, startsAt, now) // Set Resolved property so the scheduler knows to send a postable alert // to Alertmanager. - s.Resolved = oldState == eval.Alerting || oldState == eval.Error || oldState == eval.NoData + if oldState == eval.Alerting || oldState == eval.Error || oldState == eval.NoData { + s.ResolvedAt = &now + } else { + s.ResolvedAt = nil + } s.LastEvaluationTime = now s.Values = map[string]float64{} transitions = append(transitions, StateTransition{ @@ -418,9 +427,15 @@ func (st *Manager) setNextState(ctx context.Context, alertRule *ngModels.AlertRu // Set Resolved property so the scheduler knows to send a postable alert // to Alertmanager. - currentState.Resolved = oldState == eval.Alerting && currentState.State == eval.Normal + newlyResolved := false + if oldState == eval.Alerting && currentState.State == eval.Normal { + currentState.ResolvedAt = &result.EvaluatedAt + newlyResolved = true + } else if currentState.State != eval.Normal && currentState.State != eval.Pending { // Retain the last resolved time for Normal->Normal and Normal->Pending. + currentState.ResolvedAt = nil + } - if shouldTakeImage(currentState.State, oldState, currentState.Image, currentState.Resolved) { + if shouldTakeImage(currentState.State, oldState, currentState.Image, newlyResolved) { image, err := takeImage(ctx, st.images, alertRule) if err != nil { logger.Warn("Failed to take an image", @@ -505,7 +520,7 @@ func (st *Manager) deleteStaleStatesFromCache(ctx context.Context, logger log.Lo s.LastEvaluationTime = evaluatedAt if oldState == eval.Alerting { - s.Resolved = true + s.ResolvedAt = &evaluatedAt image, err := takeImage(ctx, st.images, alertRule) if err != nil { logger.Warn("Failed to take an image", diff --git a/pkg/services/ngalert/state/manager_private_test.go b/pkg/services/ngalert/state/manager_private_test.go index 8740186deac..6f2993d6d88 100644 --- a/pkg/services/ngalert/state/manager_private_test.go +++ b/pkg/services/ngalert/state/manager_private_test.go @@ -563,7 +563,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { StartsAt: t2, EndsAt: t2, LastEvaluationTime: t2, - Resolved: true, + ResolvedAt: &t2, }, }, }, @@ -622,7 +622,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { StartsAt: t1, EndsAt: t3, LastEvaluationTime: t3, - Resolved: true, + ResolvedAt: &t3, }, }, { @@ -1051,7 +1051,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { StartsAt: t1, EndsAt: t3, LastEvaluationTime: t3, - Resolved: true, + ResolvedAt: &t3, }, }, { @@ -1091,7 +1091,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { StartsAt: t1, EndsAt: t3, LastEvaluationTime: t3, - Resolved: true, + ResolvedAt: &t3, }, }, { @@ -1133,7 +1133,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { StartsAt: t1, EndsAt: t3, LastEvaluationTime: t3, - Resolved: true, + ResolvedAt: &t3, }, }, { @@ -1175,7 +1175,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { StartsAt: t1, EndsAt: t3, LastEvaluationTime: t3, - Resolved: true, + ResolvedAt: &t3, }, }, { @@ -1275,7 +1275,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { StartsAt: t2, EndsAt: t2, LastEvaluationTime: t2, - Resolved: true, + ResolvedAt: &t2, }, }, }, @@ -1304,6 +1304,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { StartsAt: t2, EndsAt: t2, LastEvaluationTime: t3, + ResolvedAt: &t2, }, }, }, @@ -1935,7 +1936,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { StartsAt: t1, EndsAt: t3, LastEvaluationTime: t3, - Resolved: true, + ResolvedAt: &t3, }, }, }, @@ -2163,7 +2164,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { StartsAt: t1, EndsAt: t3, LastEvaluationTime: t3, - Resolved: true, + ResolvedAt: &t3, }, }, { @@ -2191,7 +2192,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { StartsAt: t1, EndsAt: t3, LastEvaluationTime: t3, - Resolved: true, + ResolvedAt: &t3, }, }, { @@ -2221,7 +2222,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { StartsAt: t1, EndsAt: t3, LastEvaluationTime: t3, - Resolved: true, + ResolvedAt: &t3, }, }, { @@ -2251,7 +2252,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { StartsAt: t1, EndsAt: t3, LastEvaluationTime: t3, - Resolved: true, + ResolvedAt: &t3, }, }, { @@ -2314,7 +2315,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { StartsAt: t2, EndsAt: t2, LastEvaluationTime: t2, - Resolved: true, + ResolvedAt: &t2, }, }, }, @@ -2330,6 +2331,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { StartsAt: t2, EndsAt: t2, LastEvaluationTime: t3, + ResolvedAt: &t2, }, }, }, @@ -3060,7 +3062,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { StartsAt: t1, EndsAt: t3, LastEvaluationTime: t3, - Resolved: true, + ResolvedAt: &t3, }, }, }, @@ -3480,7 +3482,7 @@ func TestProcessEvalResults_StateTransitions(t *testing.T) { StartsAt: t2, EndsAt: t2, LastEvaluationTime: t2, - Resolved: true, + ResolvedAt: &t2, }, }, }, diff --git a/pkg/services/ngalert/state/manager_test.go b/pkg/services/ngalert/state/manager_test.go index edb836b4dac..b31eab7402d 100644 --- a/pkg/services/ngalert/state/manager_test.go +++ b/pkg/services/ngalert/state/manager_test.go @@ -7,6 +7,7 @@ import ( "fmt" "math" "math/rand" + "net/url" "sort" "strings" "testing" @@ -294,7 +295,7 @@ func TestProcessEvalResults(t *testing.T) { evaluationDuration := 10 * time.Millisecond evaluationInterval := 10 * time.Second - t1 := time.Time{}.Add(evaluationInterval) + t1 := time.Unix(0, 0).Add(evaluationInterval) tn := func(n int) time.Time { return t1.Add(time.Duration(n-1) * evaluationInterval) @@ -424,6 +425,7 @@ func TestProcessEvalResults(t *testing.T) { StartsAt: t1, EndsAt: t1.Add(state.ResendDelay * 4), LastEvaluationTime: t1, + LastSentAt: &t1, }, }, }, @@ -471,6 +473,7 @@ func TestProcessEvalResults(t *testing.T) { StartsAt: t2, EndsAt: t2.Add(state.ResendDelay * 4), LastEvaluationTime: t2, + LastSentAt: &t2, }, }, }, @@ -501,6 +504,94 @@ func TestProcessEvalResults(t *testing.T) { StartsAt: tn(4), EndsAt: tn(4).Add(state.ResendDelay * 4), LastEvaluationTime: tn(4), + LastSentAt: util.Pointer(tn(4)), + }, + }, + }, + { + desc: "alerting -> normal resolves and sets ResolvedAt", + alertRule: baseRule, + evalResults: map[time.Time]eval.Results{ + t1: { + newResult(eval.WithState(eval.Alerting), eval.WithLabels(labels1)), + }, + t2: { + newResult(eval.WithState(eval.Normal), eval.WithLabels(labels1)), + }, + }, + expectedAnnotations: 2, + expectedStates: []*state.State{ + { + Labels: labels["system + rule + labels1"], + ResultFingerprint: labels1.Fingerprint(), + State: eval.Normal, + LatestResult: newEvaluation(t2, eval.Normal), + StartsAt: t2, + EndsAt: t2, + LastEvaluationTime: t2, + ResolvedAt: &t2, + LastSentAt: &t2, + }, + }, + }, + { + desc: "alerting -> normal -> normal resolves and maintains ResolvedAt", + alertRule: baseRule, + evalResults: map[time.Time]eval.Results{ + t1: { + newResult(eval.WithState(eval.Alerting), eval.WithLabels(labels1)), + }, + t2: { + newResult(eval.WithState(eval.Normal), eval.WithLabels(labels1)), + }, + t3: { + newResult(eval.WithState(eval.Normal), eval.WithLabels(labels1)), + }, + }, + expectedAnnotations: 2, + expectedStates: []*state.State{ + { + Labels: labels["system + rule + labels1"], + ResultFingerprint: labels1.Fingerprint(), + State: eval.Normal, + LatestResult: newEvaluation(t3, eval.Normal), + StartsAt: t2, + EndsAt: t2, + LastEvaluationTime: t3, + ResolvedAt: &t2, + LastSentAt: &t2, + }, + }, + }, + { + desc: "pending -> alerting -> normal -> pending resolves and resets ResolvedAt at t4", + alertRule: baseRuleWith(m.WithForNTimes(1)), + evalResults: map[time.Time]eval.Results{ + t1: { + newResult(eval.WithState(eval.Alerting), eval.WithLabels(labels1)), + }, + t2: { + newResult(eval.WithState(eval.Alerting), eval.WithLabels(labels1)), // Alerting. + }, + t3: { + newResult(eval.WithState(eval.Normal), eval.WithLabels(labels1)), + }, + tn(4): { + newResult(eval.WithState(eval.Alerting), eval.WithLabels(labels1)), // Pending. + }, + }, + expectedAnnotations: 4, + expectedStates: []*state.State{ + { + Labels: labels["system + rule + labels1"], + ResultFingerprint: labels1.Fingerprint(), + State: eval.Pending, + LatestResult: newEvaluation(tn(4), eval.Alerting), + StartsAt: tn(4), + EndsAt: tn(4).Add(state.ResendDelay * 4), + LastEvaluationTime: tn(4), + ResolvedAt: &t3, + LastSentAt: &t3, }, }, }, @@ -534,6 +625,7 @@ func TestProcessEvalResults(t *testing.T) { StartsAt: tn(4), EndsAt: tn(4).Add(state.ResendDelay * 4), LastEvaluationTime: tn(5), + LastSentAt: util.Pointer(tn(3)), // 30s resend delay causing the last sent at to be t3. }, }, }, @@ -564,6 +656,7 @@ func TestProcessEvalResults(t *testing.T) { StartsAt: tn(4), EndsAt: tn(4).Add(state.ResendDelay * 4), LastEvaluationTime: tn(4), + LastSentAt: &t3, // Resend delay is 30s, so last sent at is t3. }, }, }, @@ -672,6 +765,7 @@ func TestProcessEvalResults(t *testing.T) { StartsAt: tn(5), EndsAt: tn(5).Add(state.ResendDelay * 4), LastEvaluationTime: tn(5), + LastSentAt: util.Pointer(tn(5)), }, }, }, @@ -696,6 +790,7 @@ func TestProcessEvalResults(t *testing.T) { StartsAt: t2, EndsAt: t2.Add(state.ResendDelay * 4), LastEvaluationTime: t2, + LastSentAt: &t2, }, }, }, @@ -729,6 +824,7 @@ func TestProcessEvalResults(t *testing.T) { StartsAt: t2, EndsAt: t2.Add(state.ResendDelay * 4), LastEvaluationTime: t2, + LastSentAt: &t2, }, }, }, @@ -772,6 +868,7 @@ func TestProcessEvalResults(t *testing.T) { StartsAt: t2, EndsAt: t2.Add(state.ResendDelay * 4), LastEvaluationTime: t2, + LastSentAt: &t2, }, }, }, @@ -808,6 +905,7 @@ func TestProcessEvalResults(t *testing.T) { StartsAt: t2, EndsAt: t2.Add(state.ResendDelay * 4), LastEvaluationTime: t2, + LastSentAt: &t2, }, }, }, @@ -839,6 +937,7 @@ func TestProcessEvalResults(t *testing.T) { StartsAt: t3, EndsAt: tn(4).Add(state.ResendDelay * 4), LastEvaluationTime: tn(4), + LastSentAt: &t3, // Resend delay is 30s, so last sent at is t3. }, }, }, @@ -870,6 +969,7 @@ func TestProcessEvalResults(t *testing.T) { StartsAt: tn(4), EndsAt: tn(4).Add(state.ResendDelay * 4), LastEvaluationTime: tn(4), + LastSentAt: util.Pointer(tn(4)), }, }, }, @@ -956,6 +1056,7 @@ func TestProcessEvalResults(t *testing.T) { StartsAt: tn(5), EndsAt: tn(5).Add(state.ResendDelay * 4), LastEvaluationTime: tn(5), + LastSentAt: util.Pointer(tn(5)), }, }, }, @@ -988,6 +1089,7 @@ func TestProcessEvalResults(t *testing.T) { StartsAt: t2, EndsAt: t2.Add(state.ResendDelay * 4), LastEvaluationTime: t2, + LastSentAt: &t2, EvaluationDuration: evaluationDuration, Annotations: map[string]string{"annotation": "test", "Error": "[sse.dataQueryError] failed to execute query [A]: this is an error"}, }, @@ -1021,6 +1123,7 @@ func TestProcessEvalResults(t *testing.T) { StartsAt: t3, EndsAt: tn(4).Add(state.ResendDelay * 4), LastEvaluationTime: tn(4), + LastSentAt: &t3, // Resend delay is 30s, so last sent at is t3. }, }, }, @@ -1052,6 +1155,7 @@ func TestProcessEvalResults(t *testing.T) { StartsAt: tn(4), EndsAt: tn(4).Add(state.ResendDelay * 4), LastEvaluationTime: tn(4), + LastSentAt: util.Pointer(tn(4)), }, }, }, @@ -1139,6 +1243,7 @@ func TestProcessEvalResults(t *testing.T) { StartsAt: tn(4), EndsAt: tn(6).Add(state.ResendDelay * 4), LastEvaluationTime: tn(6), + LastSentAt: util.Pointer(tn(6)), // After 30s resend delay, last sent at is t6. }, }, }, @@ -1169,6 +1274,7 @@ func TestProcessEvalResults(t *testing.T) { StartsAt: tn(8), EndsAt: tn(8).Add(state.ResendDelay * 4), LastEvaluationTime: tn(8), + LastSentAt: util.Pointer(tn(5)), }, }, }, @@ -1199,6 +1305,7 @@ func TestProcessEvalResults(t *testing.T) { StartsAt: tn(6), EndsAt: tn(6).Add(state.ResendDelay * 4), LastEvaluationTime: tn(6), + LastSentAt: util.Pointer(tn(5)), }, }, }, @@ -1265,6 +1372,7 @@ func TestProcessEvalResults(t *testing.T) { StartsAt: t3, EndsAt: t3.Add(state.ResendDelay * 4), LastEvaluationTime: t3, + LastSentAt: &t1, // Resend delay is 30s, so last sent at is t1. }, }, }, @@ -1306,8 +1414,9 @@ func TestProcessEvalResults(t *testing.T) { res[i].EvaluatedAt = evalTime } clk.Set(evalTime) - _ = st.ProcessEvalResults(context.Background(), evalTime, tc.alertRule, res, systemLabels) + processedStates := st.ProcessEvalResults(context.Background(), evalTime, tc.alertRule, res, systemLabels) results += len(res) + _ = state.FromStateTransitionToPostableAlerts(evalTime, processedStates, st, &url.URL{}) // Set LastSentAt. } states := st.GetStatesForRuleUID(tc.alertRule.OrgID, tc.alertRule.UID) @@ -1670,7 +1779,7 @@ func TestStaleResults(t *testing.T) { assert.Equal(t, models.StateReasonMissingSeries, s.StateReason) assert.Equal(t, clk.Now(), s.EndsAt) if s.CacheID == state2 { - assert.Truef(t, s.Resolved, "Returned stale state should have Resolved set to true") + assert.Equalf(t, clk.Now(), *s.ResolvedAt, "Returned stale state should have ResolvedAt set") } key, err := s.GetAlertInstanceKey() require.NoError(t, err) @@ -1819,11 +1928,11 @@ func TestDeleteStateByRuleUID(t *testing.T) { assert.Equal(t, expectedReason, s.StateReason) if oldState.State == eval.Normal { assert.Equal(t, oldState.StartsAt, s.StartsAt) - assert.False(t, s.Resolved) + assert.Zero(t, s.ResolvedAt) } else { assert.Equal(t, clk.Now(), s.StartsAt) if oldState.State == eval.Alerting { - assert.True(t, s.Resolved) + assert.Equal(t, clk.Now(), *s.ResolvedAt) } } assert.Equal(t, clk.Now(), s.EndsAt) @@ -1959,11 +2068,11 @@ func TestResetStateByRuleUID(t *testing.T) { assert.Equal(t, models.StateReasonPaused, s.StateReason) if oldState.State == eval.Normal { assert.Equal(t, oldState.StartsAt, s.StartsAt) - assert.False(t, s.Resolved) + assert.Zero(t, s.ResolvedAt) } else { assert.Equal(t, clk.Now(), s.StartsAt) if oldState.State == eval.Alerting { - assert.True(t, s.Resolved) + assert.Equal(t, clk.Now(), *s.ResolvedAt) } } assert.Equal(t, clk.Now(), s.EndsAt) diff --git a/pkg/services/ngalert/state/state.go b/pkg/services/ngalert/state/state.go index 206747c652e..260dd92aeef 100644 --- a/pkg/services/ngalert/state/state.go +++ b/pkg/services/ngalert/state/state.go @@ -45,10 +45,6 @@ type State struct { // can still contain the results of previous evaluations. Error error - // Resolved is set to true if this state is the transitional state between Firing and Normal. - // All subsequent states will be false until the next transition from Firing to Normal. - Resolved bool - // Image contains an optional image for the state. It tends to be included in notifications // as a visualization to show why the alert fired. Image *models.Image @@ -65,9 +61,15 @@ type State struct { // conditions. Values map[string]float64 - StartsAt time.Time - EndsAt time.Time - LastSentAt time.Time + StartsAt time.Time + // EndsAt is different from the Prometheus EndsAt as EndsAt is updated for both Normal states + // and states that have been resolved. It cannot be used to determine when a state was resolved. + EndsAt time.Time + // ResolvedAt is set when the state is first resolved. That is to say, when the state first transitions + // from Alerting, NoData, or Error to Normal. It is reset to zero when the state transitions from Normal + // to any other state. + ResolvedAt *time.Time + LastSentAt *time.Time LastEvaluationString string LastEvaluationTime time.Time EvaluationDuration time.Duration @@ -134,14 +136,6 @@ func (a *State) SetNormal(reason string, startsAt, endsAt time.Time) { a.Error = nil } -// Resolve sets the State to Normal. It updates the StateReason, the end time, and sets Resolved to true. -func (a *State) Resolve(reason string, endsAt time.Time) { - a.State = eval.Normal - a.StateReason = reason - a.Resolved = true - a.EndsAt = endsAt -} - // Maintain updates the end time using the most recent evaluation. func (a *State) Maintain(interval int64, evaluatedAt time.Time) { a.EndsAt = nextEndsTime(interval, evaluatedAt) @@ -400,19 +394,31 @@ func resultKeepLast(state *State, rule *models.AlertRule, result eval.Result, lo } } -func (a *State) NeedsSending(resendDelay time.Duration) bool { - switch a.State { - case eval.Pending: - // We do not send notifications for pending states +// NeedsSending returns true if the given state needs to be sent to the Alertmanager. +// Reasons for sending include: +// - The state has been resolved since the last notification. +// - The state is firing and the last notification was sent at least resendDelay ago. +// - The state was resolved within the resolvedRetention period, and the last notification was sent at least resendDelay ago. +func (a *State) NeedsSending(resendDelay time.Duration, resolvedRetention time.Duration) bool { + if a.State == eval.Pending { + // We do not send notifications for pending states. return false - case eval.Normal: - // We should send a notification if the state is Normal because it was resolved - return a.Resolved - default: - // We should send, and re-send notifications, each time LastSentAt is <= LastEvaluationTime + resendDelay - nextSent := a.LastSentAt.Add(resendDelay) - return nextSent.Before(a.LastEvaluationTime) || nextSent.Equal(a.LastEvaluationTime) } + + // We should send a notification if the state has been resolved since the last notification. + if a.ResolvedAt != nil && (a.LastSentAt == nil || a.ResolvedAt.After(*a.LastSentAt)) { + return true + } + + // For normal states, we should only be sending if this is a resolved notification or a re-send of the resolved + // notification within the resolvedRetention period. + if a.State == eval.Normal && (a.ResolvedAt == nil || a.LastEvaluationTime.Sub(*a.ResolvedAt) > resolvedRetention) { + return false + } + + // We should send, and re-send notifications, each time LastSentAt is <= LastEvaluationTime + resendDelay. + // This can include normal->normal transitions that were resolved in recent past evaluations. + return a.LastSentAt == nil || !a.LastSentAt.Add(resendDelay).After(a.LastEvaluationTime) } func (a *State) Equals(b *State) bool { diff --git a/pkg/services/ngalert/state/state_test.go b/pkg/services/ngalert/state/state_test.go index bc00efe875a..f079d90e391 100644 --- a/pkg/services/ngalert/state/state_test.go +++ b/pkg/services/ngalert/state/state_test.go @@ -11,11 +11,12 @@ import ( "github.com/benbjohnson/clock" "github.com/golang/mock/gomock" "github.com/google/uuid" - "github.com/grafana/alerting/models" "github.com/prometheus/common/model" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + "github.com/grafana/alerting/models" + "github.com/grafana/grafana/pkg/infra/log" "github.com/grafana/grafana/pkg/services/ngalert/eval" ngmodels "github.com/grafana/grafana/pkg/services/ngalert/models" @@ -350,10 +351,11 @@ func TestEnd(t *testing.T) { func TestNeedsSending(t *testing.T) { evaluationTime, _ := time.Parse("2006-01-02", "2021-03-25") testCases := []struct { - name string - resendDelay time.Duration - expected bool - testState *State + name string + resendDelay time.Duration + resolvedRetention time.Duration + expected bool + testState *State }{ { name: "state: alerting and LastSentAt before LastEvaluationTime + ResendDelay", @@ -362,7 +364,7 @@ func TestNeedsSending(t *testing.T) { testState: &State{ State: eval.Alerting, LastEvaluationTime: evaluationTime, - LastSentAt: evaluationTime.Add(-2 * time.Minute), + LastSentAt: util.Pointer(evaluationTime.Add(-2 * time.Minute)), }, }, { @@ -372,7 +374,7 @@ func TestNeedsSending(t *testing.T) { testState: &State{ State: eval.Alerting, LastEvaluationTime: evaluationTime, - LastSentAt: evaluationTime, + LastSentAt: util.Pointer(evaluationTime), }, }, { @@ -382,7 +384,7 @@ func TestNeedsSending(t *testing.T) { testState: &State{ State: eval.Alerting, LastEvaluationTime: evaluationTime, - LastSentAt: evaluationTime.Add(-1 * time.Minute), + LastSentAt: util.Pointer(evaluationTime.Add(-1 * time.Minute)), }, }, { @@ -400,18 +402,54 @@ func TestNeedsSending(t *testing.T) { testState: &State{ State: eval.Alerting, LastEvaluationTime: evaluationTime, - LastSentAt: evaluationTime, + LastSentAt: util.Pointer(evaluationTime), }, }, { - name: "state: normal + resolved should send without waiting", + name: "state: normal + resolved should send without waiting if ResolvedAt > LastSentAt", resendDelay: 1 * time.Minute, expected: true, testState: &State{ State: eval.Normal, - Resolved: true, + ResolvedAt: util.Pointer(evaluationTime), LastEvaluationTime: evaluationTime, - LastSentAt: evaluationTime, + LastSentAt: util.Pointer(evaluationTime.Add(-1 * time.Minute)), + }, + }, + { + name: "state: normal + recently resolved should send with wait", + resendDelay: 1 * time.Minute, + resolvedRetention: 15 * time.Minute, + expected: true, + testState: &State{ + State: eval.Normal, + ResolvedAt: util.Pointer(evaluationTime.Add(-2 * time.Minute)), + LastEvaluationTime: evaluationTime, + LastSentAt: util.Pointer(evaluationTime.Add(-1 * time.Minute)), + }, + }, + { + name: "state: normal + recently resolved should not send without wait", + resendDelay: 2 * time.Minute, + resolvedRetention: 15 * time.Minute, + expected: false, + testState: &State{ + State: eval.Normal, + ResolvedAt: util.Pointer(evaluationTime.Add(-2 * time.Minute)), + LastEvaluationTime: evaluationTime, + LastSentAt: util.Pointer(evaluationTime.Add(-1 * time.Minute)), + }, + }, + { + name: "state: normal + not recently resolved should not send even with wait", + resendDelay: 1 * time.Minute, + resolvedRetention: 15 * time.Minute, + expected: false, + testState: &State{ + State: eval.Normal, + ResolvedAt: util.Pointer(evaluationTime.Add(-16 * time.Minute)), + LastEvaluationTime: evaluationTime, + LastSentAt: util.Pointer(evaluationTime.Add(-1 * time.Minute)), }, }, { @@ -420,9 +458,9 @@ func TestNeedsSending(t *testing.T) { expected: false, testState: &State{ State: eval.Normal, - Resolved: false, + ResolvedAt: util.Pointer(time.Time{}), LastEvaluationTime: evaluationTime, - LastSentAt: evaluationTime.Add(-1 * time.Minute), + LastSentAt: util.Pointer(evaluationTime.Add(-1 * time.Minute)), }, }, { @@ -432,7 +470,7 @@ func TestNeedsSending(t *testing.T) { testState: &State{ State: eval.NoData, LastEvaluationTime: evaluationTime, - LastSentAt: evaluationTime.Add(-1 * time.Minute), + LastSentAt: util.Pointer(evaluationTime.Add(-1 * time.Minute)), }, }, { @@ -442,7 +480,7 @@ func TestNeedsSending(t *testing.T) { testState: &State{ State: eval.NoData, LastEvaluationTime: evaluationTime, - LastSentAt: evaluationTime.Add(-time.Duration(rand.Int63n(59)+1) * time.Second), + LastSentAt: util.Pointer(evaluationTime.Add(-time.Duration(rand.Int63n(59)+1) * time.Second)), }, }, { @@ -452,7 +490,7 @@ func TestNeedsSending(t *testing.T) { testState: &State{ State: eval.Error, LastEvaluationTime: evaluationTime, - LastSentAt: evaluationTime.Add(-1 * time.Minute), + LastSentAt: util.Pointer(evaluationTime.Add(-1 * time.Minute)), }, }, { @@ -462,14 +500,14 @@ func TestNeedsSending(t *testing.T) { testState: &State{ State: eval.Error, LastEvaluationTime: evaluationTime, - LastSentAt: evaluationTime.Add(-time.Duration(rand.Int63n(59)+1) * time.Second), + LastSentAt: util.Pointer(evaluationTime.Add(-time.Duration(rand.Int63n(59)+1) * time.Second)), }, }, } for _, tc := range testCases { t.Run(tc.name, func(t *testing.T) { - assert.Equal(t, tc.expected, tc.testState.NeedsSending(tc.resendDelay)) + assert.Equal(t, tc.expected, tc.testState.NeedsSending(tc.resendDelay, tc.resolvedRetention)) }) } } @@ -531,13 +569,6 @@ func TestGetLastEvaluationValuesForCondition(t *testing.T) { }) } -func TestResolve(t *testing.T) { - s := State{State: eval.Alerting, EndsAt: time.Now().Add(time.Minute)} - expected := State{State: eval.Normal, StateReason: "This is a reason", EndsAt: time.Now(), Resolved: true} - s.Resolve("This is a reason", expected.EndsAt) - assert.Equal(t, expected, s) -} - func TestShouldTakeImage(t *testing.T) { tests := []struct { name string diff --git a/pkg/setting/setting_unified_alerting.go b/pkg/setting/setting_unified_alerting.go index 4cd4e9dde85..1d205dc6faa 100644 --- a/pkg/setting/setting_unified_alerting.go +++ b/pkg/setting/setting_unified_alerting.go @@ -6,11 +6,12 @@ import ( "strings" "time" - alertingCluster "github.com/grafana/alerting/cluster" dstls "github.com/grafana/dskit/crypto/tls" "github.com/grafana/grafana-plugin-sdk-go/backend/gtime" "gopkg.in/ini.v1" + alertingCluster "github.com/grafana/alerting/cluster" + "github.com/grafana/grafana/pkg/util" ) @@ -113,6 +114,9 @@ type UnifiedAlertingSettings struct { // Retention period for Alertmanager notification log entries. NotificationLogRetention time.Duration + + // Duration for which a resolved alert state transition will continue to be sent to the Alertmanager. + ResolvedAlertRetention time.Duration } type RecordingRuleSettings struct { @@ -435,6 +439,11 @@ func (cfg *Cfg) ReadUnifiedAlertingSettings(iniFile *ini.File) error { return err } + uaCfg.ResolvedAlertRetention, err = gtime.ParseDuration(valueAsString(ua, "resolved_alert_retention", (15 * time.Minute).String())) + if err != nil { + return err + } + cfg.UnifiedAlerting = uaCfg return nil }