executor.go 4.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192
  1. package alerting
  2. import (
  3. "fmt"
  4. "strconv"
  5. "math"
  6. "github.com/grafana/grafana/pkg/bus"
  7. "github.com/grafana/grafana/pkg/log"
  8. m "github.com/grafana/grafana/pkg/models"
  9. "github.com/grafana/grafana/pkg/services/alerting/alertstates"
  10. "github.com/grafana/grafana/pkg/tsdb"
  11. )
  12. var (
  13. resultLogFmt = "Alerting: executor %s %1.2f %s %1.2f : %v"
  14. descriptionFmt = "Actual value: %1.2f for %s"
  15. )
  16. type ExecutorImpl struct {
  17. log log.Logger
  18. }
  19. func NewExecutor() *ExecutorImpl {
  20. return &ExecutorImpl{
  21. log: log.New("alerting.executor"),
  22. }
  23. }
  24. type compareFn func(float64, float64) bool
  25. type aggregationFn func(*tsdb.TimeSeries) float64
  26. var operators = map[string]compareFn{
  27. ">": func(num1, num2 float64) bool { return num1 > num2 },
  28. ">=": func(num1, num2 float64) bool { return num1 >= num2 },
  29. "<": func(num1, num2 float64) bool { return num1 < num2 },
  30. "<=": func(num1, num2 float64) bool { return num1 <= num2 },
  31. "": func(num1, num2 float64) bool { return false },
  32. }
  33. var aggregator = map[string]aggregationFn{
  34. "avg": func(series *tsdb.TimeSeries) float64 {
  35. sum := float64(0)
  36. for _, v := range series.Points {
  37. sum += v[0]
  38. }
  39. return sum / float64(len(series.Points))
  40. },
  41. "sum": func(series *tsdb.TimeSeries) float64 {
  42. sum := float64(0)
  43. for _, v := range series.Points {
  44. sum += v[0]
  45. }
  46. return sum
  47. },
  48. "min": func(series *tsdb.TimeSeries) float64 {
  49. min := series.Points[0][0]
  50. for _, v := range series.Points {
  51. if v[0] < min {
  52. min = v[0]
  53. }
  54. }
  55. return min
  56. },
  57. "max": func(series *tsdb.TimeSeries) float64 {
  58. max := series.Points[0][0]
  59. for _, v := range series.Points {
  60. if v[0] > max {
  61. max = v[0]
  62. }
  63. }
  64. return max
  65. },
  66. "mean": func(series *tsdb.TimeSeries) float64 {
  67. midPosition := int64(math.Floor(float64(len(series.Points)) / float64(2)))
  68. return series.Points[midPosition][0]
  69. },
  70. }
  71. func (e *ExecutorImpl) Execute(job *AlertJob, resultQueue chan *AlertResult) {
  72. timeSeries, err := e.executeQuery(job)
  73. if err != nil {
  74. resultQueue <- &AlertResult{
  75. Error: err,
  76. State: alertstates.Pending,
  77. AlertJob: job,
  78. }
  79. }
  80. result := e.evaluateRule(job.Rule, timeSeries)
  81. result.AlertJob = job
  82. resultQueue <- result
  83. }
  84. func (e *ExecutorImpl) executeQuery(job *AlertJob) (tsdb.TimeSeriesSlice, error) {
  85. getDsInfo := &m.GetDataSourceByIdQuery{
  86. Id: job.Rule.DatasourceId,
  87. OrgId: job.Rule.OrgId,
  88. }
  89. if err := bus.Dispatch(getDsInfo); err != nil {
  90. return nil, fmt.Errorf("Could not find datasource for %d", job.Rule.DatasourceId)
  91. }
  92. req := e.GetRequestForAlertRule(job.Rule, getDsInfo.Result)
  93. result := make(tsdb.TimeSeriesSlice, 0)
  94. resp, err := tsdb.HandleRequest(req)
  95. if err != nil {
  96. return nil, fmt.Errorf("Alerting: GetSeries() tsdb.HandleRequest() error %v", err)
  97. }
  98. for _, v := range resp.Results {
  99. if v.Error != nil {
  100. return nil, fmt.Errorf("Alerting: GetSeries() tsdb.HandleRequest() response error %v", v)
  101. }
  102. result = append(result, v.Series...)
  103. }
  104. return result, nil
  105. }
  106. func (e *ExecutorImpl) GetRequestForAlertRule(rule *AlertRule, datasource *m.DataSource) *tsdb.Request {
  107. req := &tsdb.Request{
  108. TimeRange: tsdb.TimeRange{
  109. From: "-" + strconv.Itoa(rule.QueryRange) + "s",
  110. To: "now",
  111. },
  112. Queries: tsdb.QuerySlice{
  113. {
  114. RefId: rule.QueryRefId,
  115. Query: rule.Query,
  116. DataSource: &tsdb.DataSourceInfo{
  117. Id: datasource.Id,
  118. Name: datasource.Name,
  119. PluginId: datasource.Type,
  120. Url: datasource.Url,
  121. },
  122. },
  123. },
  124. }
  125. return req
  126. }
  127. func (e *ExecutorImpl) evaluateRule(rule *AlertRule, series tsdb.TimeSeriesSlice) *AlertResult {
  128. e.log.Debug("Evaluating Alerting Rule", "seriesCount", len(series), "ruleName", rule.Name)
  129. for _, serie := range series {
  130. log.Debug("Evaluating series", "series", serie.Name)
  131. if aggregator[rule.Aggregator] == nil {
  132. continue
  133. }
  134. var aggValue = aggregator[rule.Aggregator](serie)
  135. var critOperartor = operators[rule.CritOperator]
  136. var critResult = critOperartor(aggValue, rule.CritLevel)
  137. log.Trace(resultLogFmt, "Crit", serie.Name, aggValue, rule.CritOperator, rule.CritLevel, critResult)
  138. if critResult {
  139. return &AlertResult{
  140. State: alertstates.Critical,
  141. ActualValue: aggValue,
  142. Description: fmt.Sprintf(descriptionFmt, aggValue, serie.Name),
  143. }
  144. }
  145. var warnOperartor = operators[rule.CritOperator]
  146. var warnResult = warnOperartor(aggValue, rule.CritLevel)
  147. log.Trace(resultLogFmt, "Warn", serie.Name, aggValue, rule.WarnOperator, rule.WarnLevel, warnResult)
  148. if warnResult {
  149. return &AlertResult{
  150. State: alertstates.Warn,
  151. Description: fmt.Sprintf(descriptionFmt, aggValue, serie.Name),
  152. ActualValue: aggValue,
  153. }
  154. }
  155. }
  156. return &AlertResult{State: alertstates.Ok, Description: "Alert is OK!"}
  157. }