Alerting: add state tracker to alerting evaluation (#32298)

* Initial commit for state tracking

* basic state transition logic and tests

* constructor. test and interface fixup

* use new sig for sch.definitionRoutine()

* test fixup

* make the linter happy

* more minor linting cleanup
This commit is contained in:
David Parrott
2021-03-24 15:34:18 -07:00
committed by GitHub
parent 58b814bd7d
commit d33a77a67f
6 changed files with 252 additions and 20 deletions
+11 -11
View File
@@ -54,26 +54,26 @@ type ExecutionResults struct {
} }
// Results is a slice of evaluated alert instances states. // Results is a slice of evaluated alert instances states.
type Results []result type Results []Result
// result contains the evaluated state of an alert instance // Result contains the evaluated State of an alert instance
// identified by its labels. // identified by its labels.
type result struct { type Result struct {
Instance data.Labels Instance data.Labels
State state // Enum State State // Enum
// StartAt is the time at which we first saw this state // StartAt is the time at which we first saw this state
StartAt time.Time StartAt time.Time
// FiredAt is the time at which we first transitioned to a firing state // FiredAt is the time at which we first transitioned to a firing state
FiredAt time.Time FiredAt time.Time
} }
// state is an enum of the evaluation state for an alert instance. // State is an enum of the evaluation State for an alert instance.
type state int type State int
const ( const (
// Normal is the eval state for an alert instance condition // Normal is the eval state for an alert instance condition
// that evaluated to false. // that evaluated to false.
Normal state = iota Normal State = iota
// Alerting is the eval state for an alert instance condition // Alerting is the eval state for an alert instance condition
// that evaluated to true (Alerting). // that evaluated to true (Alerting).
@@ -88,7 +88,7 @@ const (
Error Error
) )
func (s state) String() string { func (s State) String() string {
return [...]string{"Normal", "Alerting", "NoData", "Error"}[s] return [...]string{"Normal", "Alerting", "NoData", "Error"}[s]
} }
@@ -177,9 +177,9 @@ func execute(ctx AlertExecCtx, c *models.Condition, now time.Time, dataService *
} }
// evaluateExecutionResult takes the ExecutionResult, and returns a frame where // evaluateExecutionResult takes the ExecutionResult, and returns a frame where
// each column is a string type that holds a string representing its state. // each column is a string type that holds a string representing its State.
func evaluateExecutionResult(results *ExecutionResults, ts time.Time) (Results, error) { func evaluateExecutionResult(results *ExecutionResults, ts time.Time) (Results, error) {
evalResults := make([]result, 0) evalResults := make([]Result, 0)
labels := make(map[string]bool) labels := make(map[string]bool)
for _, f := range results.Results { for _, f := range results.Results {
rowLen, err := f.RowLen() rowLen, err := f.RowLen()
@@ -210,7 +210,7 @@ func evaluateExecutionResult(results *ExecutionResults, ts time.Time) (Results,
return nil, &invalidEvalResultFormatError{refID: f.RefID, reason: fmt.Sprintf("expected nullable float64 but got type %T", f.Fields[0].Type())} return nil, &invalidEvalResultFormatError{refID: f.RefID, reason: fmt.Sprintf("expected nullable float64 but got type %T", f.Fields[0].Type())}
} }
r := result{ r := Result{
Instance: f.Fields[0].Labels, Instance: f.Fields[0].Labels,
StartAt: ts, StartAt: ts,
} }
+5 -2
View File
@@ -4,6 +4,8 @@ import (
"context" "context"
"time" "time"
"github.com/grafana/grafana/pkg/services/ngalert/state"
"github.com/benbjohnson/clock" "github.com/benbjohnson/clock"
"github.com/grafana/grafana/pkg/api/routing" "github.com/grafana/grafana/pkg/api/routing"
@@ -45,6 +47,7 @@ type AlertNG struct {
DataProxy *datasourceproxy.DatasourceProxyService `inject:""` DataProxy *datasourceproxy.DatasourceProxyService `inject:""`
Log log.Logger Log log.Logger
schedule schedule.ScheduleService schedule schedule.ScheduleService
stateTracker *state.StateTracker
} }
func init() { func init() {
@@ -54,7 +57,7 @@ func init() {
// Init initializes the AlertingService. // Init initializes the AlertingService.
func (ng *AlertNG) Init() error { func (ng *AlertNG) Init() error {
ng.Log = log.New("ngalert") ng.Log = log.New("ngalert")
ng.stateTracker = state.NewStateTracker()
baseInterval := baseIntervalSeconds * time.Second baseInterval := baseIntervalSeconds * time.Second
store := store.DBstore{BaseInterval: baseInterval, DefaultIntervalSeconds: defaultIntervalSeconds, SQLStore: ng.SQLStore} store := store.DBstore{BaseInterval: baseInterval, DefaultIntervalSeconds: defaultIntervalSeconds, SQLStore: ng.SQLStore}
@@ -87,7 +90,7 @@ func (ng *AlertNG) Init() error {
// Run starts the scheduler // Run starts the scheduler
func (ng *AlertNG) Run(ctx context.Context) error { func (ng *AlertNG) Run(ctx context.Context) error {
ng.Log.Debug("ngalert starting") ng.Log.Debug("ngalert starting")
return ng.schedule.Ticker(ctx) return ng.schedule.Ticker(ctx, ng.stateTracker)
} }
// IsDisabled returns true if the alerting service is disable for this instance. // IsDisabled returns true if the alerting service is disable for this instance.
+8 -6
View File
@@ -6,6 +6,8 @@ import (
"sync" "sync"
"time" "time"
"github.com/grafana/grafana/pkg/services/ngalert/state"
"github.com/grafana/grafana/pkg/services/ngalert/store" "github.com/grafana/grafana/pkg/services/ngalert/store"
"github.com/grafana/grafana/pkg/services/ngalert/models" "github.com/grafana/grafana/pkg/services/ngalert/models"
@@ -23,7 +25,7 @@ var timeNow = time.Now
// ScheduleService handles scheduling // ScheduleService handles scheduling
type ScheduleService interface { type ScheduleService interface {
Ticker(context.Context) error Ticker(context.Context, *state.StateTracker) error
Pause() error Pause() error
Unpause() error Unpause() error
@@ -33,8 +35,7 @@ type ScheduleService interface {
overrideCfg(cfg SchedulerCfg) overrideCfg(cfg SchedulerCfg)
} }
func (sch *schedule) definitionRoutine(grafanaCtx context.Context, key models.AlertDefinitionKey, func (sch *schedule) definitionRoutine(grafanaCtx context.Context, key models.AlertDefinitionKey, evalCh <-chan *evalContext, stopCh <-chan struct{}, stateTracker *state.StateTracker) error {
evalCh <-chan *evalContext, stopCh <-chan struct{}) error {
sch.log.Debug("alert definition routine started", "key", key) sch.log.Debug("alert definition routine started", "key", key)
evalRunning := false evalRunning := false
@@ -77,13 +78,14 @@ func (sch *schedule) definitionRoutine(grafanaCtx context.Context, key models.Al
return err return err
} }
for _, r := range results { for _, r := range results {
sch.log.Debug("alert definition result", "title", alertDefinition.Title, "key", key, "attempt", attempt, "now", ctx.now, "duration", end.Sub(start), "instance", r.Instance, "state", r.State.String()) sch.log.Info("alert definition result", "title", alertDefinition.Title, "key", key, "attempt", attempt, "now", ctx.now, "duration", end.Sub(start), "instance", r.Instance, "state", r.State.String())
cmd := models.SaveAlertInstanceCommand{DefinitionOrgID: key.OrgID, DefinitionUID: key.DefinitionUID, State: models.InstanceStateType(r.State.String()), Labels: models.InstanceLabels(r.Instance), LastEvalTime: ctx.now} cmd := models.SaveAlertInstanceCommand{DefinitionOrgID: key.OrgID, DefinitionUID: key.DefinitionUID, State: models.InstanceStateType(r.State.String()), Labels: models.InstanceLabels(r.Instance), LastEvalTime: ctx.now}
err := sch.store.SaveAlertInstance(&cmd) err := sch.store.SaveAlertInstance(&cmd)
if err != nil { if err != nil {
sch.log.Error("failed saving alert instance", "title", alertDefinition.Title, "key", key, "attempt", attempt, "now", ctx.now, "instance", r.Instance, "state", r.State.String(), "error", err) sch.log.Error("failed saving alert instance", "title", alertDefinition.Title, "key", key, "attempt", attempt, "now", ctx.now, "instance", r.Instance, "state", r.State.String(), "error", err)
} }
} }
_ = stateTracker.ProcessEvalResults(key.DefinitionUID, results, condition)
return nil return nil
} }
@@ -217,7 +219,7 @@ func (sch *schedule) Unpause() error {
return nil return nil
} }
func (sch *schedule) Ticker(grafanaCtx context.Context) error { func (sch *schedule) Ticker(grafanaCtx context.Context, stateTracker *state.StateTracker) error {
dispatcherGroup, ctx := errgroup.WithContext(grafanaCtx) dispatcherGroup, ctx := errgroup.WithContext(grafanaCtx)
for { for {
select { select {
@@ -250,7 +252,7 @@ func (sch *schedule) Ticker(grafanaCtx context.Context) error {
if newRoutine && !invalidInterval { if newRoutine && !invalidInterval {
dispatcherGroup.Go(func() error { dispatcherGroup.Go(func() error {
return sch.definitionRoutine(ctx, key, definitionInfo.evalCh, definitionInfo.stopCh) return sch.definitionRoutine(ctx, key, definitionInfo.evalCh, definitionInfo.stopCh, stateTracker)
}) })
} }
+100
View File
@@ -0,0 +1,100 @@
package state
import (
"fmt"
"sync"
"github.com/grafana/grafana-plugin-sdk-go/data"
"github.com/grafana/grafana/pkg/services/ngalert/eval"
"github.com/grafana/grafana/pkg/services/ngalert/models"
)
type AlertState struct {
UID string
CacheId string
Labels data.Labels
State eval.State
Results []eval.State
}
type cache struct {
cacheMap map[string]AlertState
mu sync.Mutex
}
type StateTracker struct {
stateCache cache
}
func NewStateTracker() *StateTracker {
return &StateTracker{
stateCache: cache{
cacheMap: make(map[string]AlertState),
mu: sync.Mutex{},
},
}
}
func (c *cache) getOrCreate(uid string, result eval.Result) AlertState {
c.mu.Lock()
defer c.mu.Unlock()
idString := fmt.Sprintf("%s %s", uid, result.Instance.String())
if state, ok := c.cacheMap[idString]; ok {
return state
}
newState := AlertState{
UID: uid,
CacheId: idString,
Labels: result.Instance,
State: result.State,
Results: []eval.State{result.State},
}
c.cacheMap[idString] = newState
return newState
}
func (c *cache) update(stateEntry AlertState) {
c.mu.Lock()
defer c.mu.Unlock()
c.cacheMap[stateEntry.CacheId] = stateEntry
}
func (c *cache) getStateForEntry(stateId string) eval.State {
c.mu.Lock()
defer c.mu.Unlock()
return c.cacheMap[stateId].State
}
func (st *StateTracker) ProcessEvalResults(uid string, results eval.Results, condition models.Condition) []AlertState {
var changedStates []AlertState
for _, result := range results {
currentState := st.stateCache.getOrCreate(uid, result)
currentState.Results = append(currentState.Results, result.State)
newState := st.getNextState(uid, result)
if newState != currentState.State {
currentState.State = newState
changedStates = append(changedStates, currentState)
}
st.stateCache.update(currentState)
}
return changedStates
}
func (st *StateTracker) getNextState(uid string, result eval.Result) eval.State {
currentState := st.stateCache.getOrCreate(uid, result)
if currentState.State == result.State {
return currentState.State
}
switch {
case currentState.State == result.State:
return currentState.State
case currentState.State == eval.Normal && result.State == eval.Alerting:
return eval.Alerting
case currentState.State == eval.Alerting && result.State == eval.Normal:
return eval.Normal
default:
return eval.Alerting
}
}
@@ -0,0 +1,123 @@
package state
import (
"testing"
"github.com/grafana/grafana-plugin-sdk-go/data"
"github.com/grafana/grafana/pkg/services/ngalert/eval"
"github.com/grafana/grafana/pkg/services/ngalert/models"
"github.com/stretchr/testify/assert"
)
func TestProcessEvalResults(t *testing.T) {
testCases := []struct {
desc string
uid string
evalResults eval.Results
condition models.Condition
expectedCacheEntries int
expectedState eval.State
expectedResultCount int
}{
{
desc: "given a single evaluation result",
uid: "test_uid",
evalResults: eval.Results{
eval.Result{
Instance: data.Labels{"label1": "value1", "label2": "value2"},
},
},
expectedCacheEntries: 1,
expectedState: eval.Normal,
expectedResultCount: 0,
},
{
desc: "given a state change from normal to alerting",
uid: "test_uid",
evalResults: eval.Results{
eval.Result{
Instance: data.Labels{"label1": "value1", "label2": "value2"},
State: eval.Normal,
},
eval.Result{
Instance: data.Labels{"label1": "value1", "label2": "value2"},
State: eval.Alerting,
},
},
expectedCacheEntries: 1,
expectedState: eval.Alerting,
expectedResultCount: 1,
},
{
desc: "given a state change from alerting to normal",
uid: "test_uid",
evalResults: eval.Results{
eval.Result{
Instance: data.Labels{"label1": "value1", "label2": "value2"},
State: eval.Alerting,
},
eval.Result{
Instance: data.Labels{"label1": "value1", "label2": "value2"},
State: eval.Normal,
},
},
expectedCacheEntries: 1,
expectedState: eval.Normal,
expectedResultCount: 1,
},
{
desc: "given a constant alerting state",
uid: "test_uid",
evalResults: eval.Results{
eval.Result{
Instance: data.Labels{"label1": "value1", "label2": "value2"},
State: eval.Alerting,
},
eval.Result{
Instance: data.Labels{"label1": "value1", "label2": "value2"},
State: eval.Alerting,
},
},
expectedCacheEntries: 1,
expectedState: eval.Alerting,
expectedResultCount: 0,
},
{
desc: "given a constant normal state",
uid: "test_uid",
evalResults: eval.Results{
eval.Result{
Instance: data.Labels{"label1": "value1", "label2": "value2"},
State: eval.Normal,
},
eval.Result{
Instance: data.Labels{"label1": "value1", "label2": "value2"},
State: eval.Normal,
},
},
expectedCacheEntries: 1,
expectedState: eval.Normal,
expectedResultCount: 0,
},
}
for _, tc := range testCases {
t.Run("the correct number of entries are added to the cache", func(t *testing.T) {
st := NewStateTracker()
st.ProcessEvalResults(tc.uid, tc.evalResults, tc.condition)
assert.Equal(t, len(st.stateCache.cacheMap), tc.expectedCacheEntries)
})
t.Run("the correct state is set", func(t *testing.T) {
st := NewStateTracker()
st.ProcessEvalResults(tc.uid, tc.evalResults, tc.condition)
assert.Equal(t, st.stateCache.getStateForEntry("test_uid label1=value1, label2=value2"), tc.expectedState)
})
t.Run("the correct number of results are returned", func(t *testing.T) {
st := NewStateTracker()
results := st.ProcessEvalResults(tc.uid, tc.evalResults, tc.condition)
assert.Equal(t, len(results), tc.expectedResultCount)
})
}
}
+5 -1
View File
@@ -8,6 +8,8 @@ import (
"testing" "testing"
"time" "time"
"github.com/grafana/grafana/pkg/services/ngalert/state"
"github.com/grafana/grafana/pkg/infra/log" "github.com/grafana/grafana/pkg/infra/log"
"github.com/grafana/grafana/pkg/services/ngalert/schedule" "github.com/grafana/grafana/pkg/services/ngalert/schedule"
@@ -57,8 +59,10 @@ func TestAlertingTicker(t *testing.T) {
sched := schedule.NewScheduler(schefCfg, nil) sched := schedule.NewScheduler(schefCfg, nil)
ctx := context.Background() ctx := context.Background()
st := state.NewStateTracker()
go func() { go func() {
err := sched.Ticker(ctx) err := sched.Ticker(ctx, st)
require.NoError(t, err) require.NoError(t, err)
}() }()
runtime.Gosched() runtime.Gosched()