Alerting: Initialize rule routine with initial alert rule fingerprint (#114979)
Alerting: Initialize rule routine with initial fingerprint
This commit is contained in:
@@ -47,10 +47,10 @@ type Rule interface {
|
||||
Identifier() ngmodels.AlertRuleKeyWithGroup
|
||||
}
|
||||
|
||||
type ruleFactoryFunc func(context.Context, *ngmodels.AlertRule) Rule
|
||||
type ruleFactoryFunc func(context.Context, ruleWithFolder) Rule
|
||||
|
||||
func (f ruleFactoryFunc) new(ctx context.Context, rule *ngmodels.AlertRule) Rule {
|
||||
return f(ctx, rule)
|
||||
func (f ruleFactoryFunc) new(ctx context.Context, rf ruleWithFolder) Rule {
|
||||
return f(ctx, rf)
|
||||
}
|
||||
|
||||
func newRuleFactory(
|
||||
@@ -70,11 +70,11 @@ func newRuleFactory(
|
||||
evalAppliedHook evalAppliedFunc,
|
||||
stopAppliedHook stopAppliedFunc,
|
||||
) ruleFactoryFunc {
|
||||
return func(ctx context.Context, rule *ngmodels.AlertRule) Rule {
|
||||
if rule.Type() == ngmodels.RuleTypeRecording {
|
||||
return func(ctx context.Context, rf ruleWithFolder) Rule {
|
||||
if rf.rule.Type() == ngmodels.RuleTypeRecording {
|
||||
return newRecordingRule(
|
||||
ctx,
|
||||
rule.GetKeyWithGroup(),
|
||||
rf.rule.GetKeyWithGroup(),
|
||||
retryConfig,
|
||||
clock,
|
||||
evalFactory,
|
||||
@@ -89,7 +89,7 @@ func newRuleFactory(
|
||||
}
|
||||
return newAlertRule(
|
||||
ctx,
|
||||
rule.GetKeyWithGroup(),
|
||||
rf,
|
||||
appURL,
|
||||
disableGrafanaFolder,
|
||||
retryConfig,
|
||||
@@ -111,7 +111,8 @@ type evalAppliedFunc = func(ngmodels.AlertRuleKey, time.Time)
|
||||
type stopAppliedFunc = func(ngmodels.AlertRuleKey)
|
||||
|
||||
type alertRule struct {
|
||||
key ngmodels.AlertRuleKeyWithGroup
|
||||
key ngmodels.AlertRuleKeyWithGroup
|
||||
currentFingerprint fingerprint
|
||||
|
||||
evalCh chan *Evaluation
|
||||
updateCh chan *Evaluation
|
||||
@@ -139,7 +140,7 @@ type alertRule struct {
|
||||
|
||||
func newAlertRule(
|
||||
parent context.Context,
|
||||
key ngmodels.AlertRuleKeyWithGroup,
|
||||
rf ruleWithFolder,
|
||||
appURL *url.URL,
|
||||
disableGrafanaFolder bool,
|
||||
retryConfig RetryConfig,
|
||||
@@ -154,10 +155,13 @@ func newAlertRule(
|
||||
evalAppliedHook func(ngmodels.AlertRuleKey, time.Time),
|
||||
stopAppliedHook func(ngmodels.AlertRuleKey),
|
||||
) *alertRule {
|
||||
key := rf.rule.GetKeyWithGroup()
|
||||
initialFingerprint := rf.Fingerprint()
|
||||
ctx, stop := util.WithCancelCause(ngmodels.WithRuleKey(parent, key.AlertRuleKey))
|
||||
|
||||
return &alertRule{
|
||||
a := &alertRule{
|
||||
key: key,
|
||||
currentFingerprint: initialFingerprint,
|
||||
evalCh: make(chan *Evaluation),
|
||||
updateCh: make(chan *Evaluation),
|
||||
ctx: ctx,
|
||||
@@ -176,6 +180,8 @@ func newAlertRule(
|
||||
tracer: tracer,
|
||||
featureToggles: featureToggles,
|
||||
}
|
||||
|
||||
return a
|
||||
}
|
||||
|
||||
func (a *alertRule) Identifier() ngmodels.AlertRuleKeyWithGroup {
|
||||
@@ -246,22 +252,23 @@ func (a *alertRule) Run() error {
|
||||
grafanaCtx := a.ctx
|
||||
a.logger.Debug("Alert rule routine started")
|
||||
|
||||
var currentFingerprint fingerprint
|
||||
firstEvalDone := false
|
||||
|
||||
defer a.stopApplied()
|
||||
for {
|
||||
select {
|
||||
// used by external services (API) to notify that rule is updated.
|
||||
case ctx := <-a.updateCh:
|
||||
fp := ctx.Fingerprint()
|
||||
if currentFingerprint == fp {
|
||||
a.logger.Info("Rule's fingerprint has not changed. Skip resetting the state", "currentFingerprint", currentFingerprint)
|
||||
if a.currentFingerprint == fp {
|
||||
a.logger.Info("Rule's fingerprint has not changed. Skip resetting the state", "currentFingerprint", a.currentFingerprint)
|
||||
continue
|
||||
}
|
||||
|
||||
a.logger.Info("Clearing the state of the rule because it was updated", "isPaused", ctx.rule.IsPaused, "fingerprint", fp)
|
||||
// clear the state. So the next evaluation will start from the scratch.
|
||||
a.resetState(grafanaCtx, ctx.rule, ctx.rule.IsPaused)
|
||||
currentFingerprint = fp
|
||||
a.currentFingerprint = fp
|
||||
// evalCh - used by the scheduler to signal that evaluation is needed.
|
||||
case ctx, ok := <-a.evalCh:
|
||||
if !ok {
|
||||
@@ -295,20 +302,21 @@ func (a *alertRule) Run() error {
|
||||
for {
|
||||
isPaused := ctx.rule.IsPaused
|
||||
|
||||
// Do not clean up state if the eval loop has just started.
|
||||
var needReset bool
|
||||
if currentFingerprint != 0 && currentFingerprint != f {
|
||||
logger.Debug("Got a new version of alert rule. Clear up the state", "current_fingerprint", currentFingerprint, "fingerprint", f)
|
||||
if a.currentFingerprint != f {
|
||||
logger.Debug("Got a new version of alert rule. Clear up the state", "current_fingerprint", a.currentFingerprint, "fingerprint", f)
|
||||
needReset = true
|
||||
}
|
||||
// We need to reset state if the loop has started and the alert is already paused. It can happen,
|
||||
// if we have an alert with state and we do file provision with stateful Grafana, that state
|
||||
// lingers in DB and won't be cleaned up until next alert rule update.
|
||||
needReset = needReset || (currentFingerprint == 0 && isPaused)
|
||||
needReset = needReset || (!firstEvalDone && isPaused)
|
||||
if needReset {
|
||||
a.resetState(grafanaCtx, ctx.rule, isPaused)
|
||||
}
|
||||
currentFingerprint = f
|
||||
|
||||
firstEvalDone = true
|
||||
a.currentFingerprint = f
|
||||
if isPaused {
|
||||
logger.Debug("Skip rule evaluation because it is paused")
|
||||
return
|
||||
@@ -319,7 +327,7 @@ func (a *alertRule) Run() error {
|
||||
evalTotal.Inc()
|
||||
}
|
||||
|
||||
fpStr := currentFingerprint.String()
|
||||
fpStr := a.currentFingerprint.String()
|
||||
utcTick := ctx.scheduledAt.UTC().Format(time.RFC3339Nano)
|
||||
tracingCtx, span := a.tracer.Start(grafanaCtx, "alert rule execution", trace.WithAttributes(
|
||||
attribute.String("rule_uid", ctx.rule.UID),
|
||||
|
||||
@@ -334,7 +334,7 @@ func TestAlertRuleAfterEval(t *testing.T) {
|
||||
ruleStore.PutRule(context.Background(), rule)
|
||||
ruleFactory := ruleFactoryFromScheduler(sch)
|
||||
|
||||
process := ruleFactory.new(context.Background(), rule)
|
||||
process := ruleFactory.new(context.Background(), ruleWithFolder{rule: rule, folderTitle: ""})
|
||||
|
||||
return &testContext{
|
||||
rule: rule,
|
||||
@@ -503,7 +503,14 @@ func blankRuleForTests(ctx context.Context, key models.AlertRuleKeyWithGroup) *a
|
||||
Log: log.NewNopLogger(),
|
||||
}
|
||||
st := state.NewManager(managerCfg, state.NewNoopPersister())
|
||||
return newAlertRule(ctx, key, nil, false, RetryConfig{}, nil, st, nil, nil, nil, log.NewNopLogger(), nil, featuremgmt.WithFeatures(), nil, nil)
|
||||
// Create a minimal rule from the key
|
||||
rule := &models.AlertRule{
|
||||
OrgID: key.OrgID,
|
||||
UID: key.UID,
|
||||
RuleGroup: key.RuleGroup,
|
||||
}
|
||||
rf := ruleWithFolder{rule: rule, folderTitle: ""}
|
||||
return newAlertRule(ctx, rf, nil, false, RetryConfig{}, nil, st, nil, nil, nil, log.NewNopLogger(), nil, featuremgmt.WithFeatures(), nil, nil)
|
||||
}
|
||||
|
||||
func TestRuleRoutine(t *testing.T) {
|
||||
@@ -540,7 +547,7 @@ func TestRuleRoutine(t *testing.T) {
|
||||
factory := ruleFactoryFromScheduler(sch)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
t.Cleanup(cancel)
|
||||
ruleInfo := factory.new(ctx, rule)
|
||||
ruleInfo := factory.new(ctx, ruleWithFolder{rule: rule, folderTitle: folderTitle})
|
||||
go func() {
|
||||
_ = ruleInfo.Run()
|
||||
}()
|
||||
@@ -728,7 +735,8 @@ func TestRuleRoutine(t *testing.T) {
|
||||
|
||||
factory := ruleFactoryFromScheduler(sch)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
ruleInfo := factory.new(ctx, rule)
|
||||
folderTitle := ""
|
||||
ruleInfo := factory.new(ctx, ruleWithFolder{rule: rule, folderTitle: folderTitle})
|
||||
go func() {
|
||||
err := ruleInfo.Run()
|
||||
stoppedChan <- err
|
||||
@@ -750,7 +758,7 @@ func TestRuleRoutine(t *testing.T) {
|
||||
require.NotEmpty(t, sch.stateManager.GetStatesForRuleUID(rule.OrgID, rule.UID))
|
||||
|
||||
factory := ruleFactoryFromScheduler(sch)
|
||||
ruleInfo := factory.new(context.Background(), rule)
|
||||
ruleInfo := factory.new(context.Background(), ruleWithFolder{rule: rule, folderTitle: ""})
|
||||
go func() {
|
||||
err := ruleInfo.Run()
|
||||
stoppedChan <- err
|
||||
@@ -774,7 +782,7 @@ func TestRuleRoutine(t *testing.T) {
|
||||
require.NotEmpty(t, sch.stateManager.GetStatesForRuleUID(rule.OrgID, rule.UID))
|
||||
|
||||
factory := ruleFactoryFromScheduler(sch)
|
||||
ruleInfo := factory.new(context.Background(), rule)
|
||||
ruleInfo := factory.new(context.Background(), ruleWithFolder{rule: rule, folderTitle: ""})
|
||||
go func() {
|
||||
err := ruleInfo.Run()
|
||||
stoppedChan <- err
|
||||
@@ -804,7 +812,7 @@ func TestRuleRoutine(t *testing.T) {
|
||||
factory := ruleFactoryFromScheduler(sch)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
t.Cleanup(cancel)
|
||||
ruleInfo := factory.new(ctx, rule)
|
||||
ruleInfo := factory.new(ctx, ruleWithFolder{rule: rule, folderTitle: folderTitle})
|
||||
|
||||
go func() {
|
||||
_ = ruleInfo.Run()
|
||||
@@ -871,6 +879,140 @@ func TestRuleRoutine(t *testing.T) {
|
||||
})
|
||||
})
|
||||
|
||||
t.Run("when update is sent before first evaluation", func(t *testing.T) {
|
||||
rule := gen.With(withQueryForState(t, eval.Normal)).GenerateRef()
|
||||
folderTitle := "folderName"
|
||||
|
||||
evalAppliedChan := make(chan time.Time)
|
||||
|
||||
sender := NewSyncAlertsSenderMock()
|
||||
sender.EXPECT().Send(mock.Anything, rule.GetKey(), mock.Anything).Return()
|
||||
|
||||
sch, ruleStore, _, _ := createSchedule(evalAppliedChan, sender, clock.NewMock())
|
||||
ruleStore.PutRule(context.Background(), rule)
|
||||
sch.schedulableAlertRules.set([]*models.AlertRule{rule}, map[models.FolderKey]string{rule.GetFolderKey(): folderTitle})
|
||||
|
||||
// Add state to verify it's not cleared
|
||||
states := []*state.State{
|
||||
{
|
||||
AlertRuleUID: rule.UID,
|
||||
CacheID: data.Labels(rule.Labels).Fingerprint(),
|
||||
OrgID: rule.OrgID,
|
||||
State: eval.Alerting,
|
||||
StartsAt: sch.clock.Now(),
|
||||
EndsAt: sch.clock.Now().Add(5 * time.Second),
|
||||
Labels: rule.Labels,
|
||||
},
|
||||
}
|
||||
sch.stateManager.Put(states)
|
||||
|
||||
t.Run("should not reset state if fingerprint is the same", func(t *testing.T) {
|
||||
factory := ruleFactoryFromScheduler(sch)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
t.Cleanup(cancel)
|
||||
ruleInfo := factory.new(ctx, ruleWithFolder{rule: rule, folderTitle: folderTitle})
|
||||
|
||||
go func() {
|
||||
_ = ruleInfo.Run()
|
||||
}()
|
||||
|
||||
// Send update before first evaluation - same rule, same fingerprint
|
||||
// This should not reset state since fingerprint is the same
|
||||
ruleInfo.Update(&Evaluation{rule: rule, folderTitle: folderTitle})
|
||||
|
||||
// Give time for update to be processed
|
||||
time.Sleep(100 * time.Millisecond)
|
||||
|
||||
actualStates := sch.stateManager.GetStatesForRuleUID(rule.OrgID, rule.UID)
|
||||
require.NotEmpty(t, actualStates)
|
||||
})
|
||||
|
||||
t.Run("should reset state if fingerprint is different", func(t *testing.T) {
|
||||
// Re-add state for this test
|
||||
sch.stateManager.Put(states)
|
||||
|
||||
sender := NewSyncAlertsSenderMock()
|
||||
sender.EXPECT().Send(mock.Anything, rule.GetKey(), mock.Anything).Return()
|
||||
|
||||
sch2, ruleStore2, _, _ := createSchedule(make(chan time.Time), sender, clock.NewMock())
|
||||
ruleStore2.PutRule(context.Background(), rule)
|
||||
sch2.schedulableAlertRules.set([]*models.AlertRule{rule}, map[models.FolderKey]string{rule.GetFolderKey(): folderTitle})
|
||||
sch2.stateManager.Put(states)
|
||||
|
||||
factory := ruleFactoryFromScheduler(sch2)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
t.Cleanup(cancel)
|
||||
ruleInfo := factory.new(ctx, ruleWithFolder{rule: rule, folderTitle: folderTitle})
|
||||
|
||||
go func() {
|
||||
_ = ruleInfo.Run()
|
||||
}()
|
||||
|
||||
// Send update before first eval, with a changed alert rule title.
|
||||
// This should reset state and send resolved alerts
|
||||
updatedRule := models.CopyRule(rule, gen.WithTitle(util.GenerateShortUID()))
|
||||
ruleInfo.Update(&Evaluation{rule: updatedRule, folderTitle: folderTitle})
|
||||
|
||||
// Wait for sender to be called (which happens when state is cleared and resolved alerts are sent)
|
||||
require.Eventually(t, func() bool {
|
||||
return len(sender.Calls()) > 0
|
||||
}, 5*time.Second, 100*time.Millisecond)
|
||||
|
||||
// State should be cleared because fingerprint changed
|
||||
actualStates := sch2.stateManager.GetStatesForRuleUID(rule.OrgID, rule.UID)
|
||||
require.Empty(t, actualStates)
|
||||
})
|
||||
})
|
||||
|
||||
t.Run("paused rule should reset state on first evaluation", func(t *testing.T) {
|
||||
rule := gen.With(withQueryForState(t, eval.Normal)).GenerateRef()
|
||||
rule.IsPaused = true
|
||||
folderTitle := "folderName"
|
||||
|
||||
sender := NewSyncAlertsSenderMock()
|
||||
sender.EXPECT().Send(mock.Anything, rule.GetKey(), mock.Anything).Return()
|
||||
|
||||
sch, ruleStore, _, _ := createSchedule(make(chan time.Time), sender, clock.NewMock())
|
||||
ruleStore.PutRule(context.Background(), rule)
|
||||
sch.schedulableAlertRules.set([]*models.AlertRule{rule}, map[models.FolderKey]string{rule.GetFolderKey(): folderTitle})
|
||||
|
||||
states := []*state.State{
|
||||
{
|
||||
AlertRuleUID: rule.UID,
|
||||
CacheID: data.Labels(rule.Labels).Fingerprint(),
|
||||
OrgID: rule.OrgID,
|
||||
State: eval.Alerting,
|
||||
StartsAt: sch.clock.Now(),
|
||||
EndsAt: sch.clock.Now().Add(5 * time.Second),
|
||||
Labels: rule.Labels,
|
||||
},
|
||||
}
|
||||
sch.stateManager.Put(states)
|
||||
require.NotEmpty(t, sch.stateManager.GetStatesForRuleUID(rule.OrgID, rule.UID))
|
||||
|
||||
factory := ruleFactoryFromScheduler(sch)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
t.Cleanup(cancel)
|
||||
ruleInfo := factory.new(ctx, ruleWithFolder{rule: rule, folderTitle: folderTitle})
|
||||
|
||||
go func() {
|
||||
_ = ruleInfo.Run()
|
||||
}()
|
||||
|
||||
ruleInfo.Eval(&Evaluation{
|
||||
scheduledAt: sch.clock.Now(),
|
||||
rule: rule,
|
||||
folderTitle: folderTitle,
|
||||
})
|
||||
|
||||
require.Eventually(t, func() bool {
|
||||
return len(sender.Calls()) > 0
|
||||
}, 5*time.Second, 100*time.Millisecond)
|
||||
|
||||
actualStates := sch.stateManager.GetStatesForRuleUID(rule.OrgID, rule.UID)
|
||||
require.Empty(t, actualStates)
|
||||
})
|
||||
|
||||
t.Run("when evaluation fails", func(t *testing.T) {
|
||||
rule := gen.With(withQueryForState(t, eval.Error)).GenerateRef()
|
||||
rule.ExecErrState = models.ErrorErrState
|
||||
@@ -910,7 +1052,7 @@ func TestRuleRoutine(t *testing.T) {
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
t.Cleanup(cancel)
|
||||
ruleInfo := factory.new(ctx, rule)
|
||||
ruleInfo := factory.new(ctx, ruleWithFolder{rule: rule, folderTitle: ""})
|
||||
|
||||
go func() {
|
||||
_ = ruleInfo.Run()
|
||||
@@ -1044,7 +1186,7 @@ func TestRuleRoutine(t *testing.T) {
|
||||
factory := ruleFactoryFromScheduler(sch)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
t.Cleanup(cancel)
|
||||
ruleInfo := factory.new(ctx, rule)
|
||||
ruleInfo := factory.new(ctx, ruleWithFolder{rule: rule, folderTitle: ""})
|
||||
|
||||
go func() {
|
||||
_ = ruleInfo.Run()
|
||||
@@ -1078,7 +1220,7 @@ func TestRuleRoutine(t *testing.T) {
|
||||
factory := ruleFactoryFromScheduler(sch)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
t.Cleanup(cancel)
|
||||
ruleInfo := factory.new(ctx, rule)
|
||||
ruleInfo := factory.new(ctx, ruleWithFolder{rule: rule, folderTitle: ""})
|
||||
|
||||
go func() {
|
||||
_ = ruleInfo.Run()
|
||||
@@ -1119,7 +1261,7 @@ func TestRuleRoutine(t *testing.T) {
|
||||
factory := ruleFactoryFromScheduler(sch)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
t.Cleanup(cancel)
|
||||
ruleInfo := factory.new(ctx, rule)
|
||||
ruleInfo := factory.new(ctx, ruleWithFolder{rule: rule, folderTitle: ""})
|
||||
|
||||
go func() {
|
||||
_ = ruleInfo.Run()
|
||||
@@ -1214,7 +1356,7 @@ func TestAlertRuleRetry(t *testing.T) {
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
t.Cleanup(cancel)
|
||||
ruleInfo := factory.new(ctx, rule)
|
||||
ruleInfo := factory.new(ctx, ruleWithFolder{rule: rule, folderTitle: ""})
|
||||
|
||||
go func() {
|
||||
_ = ruleInfo.Run()
|
||||
|
||||
@@ -239,7 +239,7 @@ func TestRecordingRuleAfterEval(t *testing.T) {
|
||||
ruleStore.PutRule(context.Background(), rule)
|
||||
ruleFactory := ruleFactoryFromScheduler(sch)
|
||||
|
||||
process := ruleFactory.new(context.Background(), rule)
|
||||
process := ruleFactory.new(context.Background(), ruleWithFolder{rule: rule, folderTitle: ""})
|
||||
|
||||
evalDoneChan := make(chan time.Time, 1) // Buffer to avoid blocking
|
||||
afterEvalCh := make(chan struct{}, 1) // Buffer to avoid blocking
|
||||
@@ -444,7 +444,7 @@ func testRecordingRule_Integration(t *testing.T, writeTarget *writer.TestRemoteW
|
||||
folderTitle := ruleStore.getNamespaceTitle(rule.NamespaceUID)
|
||||
ruleFactory := ruleFactoryFromScheduler(sch)
|
||||
|
||||
process := ruleFactory.new(context.Background(), rule)
|
||||
process := ruleFactory.new(context.Background(), ruleWithFolder{rule: rule, folderTitle: ""})
|
||||
evalDoneChan := make(chan time.Time)
|
||||
process.(*recordingRule).evalAppliedHook = func(_ models.AlertRuleKey, t time.Time) {
|
||||
evalDoneChan <- t
|
||||
@@ -583,7 +583,7 @@ func testRecordingRule_Integration(t *testing.T, writeTarget *writer.TestRemoteW
|
||||
folderTitle := ruleStore.getNamespaceTitle(rule.NamespaceUID)
|
||||
ruleFactory := ruleFactoryFromScheduler(sch)
|
||||
|
||||
process := ruleFactory.new(context.Background(), rule)
|
||||
process := ruleFactory.new(context.Background(), ruleWithFolder{rule: rule, folderTitle: ""})
|
||||
evalDoneChan := make(chan time.Time)
|
||||
process.(*recordingRule).evalAppliedHook = func(_ models.AlertRuleKey, t time.Time) {
|
||||
evalDoneChan <- t
|
||||
@@ -728,7 +728,7 @@ func testRecordingRule_Integration(t *testing.T, writeTarget *writer.TestRemoteW
|
||||
folderTitle := ruleStore.getNamespaceTitle(rule.NamespaceUID)
|
||||
ruleFactory := ruleFactoryFromScheduler(sch)
|
||||
|
||||
process := ruleFactory.new(context.Background(), rule)
|
||||
process := ruleFactory.new(context.Background(), ruleWithFolder{rule: rule, folderTitle: ""})
|
||||
evalDoneChan := make(chan time.Time)
|
||||
process.(*recordingRule).evalAppliedHook = func(_ models.AlertRuleKey, t time.Time) {
|
||||
evalDoneChan <- t
|
||||
@@ -787,7 +787,7 @@ func testRecordingRule_Integration(t *testing.T, writeTarget *writer.TestRemoteW
|
||||
folderTitle := ruleStore.getNamespaceTitle(rule.NamespaceUID)
|
||||
ruleFactory := ruleFactoryFromScheduler(sch)
|
||||
|
||||
process := ruleFactory.new(context.Background(), rule)
|
||||
process := ruleFactory.new(context.Background(), ruleWithFolder{rule: rule, folderTitle: ""})
|
||||
evalDoneChan := make(chan time.Time)
|
||||
process.(*recordingRule).evalAppliedHook = func(_ models.AlertRuleKey, t time.Time) {
|
||||
evalDoneChan <- t
|
||||
|
||||
@@ -21,7 +21,7 @@ var (
|
||||
)
|
||||
|
||||
type ruleFactory interface {
|
||||
new(context.Context, *models.AlertRule) Rule
|
||||
new(context.Context, ruleWithFolder) Rule
|
||||
}
|
||||
|
||||
type ruleRegistry struct {
|
||||
@@ -35,14 +35,14 @@ func newRuleRegistry() ruleRegistry {
|
||||
|
||||
// getOrCreate gets a rule routine from registry for the provided rule. If it does not exist, it creates a new one.
|
||||
// Returns a pointer to the rule routine and a flag that indicates whether it is a new struct or not.
|
||||
func (r *ruleRegistry) getOrCreate(context context.Context, item *models.AlertRule, factory ruleFactory) (Rule, bool) {
|
||||
func (r *ruleRegistry) getOrCreate(context context.Context, rf ruleWithFolder, factory ruleFactory) (Rule, bool) {
|
||||
r.mu.Lock()
|
||||
defer r.mu.Unlock()
|
||||
|
||||
key := item.GetKey()
|
||||
key := rf.rule.GetKey()
|
||||
rule, ok := r.rules[key]
|
||||
if !ok {
|
||||
rule = factory.new(context, item)
|
||||
rule = factory.new(context, rf)
|
||||
r.rules[key] = rule
|
||||
}
|
||||
return rule, !ok
|
||||
|
||||
@@ -320,10 +320,22 @@ func (sch *schedule) processTick(ctx context.Context, dispatcherGroup *errgroup.
|
||||
sch.stopAppliedFunc,
|
||||
)
|
||||
for _, item := range alertRules {
|
||||
ruleRoutine, newRoutine := sch.registry.getOrCreate(ctx, item, ruleFactory)
|
||||
key := item.GetKey()
|
||||
logger := sch.log.FromContext(ctx).New(key.LogContext()...)
|
||||
|
||||
var folderTitle string
|
||||
if !sch.disableGrafanaFolder {
|
||||
title, ok := folderTitles[item.GetFolderKey()]
|
||||
if ok {
|
||||
folderTitle = title
|
||||
} else {
|
||||
missingFolder[item.NamespaceUID] = append(missingFolder[item.NamespaceUID], item.UID)
|
||||
}
|
||||
}
|
||||
|
||||
rf := ruleWithFolder{rule: item, folderTitle: folderTitle}
|
||||
ruleRoutine, newRoutine := sch.registry.getOrCreate(ctx, rf, ruleFactory)
|
||||
|
||||
// enforce minimum evaluation interval
|
||||
if item.IntervalSeconds < int64(sch.minRuleInterval.Seconds()) {
|
||||
logger.Debug("Interval adjusted", "originalInterval", item.IntervalSeconds, "adjustedInterval", sch.minRuleInterval.Seconds())
|
||||
@@ -337,7 +349,7 @@ func (sch *schedule) processTick(ctx context.Context, dispatcherGroup *errgroup.
|
||||
logger.Debug("Rule restarted because type changed", "old", ruleRoutine.Type(), "new", item.Type())
|
||||
restartedRules = append(restartedRules, ruleRoutine)
|
||||
sch.registry.del(key)
|
||||
ruleRoutine, newRoutine = sch.registry.getOrCreate(ctx, item, ruleFactory)
|
||||
ruleRoutine, newRoutine = sch.registry.getOrCreate(ctx, rf, ruleFactory)
|
||||
}
|
||||
|
||||
if newRoutine && !invalidInterval {
|
||||
@@ -357,16 +369,6 @@ func (sch *schedule) processTick(ctx context.Context, dispatcherGroup *errgroup.
|
||||
offset := jitterOffsetInTicks(item, sch.baseInterval, sch.jitterEvaluations)
|
||||
isReadyToRun := item.IntervalSeconds != 0 && (tickNum%itemFrequency)-offset == 0
|
||||
|
||||
var folderTitle string
|
||||
if !sch.disableGrafanaFolder {
|
||||
title, ok := folderTitles[item.GetFolderKey()]
|
||||
if ok {
|
||||
folderTitle = title
|
||||
} else {
|
||||
missingFolder[item.NamespaceUID] = append(missingFolder[item.NamespaceUID], item.UID)
|
||||
}
|
||||
}
|
||||
|
||||
if isReadyToRun {
|
||||
logger.Debug("Rule is ready to run on the current tick", "tick", tick, "frequency", itemFrequency, "offset", offset)
|
||||
readyToRun = append(readyToRun, readyToRunItem{ruleRoutine: ruleRoutine, Evaluation: Evaluation{
|
||||
@@ -378,12 +380,12 @@ func (sch *schedule) processTick(ctx context.Context, dispatcherGroup *errgroup.
|
||||
if _, isUpdated := updated[key]; isUpdated && !isReadyToRun {
|
||||
// if we do not need to eval the rule, check the whether rule was just updated and if it was, notify evaluation routine about that
|
||||
logger.Debug("Rule has been updated. Notifying evaluation routine")
|
||||
go func(routine Rule, rule *ngmodels.AlertRule) {
|
||||
go func(routine Rule, rule *ngmodels.AlertRule, folder string) {
|
||||
routine.Update(&Evaluation{
|
||||
rule: rule,
|
||||
folderTitle: folderTitle,
|
||||
folderTitle: folder,
|
||||
})
|
||||
}(ruleRoutine, item)
|
||||
}(ruleRoutine, item, folderTitle)
|
||||
updatedRules = append(updatedRules, ngmodels.AlertRuleKeyWithVersion{
|
||||
Version: item.Version,
|
||||
AlertRuleKey: item.GetKey(),
|
||||
|
||||
@@ -1107,7 +1107,7 @@ func TestSchedule_deleteAlertRule(t *testing.T) {
|
||||
rule := models.RuleGen.GenerateRef()
|
||||
ruleStore.PutRule(ctx, rule)
|
||||
key := rule.GetKey()
|
||||
info, _ := sch.registry.getOrCreate(ctx, rule, ruleFactory)
|
||||
info, _ := sch.registry.getOrCreate(ctx, ruleWithFolder{rule: rule, folderTitle: ""}, ruleFactory)
|
||||
|
||||
sch.deleteAlertRule(ctx, key)
|
||||
|
||||
@@ -1126,7 +1126,7 @@ func TestSchedule_deleteAlertRule(t *testing.T) {
|
||||
rule := models.RuleGen.GenerateRef()
|
||||
ruleStore.PutRule(ctx, rule)
|
||||
key := rule.GetKey()
|
||||
info, _ := sch.registry.getOrCreate(ctx, rule, ruleFactory)
|
||||
info, _ := sch.registry.getOrCreate(ctx, ruleWithFolder{rule: rule, folderTitle: ""}, ruleFactory)
|
||||
|
||||
_, err := sch.updateSchedulableAlertRules(ctx)
|
||||
require.NoError(t, err)
|
||||
@@ -1149,7 +1149,7 @@ func TestSchedule_deleteAlertRule(t *testing.T) {
|
||||
rule := models.RuleGen.GenerateRef()
|
||||
ruleStore.PutRule(ctx, rule)
|
||||
key := rule.GetKey()
|
||||
info, _ := sch.registry.getOrCreate(ctx, rule, ruleFactory)
|
||||
info, _ := sch.registry.getOrCreate(ctx, ruleWithFolder{rule: rule, folderTitle: ""}, ruleFactory)
|
||||
|
||||
_, err := sch.updateSchedulableAlertRules(ctx)
|
||||
require.NoError(t, err)
|
||||
@@ -1172,7 +1172,7 @@ func TestSchedule_deleteAlertRule(t *testing.T) {
|
||||
ruleFactory := ruleFactoryFromScheduler(sch)
|
||||
rule := models.RuleGen.GenerateRef()
|
||||
key := rule.GetKey()
|
||||
info, _ := sch.registry.getOrCreate(ctx, rule, ruleFactory)
|
||||
info, _ := sch.registry.getOrCreate(ctx, ruleWithFolder{rule: rule, folderTitle: ""}, ruleFactory)
|
||||
|
||||
_, err := sch.updateSchedulableAlertRules(ctx)
|
||||
require.NoError(t, err)
|
||||
|
||||
Reference in New Issue
Block a user