Alerting: Condition evaluator with cached pipeline (#57479)
* create rule evaluator * load header from the context * init one factory * update scheduler
This commit is contained in:
@@ -3,6 +3,7 @@
|
||||
package eval
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"runtime/debug"
|
||||
@@ -24,28 +25,72 @@ import (
|
||||
|
||||
var logger = log.New("ngalert.eval")
|
||||
|
||||
//go:generate mockery --name Evaluator --structname FakeEvaluator --inpackage --filename evaluator_mock.go --with-expecter
|
||||
type Evaluator interface {
|
||||
// ConditionEval executes conditions and evaluates the result.
|
||||
ConditionEval(ctx EvaluationContext, condition models.Condition) Results
|
||||
// QueriesAndExpressionsEval executes queries and expressions and returns the result.
|
||||
QueriesAndExpressionsEval(ctx EvaluationContext, data []models.AlertQuery) (*backend.QueryDataResponse, error)
|
||||
type EvaluatorFactory interface {
|
||||
// Validate validates that the condition is correct. Returns nil if the condition is correct. Otherwise, error that describes the failure
|
||||
Validate(ctx EvaluationContext, condition models.Condition) error
|
||||
// BuildRuleEvaluator build an evaluator pipeline ready to evaluate a rule's query
|
||||
Create(ctx EvaluationContext, condition models.Condition) (ConditionEvaluator, error)
|
||||
}
|
||||
|
||||
//go:generate mockery --name ConditionEvaluator --structname ConditionEvaluatorMock --with-expecter --output eval_mocks --outpkg eval_mocks
|
||||
type ConditionEvaluator interface {
|
||||
// EvaluateRaw evaluates the condition and returns raw backend response backend.QueryDataResponse
|
||||
EvaluateRaw(ctx context.Context, now time.Time) (resp *backend.QueryDataResponse, err error)
|
||||
// Evaluate evaluates the condition and converts the response to Results
|
||||
Evaluate(ctx context.Context, now time.Time) (Results, error)
|
||||
}
|
||||
|
||||
type conditionEvaluator struct {
|
||||
pipeline expr.DataPipeline
|
||||
expressionService *expr.Service
|
||||
condition models.Condition
|
||||
evalTimeout time.Duration
|
||||
}
|
||||
|
||||
func (r *conditionEvaluator) EvaluateRaw(ctx context.Context, now time.Time) (resp *backend.QueryDataResponse, err error) {
|
||||
defer func() {
|
||||
if e := recover(); e != nil {
|
||||
logger.FromContext(ctx).Error("alert rule panic", "error", e, "stack", string(debug.Stack()))
|
||||
panicErr := fmt.Errorf("alert rule panic; please check the logs for the full stack")
|
||||
if err != nil {
|
||||
err = fmt.Errorf("queries and expressions execution failed: %w; %v", err, panicErr.Error())
|
||||
} else {
|
||||
err = panicErr
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
execCtx := ctx
|
||||
if r.evalTimeout <= 0 {
|
||||
timeoutCtx, cancel := context.WithTimeout(ctx, r.evalTimeout)
|
||||
defer cancel()
|
||||
execCtx = timeoutCtx
|
||||
}
|
||||
return r.expressionService.ExecutePipeline(execCtx, now, r.pipeline)
|
||||
}
|
||||
|
||||
// Evaluate evaluates the condition and converts the response to Results
|
||||
func (r *conditionEvaluator) Evaluate(ctx context.Context, now time.Time) (Results, error) {
|
||||
response, err := r.EvaluateRaw(ctx, now)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
execResults := queryDataResponseToExecutionResults(r.condition, response)
|
||||
return evaluateExecutionResult(execResults, now), nil
|
||||
}
|
||||
|
||||
type evaluatorImpl struct {
|
||||
cfg *setting.Cfg
|
||||
evaluationTimeout time.Duration
|
||||
dataSourceCache datasources.CacheService
|
||||
expressionService *expr.Service
|
||||
}
|
||||
|
||||
func NewEvaluator(
|
||||
cfg *setting.Cfg,
|
||||
func NewEvaluatorFactory(
|
||||
cfg setting.UnifiedAlertingSettings,
|
||||
datasourceCache datasources.CacheService,
|
||||
expressionService *expr.Service) Evaluator {
|
||||
expressionService *expr.Service) EvaluatorFactory {
|
||||
return &evaluatorImpl{
|
||||
cfg: cfg,
|
||||
evaluationTimeout: cfg.EvaluationTimeout,
|
||||
dataSourceCache: datasourceCache,
|
||||
expressionService: expressionService,
|
||||
}
|
||||
@@ -122,6 +167,15 @@ type Result struct {
|
||||
EvaluationString string
|
||||
}
|
||||
|
||||
func NewResultFromError(err error, evaluatedAt time.Time, duration time.Duration) Result {
|
||||
return Result{
|
||||
State: Error,
|
||||
Error: err,
|
||||
EvaluatedAt: evaluatedAt,
|
||||
EvaluationDuration: duration,
|
||||
}
|
||||
}
|
||||
|
||||
// State is an enum of the evaluation State for an alert instance.
|
||||
type State int
|
||||
|
||||
@@ -169,8 +223,9 @@ func buildDatasourceHeaders(ctx EvaluationContext) map[string]string {
|
||||
"X-Cache-Skip": "true",
|
||||
}
|
||||
|
||||
if ctx.RuleUID != "" {
|
||||
headers["X-Rule-Uid"] = ctx.RuleUID
|
||||
key, ok := models.RuleKeyFromContext(ctx.Ctx)
|
||||
if ok {
|
||||
headers["X-Rule-Uid"] = key.UID
|
||||
}
|
||||
|
||||
return headers
|
||||
@@ -312,27 +367,6 @@ func queryDataResponseToExecutionResults(c models.Condition, execResp *backend.Q
|
||||
return result
|
||||
}
|
||||
|
||||
func executeQueriesAndExpressions(ctx EvaluationContext, data []models.AlertQuery, exprService *expr.Service, dsCacheService datasources.CacheService) (resp *backend.QueryDataResponse, err error) {
|
||||
defer func() {
|
||||
if e := recover(); e != nil {
|
||||
logger.FromContext(ctx.Ctx).Error("alert rule panic", "error", e, "stack", string(debug.Stack()))
|
||||
panicErr := fmt.Errorf("alert rule panic; please check the logs for the full stack")
|
||||
if err != nil {
|
||||
err = fmt.Errorf("queries and expressions execution failed: %w; %v", err, panicErr.Error())
|
||||
} else {
|
||||
err = panicErr
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
queryDataReq, err := getExprRequest(ctx, data, dsCacheService)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return exprService.TransformData(ctx.Ctx, ctx.At, queryDataReq)
|
||||
}
|
||||
|
||||
// datasourceUIDsToRefIDs returns a sorted slice of Ref IDs for each Datasource UID.
|
||||
//
|
||||
// If refIDsToDatasourceUIDs is nil then this function also returns nil. Likewise,
|
||||
@@ -399,12 +433,7 @@ func evaluateExecutionResult(execResults ExecutionResults, ts time.Time) Results
|
||||
evalResults := make([]Result, 0)
|
||||
|
||||
appendErrRes := func(e error) {
|
||||
evalResults = append(evalResults, Result{
|
||||
State: Error,
|
||||
Error: e,
|
||||
EvaluatedAt: ts,
|
||||
EvaluationDuration: time.Since(ts),
|
||||
})
|
||||
evalResults = append(evalResults, NewResultFromError(e, ts, time.Since(ts)))
|
||||
}
|
||||
|
||||
appendNoData := func(labels data.Labels) {
|
||||
@@ -560,55 +589,37 @@ func (evalResults Results) AsDataFrame() data.Frame {
|
||||
return *frame
|
||||
}
|
||||
|
||||
// ConditionEval executes conditions and evaluates the result.
|
||||
func (e *evaluatorImpl) ConditionEval(ctx EvaluationContext, condition models.Condition) Results {
|
||||
execResp, err := e.QueriesAndExpressionsEval(ctx, condition.Data)
|
||||
var execResults ExecutionResults
|
||||
if err != nil {
|
||||
execResults = ExecutionResults{Error: err}
|
||||
} else {
|
||||
execResults = queryDataResponseToExecutionResults(condition, execResp)
|
||||
}
|
||||
return evaluateExecutionResult(execResults, ctx.At)
|
||||
}
|
||||
|
||||
// QueriesAndExpressionsEval executes queries and expressions and returns the result.
|
||||
func (e *evaluatorImpl) QueriesAndExpressionsEval(ctx EvaluationContext, data []models.AlertQuery) (*backend.QueryDataResponse, error) {
|
||||
timeoutCtx, cancelFn := ctx.WithTimeout(e.cfg.UnifiedAlerting.EvaluationTimeout)
|
||||
defer cancelFn()
|
||||
|
||||
execResult, err := executeQueriesAndExpressions(timeoutCtx, data, e.expressionService, e.dataSourceCache)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to execute conditions: %w", err)
|
||||
}
|
||||
|
||||
return execResult, nil
|
||||
}
|
||||
|
||||
func (e *evaluatorImpl) Validate(ctx EvaluationContext, condition models.Condition) error {
|
||||
_, err := e.Create(ctx, condition)
|
||||
return err
|
||||
}
|
||||
|
||||
func (e *evaluatorImpl) Create(ctx EvaluationContext, condition models.Condition) (ConditionEvaluator, error) {
|
||||
if len(condition.Data) == 0 {
|
||||
return errors.New("expression list is empty. must be at least 1 expression")
|
||||
return nil, errors.New("expression list is empty. must be at least 1 expression")
|
||||
}
|
||||
if len(condition.Condition) == 0 {
|
||||
return errors.New("condition must not be empty")
|
||||
return nil, errors.New("condition must not be empty")
|
||||
}
|
||||
|
||||
ctx.At = time.Now()
|
||||
|
||||
req, err := getExprRequest(ctx, condition.Data, e.dataSourceCache)
|
||||
if err != nil {
|
||||
return err
|
||||
return nil, err
|
||||
}
|
||||
pipeline, err := e.expressionService.BuildPipeline(req)
|
||||
if err != nil {
|
||||
return err
|
||||
return nil, err
|
||||
}
|
||||
conditions := make([]string, 0, len(pipeline))
|
||||
for _, node := range pipeline {
|
||||
if node.RefID() == condition.Condition {
|
||||
return nil
|
||||
return &conditionEvaluator{
|
||||
pipeline: pipeline,
|
||||
expressionService: e.expressionService,
|
||||
condition: condition,
|
||||
evalTimeout: e.evaluationTimeout,
|
||||
}, nil
|
||||
}
|
||||
conditions = append(conditions, node.RefID())
|
||||
}
|
||||
return fmt.Errorf("condition %s does not exist, must be one of %v", condition.Condition, conditions)
|
||||
return nil, fmt.Errorf("condition %s does not exist, must be one of %v", condition.Condition, conditions)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user