Merge pull request #6087 from zhaoxiangchun/feature/zxc-influxdb-schema

在monitor的新增influxdb相关的功能
This commit is contained in:
Zexi Li
2020-05-13 22:59:41 +08:00
committed by GitHub
13 changed files with 644 additions and 30 deletions
+12 -7
View File
@@ -18,8 +18,11 @@ const (
EIP_UN_USED = "EIP_UNUSED"
DISK_UN_USED = "DISK_UNUSED"
LB_UN_USED = "LB_UNUSED"
SCALE_DOWN = "SCALE_DOWN"
SCALE_UP = "SCALE_UP"
DRIVER_ACTION = "DELETE"
DRIVER_ACTION = "DELETE"
SCALE_DOWN_DRIVER_ACTION = "SCALE DOWN"
EIP_UNUSED_START_DELETE = "start_delete"
EIP_UNUSED_DELETE_FAIL = "delete_fail"
@@ -30,13 +33,15 @@ type MonitorSuggest string
type MonitorResourceType string
const (
EIP_MONITOR_RES_TYPE = MonitorResourceType("弹性EIP")
DISK_MONITOR_RES_TYPE = MonitorResourceType("云硬盘")
LB_MONITOR_RES_TYPE = MonitorResourceType("负载均衡实例")
EIP_MONITOR_RES_TYPE = MonitorResourceType("弹性EIP")
DISK_MONITOR_RES_TYPE = MonitorResourceType("云硬盘")
LB_MONITOR_RES_TYPE = MonitorResourceType("负载均衡实例")
SCALE_MONTITOR_RES_TYPE = MonitorResourceType("虚拟机")
)
const (
EIP_MONITOR_SUGGEST = MonitorSuggest("释放未使用的EIP")
DISK_MONITOR_SUGGEST = MonitorSuggest("释放未使用的Disk")
LB_MONITOR_SUGGEST = MonitorSuggest("释放未使用的LB")
EIP_MONITOR_SUGGEST = MonitorSuggest("释放未使用的EIP")
DISK_MONITOR_SUGGEST = MonitorSuggest("释放未使用的Disk")
LB_MONITOR_SUGGEST = MonitorSuggest("释放未使用的LB")
SCALE_DOWN_MONITOR_SUGGEST = MonitorSuggest("缩减机器配置")
)
+44 -4
View File
@@ -20,6 +20,25 @@ import (
"yunion.io/x/onecloud/pkg/apis"
)
const (
METRIC_TAG = "TAG"
METRIC_FIELD = "FIELD"
)
var PROPERTY_TYPE = []string{"databases", "measurements", "metric-measurement"}
var METRIC_ATTRI = []string{METRIC_TAG, METRIC_FIELD}
type InfluxMeasurement struct {
apis.Meta
Database string
Measurement string
TagKey []string
TagValue map[string][]string
FieldKey []string
Unit []string
}
type SuggestSysRuleListInput struct {
apis.VirtualResourceListInput
apis.EnabledResourceBaseListInput
@@ -29,10 +48,11 @@ type SuggestSysRuleCreateInput struct {
apis.VirtualResourceCreateInput
// 查询指标周期
Period string `json:"period"`
Type string `json:"type"`
Enabled *bool `json:"enabled"`
Setting *SSuggestSysAlertSetting `json:"setting"`
Period string `json:"period"`
TimeFrom string `json:"time_from"`
Type string `json:"type"`
Enabled *bool `json:"enabled"`
Setting *SSuggestSysAlertSetting `json:"setting"`
}
type SuggestSysRuleUpdateInput struct {
@@ -60,6 +80,7 @@ type SSuggestSysAlertSetting struct {
EIPUnused *EIPUnused `json:"eip_unused"`
DiskUnused *DiskUnused `json:"disk_unused"`
LBUnused *LBUnused `json:"lb_unused"`
ScaleRule *ScaleRule `json:"scale_rule"`
}
type EIPUnused struct {
@@ -71,3 +92,22 @@ type DiskUnused struct {
type LBUnused struct {
}
type ScaleRule []Scale
type Scale struct {
Database string `json:"database"`
Measurement string `json:"measurement"`
//rule operator rule [and|or]
Operator string `json:"operator"`
Field string `json:"field"`
EvalType string `json:"eval_type"`
Threshold float64 `json:"threshold"`
Tag string `json:"tag"`
TagVal string `json:"tag_val"`
}
type ScaleEvalMatch struct {
EvalMatch
ResourceId map[string]string `json:"resource_id"`
}
@@ -21,14 +21,17 @@ import (
var (
SuggestSysRuleManager *SSuggestSysRuleManager
SuggestSysAlertManager *SSuggestSysAlertManager
InfluxdbShemaManager *SInfluxdbShemaManager
)
func init() {
SuggestSysRuleManager = NewSuggestSysRuleManager()
SuggestSysAlertManager = NewSuggestSysAlertManager()
InfluxdbShemaManager = NewInfluxdbShemaManager()
for _, m := range []modulebase.IBaseManager{
SuggestSysRuleManager,
SuggestSysAlertManager,
InfluxdbShemaManager,
} {
Register(m)
}
@@ -42,6 +45,10 @@ type SSuggestSysAlertManager struct {
*modulebase.ResourceManager
}
type SInfluxdbShemaManager struct {
*modulebase.ResourceManager
}
func NewSuggestSysRuleManager() *SSuggestSysRuleManager {
man := NewMonitorV2Manager("suggestsysrule", "suggestsysrules",
[]string{"id", "name", "type", "enabled", "setting"},
@@ -59,3 +66,12 @@ func NewSuggestSysAlertManager() *SSuggestSysAlertManager {
ResourceManager: &man,
}
}
func NewInfluxdbShemaManager() *SInfluxdbShemaManager {
man := NewMonitorV2Manager("influxdbshema", "influxdbshemas",
[]string{},
[]string{})
return &SInfluxdbShemaManager{
ResourceManager: &man,
}
}
@@ -22,14 +22,17 @@ import (
var (
SuggestSysRuleManager *SSuggestSysRuleManager
SuggestSysAlertManager *SSuggestSysAlertManager
InfluxdbShemaManager *SInfluxdbShemaManager
)
func init() {
SuggestSysRuleManager = NewSuggestSysRuleManager()
SuggestSysAlertManager = NewSuggestSysAlertManager()
InfluxdbShemaManager = NewInfluxdbShemaManager()
for _, m := range []modulebase.IBaseManager{
SuggestSysRuleManager,
SuggestSysAlertManager,
InfluxdbShemaManager,
} {
modules.Register(m)
}
@@ -43,6 +46,10 @@ type SSuggestSysAlertManager struct {
*modulebase.ResourceManager
}
type SInfluxdbShemaManager struct {
*modulebase.ResourceManager
}
func NewSuggestSysRuleManager() *SSuggestSysRuleManager {
man := modules.NewMonitorV2Manager("suggestsysrule", "suggestsysrules",
[]string{"id", "name", "type", "enabled", "setting"},
@@ -60,3 +67,12 @@ func NewSuggestSysAlertManager() *SSuggestSysAlertManager {
ResourceManager: &man,
}
}
func NewInfluxdbShemaManager() *SInfluxdbShemaManager {
man := modules.NewMonitorV2Manager("influxdbshema", "influxdbshemas",
[]string{},
[]string{})
return &SInfluxdbShemaManager{
ResourceManager: &man,
}
}
@@ -0,0 +1,23 @@
package monitor
import (
"yunion.io/x/jsonutils"
"yunion.io/x/onecloud/pkg/apis/monitor"
)
type InfluxdbShemaListOptions struct {
}
type InfluxdbShemaShowOptions struct {
ID string `help:"attribute of the inluxdb" choices:"databases|measurements|metric-measurement"`
Database string `help:influxdb database`
Measurement string `help:influxdb table`
}
func (opt InfluxdbShemaShowOptions) Params() (jsonutils.JSONObject, error) {
input := new(monitor.InfluxMeasurement)
input.Measurement = opt.Measurement
input.Database = opt.Database
return input.JSON(input), nil
}
+29 -1
View File
@@ -38,7 +38,7 @@ type SuggestSysRuleAlertSettingOptions struct {
type SuggestRuleCreateOptions struct {
SuggestSysRuleAlertSettingOptions
Name string `help:"Name of the alert"`
Type string `help:"Type of suggest rule" choices:"EIP_UNUSED|DISK_UNUSED|LB_UNUSED"`
Type string `help:"Type of suggest rule" choices:"EIP_UNUSED|DISK_UNUSED|LB_UNUSED|SCALE_DOWN"`
Enabled bool `help:"Enable rule"`
Period string `help:"Period of suggest rule e.g. '5s', '1m'" default:"30s""`
}
@@ -86,6 +86,34 @@ func newSuggestSysAlertSetting(tp string) *monitor.SSuggestSysAlertSetting {
setting = &monitor.SSuggestSysAlertSetting{
LBUnused: &monitor.LBUnused{},
}
case monitor.SCALE_DOWN:
scaleRuel := make(monitor.ScaleRule, 0)
scale := monitor.Scale{
Database: "telegraf",
Measurement: "vm_cpu",
Operator: "and",
Field: "usage_active",
EvalType: ">=",
Threshold: 50,
Tag: "",
TagVal: "",
}
scaleRuel = append(scaleRuel, scale)
scale = monitor.Scale{
Database: "telegraf",
Measurement: "vm_diskio",
Operator: "or",
Field: "read_bps",
EvalType: ">=",
Threshold: 500,
Tag: "",
TagVal: "",
}
scaleRuel = append(scaleRuel, scale)
setting = &monitor.SSuggestSysAlertSetting{
ScaleRule: &scaleRuel,
}
}
return setting
}
+8 -9
View File
@@ -17,7 +17,6 @@ package conditions
import (
gocontext "context"
"fmt"
"strings"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
@@ -76,14 +75,14 @@ func (c FormatCond) String() string {
}
func (c *QueryCondition) filterTags(tags map[string]string) map[string]string {
ret := make(map[string]string)
for key, val := range tags {
if strings.HasSuffix(key, "_id") {
continue
}
ret[key] = val
}
return ret
//ret := make(map[string]string)
//for key, val := range tags {
// if strings.HasSuffix(key, "_id") {
// continue
// }
// ret[key] = val
//}
return tags
}
// Eval evaluates te `QueryCondition`.
+4
View File
@@ -201,3 +201,7 @@ var conditionFactories = make(map[string]ConditionFactory)
func RegisterCondition(typeName string, factory ConditionFactory) {
conditionFactories[typeName] = factory
}
func GetConditionFactories() map[string]ConditionFactory {
return conditionFactories
}
+123
View File
@@ -17,10 +17,13 @@ package models
import (
"context"
"database/sql"
"fmt"
"strings"
"time"
"golang.org/x/sync/errgroup"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/tristate"
@@ -28,10 +31,12 @@ import (
"yunion.io/x/onecloud/pkg/apis/monitor"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/httperrors"
"yunion.io/x/onecloud/pkg/mcclient/auth"
"yunion.io/x/onecloud/pkg/monitor/options"
"yunion.io/x/onecloud/pkg/monitor/registry"
"yunion.io/x/onecloud/pkg/monitor/tsdb"
"yunion.io/x/onecloud/pkg/util/influxdb"
)
var (
@@ -167,3 +172,121 @@ func (ds *SDataSource) ToTSDBDataSource(db string) *tsdb.DataSource {
TimeInterval: ds.TimeInterval,*/
}
}
func (self *SDataSourceManager) GetDatabases() (jsonutils.JSONObject, error) {
ret := jsonutils.NewDict()
dataSource, err := self.GetDefaultSource()
if err != nil {
return jsonutils.JSONNull, errors.Wrap(err, "s.GetDefaultSource")
}
db := influxdb.NewInfluxdb(dataSource.Url)
//db.SetDatabase("telegraf")
databases, err := db.GetDatabases()
if err != nil {
return jsonutils.JSONNull, errors.Wrap(err, "GetDatabases")
}
ret.Add(jsonutils.NewStringArray(databases), "databases")
return ret, nil
}
func (self *SDataSourceManager) GetMeasurements(query jsonutils.JSONObject) (jsonutils.JSONObject, error) {
ret := jsonutils.NewDict()
database, _ := query.GetString("database")
if database == "" {
return jsonutils.JSONNull, httperrors.NewInputParameterError("not support database")
}
dataSource, err := self.GetDefaultSource()
if err != nil {
return jsonutils.JSONNull, errors.Wrap(err, "s.GetDefaultSource")
}
db := influxdb.NewInfluxdb(dataSource.Url)
db.SetDatabase(database)
dbRtn, err := db.Query(fmt.Sprintf("SHOW MEASUREMENTS ON %s", database))
if err != nil {
return jsonutils.JSONNull, errors.Wrap(err, "SHOW MEASUREMENTS")
}
res := dbRtn[0][0]
measurements := make([]monitor.InfluxMeasurement, len(res.Values))
for i := range res.Values {
tmpDict := jsonutils.NewDict()
tmpDict.Add(res.Values[i][0], "measurement")
err := tmpDict.Unmarshal(&measurements[i])
if err != nil {
return jsonutils.JSONNull, errors.Wrap(err, "measurement unmarshal error")
}
}
ret.Add(jsonutils.Marshal(&measurements), "measurements")
return ret, nil
}
func (self *SDataSourceManager) GetMetricMeasurement(query jsonutils.JSONObject) (jsonutils.JSONObject, error) {
database, _ := query.GetString("database")
if database == "" {
return jsonutils.JSONNull, httperrors.NewInputParameterError("not support database")
}
measurement, _ := query.GetString("measurement")
if measurement == "" {
return jsonutils.JSONNull, httperrors.NewInputParameterError("not support measurement")
}
dataSource, err := self.GetDefaultSource()
if err != nil {
return jsonutils.JSONNull, errors.Wrap(err, "s.GetDefaultSource")
}
db := influxdb.NewInfluxdb(dataSource.Url)
db.SetDatabase(database)
output := new(monitor.InfluxMeasurement)
output.Measurement = measurement
output.Database = database
for _, val := range monitor.METRIC_ATTRI {
err = getAttributesOnMeasurement(database, val, output, db)
if err != nil {
return jsonutils.JSONNull, errors.Wrap(err, "getAttributesOnMeasurement error")
}
}
err = getTagValue(database, output, db)
if err != nil {
return jsonutils.JSONNull, errors.Wrap(err, "getTagValue error")
}
return jsonutils.Marshal(output), nil
}
func getAttributesOnMeasurement(database, tp string, output *monitor.InfluxMeasurement, db *influxdb.SInfluxdb) error {
dbRtn, err := db.Query(fmt.Sprintf("SHOW %s KEYS ON %s FROM %s", tp, database, output.Measurement))
if err != nil {
return errors.Wrap(err, "SHOW MEASUREMENTS")
}
res := dbRtn[0][0]
tmpDict := jsonutils.NewDict()
tmpArr := jsonutils.NewArray()
for i := range res.Values {
tmpArr.Add(res.Values[i][0])
}
tmpDict.Add(tmpArr, res.Columns[0])
err = tmpDict.Unmarshal(output)
if err != nil {
return errors.Wrap(err, "measurement unmarshal error")
}
return nil
}
func getTagValue(database string, output *monitor.InfluxMeasurement, db *influxdb.SInfluxdb) error {
dbRtn, err := db.Query(fmt.Sprintf("SHOW TAG VALUES ON %s FROM %s WITH KEY IN (%s)", database, output.Measurement, strings.Join(output.TagKey, ",")))
if err != nil {
return errors.Wrap(err, "SHOW MEASUREMENTS")
}
res := dbRtn[0][0]
tagValue := make(map[string][]string, 0)
for i := range res.Values {
val := res.Values[i][0].(*jsonutils.JSONString)
if _, ok := tagValue[val.Value()]; !ok {
tagValue[val.Value()] = make([]string, 0)
}
tag := res.Values[i][1].(*jsonutils.JSONString)
tagValue[val.Value()] = append(tagValue[val.Value()], tag.Value())
}
output.TagValue = tagValue
return nil
}
+3 -1
View File
@@ -200,6 +200,8 @@ func (self *SSuggestSysAlert) getMoreDetails(out monitor.SuggestSysAlertDetails)
out.Suggest = string(monitor.DISK_MONITOR_SUGGEST)
case monitor.LB_UN_USED:
out.Suggest = string(monitor.LB_MONITOR_SUGGEST)
case monitor.SCALE_DOWN:
out.Suggest = string(monitor.SCALE_DOWN_MONITOR_SUGGEST)
}
return out
@@ -263,7 +265,7 @@ func (self *SSuggestSysAlert) Delete(ctx context.Context, userCred mcclient.Toke
}
func (self *SSuggestSysAlert) RealDelete(ctx context.Context, userCred mcclient.TokenCredential) error {
return self.SVirtualResourceBase.Delete(ctx, userCred)
return db.DeleteModel(ctx, userCred, self)
}
func (self *SSuggestSysAlert) StartDeleteTask(
+45 -7
View File
@@ -60,6 +60,7 @@ type SSuggestSysRule struct {
Type string `width:"256" charset:"ascii" list:"user" update:"user"`
Period string `width:"256" charset:"ascii" list:"user" update:"user"`
TimeFrom string `width:"256" charset:"ascii" list:"user" update:"user"`
Setting jsonutils.JSONObject ` list:"user" update:"user"`
ExecTime time.Time `json:"exec_time"`
}
@@ -84,13 +85,9 @@ func (man *SSuggestSysRuleManager) FetchSuggestSysAlartSettings(ruleTypes ...str
//根据数据库中查询得到的信息进行适配转换,同时更新drivers中的内容
func (dConfig *SSuggestSysRule) getSuggestSysAlertSetting() (*monitor.SSuggestSysAlertSetting, error) {
setting := new(monitor.SSuggestSysAlertSetting)
switch dConfig.Type {
case monitor.EIP_UN_USED:
setting.EIPUnused = new(monitor.EIPUnused)
err := dConfig.Setting.Unmarshal(setting.EIPUnused)
if err != nil {
return nil, errors.Wrap(err, "SSuggestSysRule getSuggestSysAlertSetting error")
}
err := dConfig.Setting.Unmarshal(setting)
if err != nil {
return nil, errors.Wrap(err, "SSuggestSysRule getSuggestSysAlertSetting error")
}
return setting, nil
}
@@ -124,6 +121,9 @@ func (man *SSuggestSysRuleManager) ValidateCreateData(
// default 30s
data.Period = "30s"
}
if data.TimeFrom == "" {
data.TimeFrom = "24h"
}
if data.Enabled == nil {
enable := true
data.Enabled = &enable
@@ -131,6 +131,9 @@ func (man *SSuggestSysRuleManager) ValidateCreateData(
if _, err := time.ParseDuration(data.Period); err != nil {
return data, httperrors.NewInputParameterError("Invalid period format: %s", data.Period)
}
if _, err := time.ParseDuration(data.TimeFrom); err != nil {
return data, httperrors.NewInputParameterError("Invalid period format: %s", data.TimeFrom)
}
if dri, ok := suggestSysRuleDrivers[data.Type]; !ok {
return data, httperrors.NewInputParameterError("not support type %q", data.Type)
} else {
@@ -139,6 +142,11 @@ func (man *SSuggestSysRuleManager) ValidateCreateData(
if err != nil {
return data, err
}
if data.Type == monitor.SCALE_DOWN || data.Type == monitor.SCALE_UP {
if data.Setting == nil {
return data, httperrors.NewInputParameterError("no found rule setting")
}
}
if data.Setting != nil {
err = dri.ValidateSetting(data.Setting)
if err != nil {
@@ -269,6 +277,36 @@ func (self *SSuggestSysRuleManager) GetPropertyRuleType(ctx context.Context, use
return ret, nil
}
func (self *SSuggestSysRuleManager) AllowGetPropertyDatabases(ctx context.Context, userCred mcclient.TokenCredential,
query jsonutils.JSONObject) bool {
return true
}
func (self *SSuggestSysRuleManager) GetPropertyDatabases(ctx context.Context, userCred mcclient.TokenCredential,
query jsonutils.JSONObject) (jsonutils.JSONObject, error) {
return DataSourceManager.GetDatabases()
}
func (self *SSuggestSysRuleManager) AllowGetPropertyMeasurements(ctx context.Context, userCred mcclient.TokenCredential,
query jsonutils.JSONObject) bool {
return true
}
func (self *SSuggestSysRuleManager) GetPropertyMeasurements(ctx context.Context, userCred mcclient.TokenCredential,
query jsonutils.JSONObject) (jsonutils.JSONObject, error) {
return DataSourceManager.GetMeasurements(query)
}
func (self *SSuggestSysRuleManager) AllowGetPropertyMetricMeasurement(ctx context.Context,
userCred mcclient.TokenCredential,
query jsonutils.JSONObject) bool {
return true
}
func (self *SSuggestSysRuleManager) GetPropertyMetricMeasurement(ctx context.Context, userCred mcclient.TokenCredential,
query jsonutils.JSONObject) (jsonutils.JSONObject, error) {
return DataSourceManager.GetMetricMeasurement(query)
}
func (self *SSuggestSysRuleManager) GetRules(tp ...string) ([]SSuggestSysRule, error) {
rules := make([]SSuggestSysRule, 0)
query := self.Query()
@@ -25,7 +25,7 @@ import (
)
func init() {
models.RegisterSuggestSysRuleDrivers(NewEIPUsedDriver(), NewDiskUnusedDriver(), NewLBUnusedDriver())
models.RegisterSuggestSysRuleDrivers(NewEIPUsedDriver(), NewDiskUnusedDriver(), NewLBUnusedDriver(), NewScaleDownDriver())
}
func InitSuggestSysRuleCronjob() {
+320
View File
@@ -0,0 +1,320 @@
package suggestsysdrivers
import (
"context"
"fmt"
"strings"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/utils"
"yunion.io/x/onecloud/pkg/apis/monitor"
"yunion.io/x/onecloud/pkg/httperrors"
"yunion.io/x/onecloud/pkg/mcclient"
"yunion.io/x/onecloud/pkg/mcclient/auth"
"yunion.io/x/onecloud/pkg/mcclient/modules"
"yunion.io/x/onecloud/pkg/monitor/alerting"
"yunion.io/x/onecloud/pkg/monitor/models"
"yunion.io/x/onecloud/pkg/monitor/validators"
)
type ScaleDown struct {
monitor.ScaleRule
}
func NewScaleDownDriver() models.ISuggestSysRuleDriver {
return &ScaleDown{
ScaleRule: []monitor.Scale{},
}
}
func (_ *ScaleDown) GetType() string {
return monitor.SCALE_DOWN
}
func (rule *ScaleDown) GetResourceType() string {
return string(monitor.SCALE_MONTITOR_RES_TYPE)
}
func (rule *ScaleDown) ValidateSetting(input *monitor.SSuggestSysAlertSetting) error {
if input.ScaleRule == nil {
return httperrors.NewInputParameterError("no found rule setting ")
}
if len(*input.ScaleRule) == 0 {
return httperrors.NewInputParameterError("no found customize monitor rule")
}
for _, scale := range *input.ScaleRule {
if scale.Database == "" {
return httperrors.NewInputParameterError("database is missing")
}
if scale.Measurement == "" {
return httperrors.NewInputParameterError("measurement is missing")
}
if scale.Field == "" {
return httperrors.NewInputParameterError("field is missing")
}
if !utils.IsInStringArray(getQueryEvalType(scale), validators.EvaluatorDefaultTypes) {
return httperrors.NewInputParameterError("the evalType is illegal")
}
if scale.Threshold == 0 {
return httperrors.NewInputParameterError("threshold is meaningless")
}
}
return nil
}
func getQueryEvalType(scale monitor.Scale) string {
typ := ""
switch scale.EvalType {
case ">=", ">":
typ = "gt"
case "<=", "<":
typ = "lt"
}
return typ
}
func (rule *ScaleDown) DoSuggestSysRule(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) {
doSuggestSysRule(ctx, userCred, isStart, rule)
}
func (rule *ScaleDown) Run(instance *monitor.SSuggestSysAlertSetting) {
oldAlert, err := getLastAlerts(rule)
if err != nil {
log.Errorln(err)
return
}
newAlert, err := rule.getLatestAlerts(instance)
if err != nil {
log.Errorln(err)
return
}
DealAlertData(oldAlert, newAlert.Value())
}
func (rule *ScaleDown) getLatestAlerts(instance *monitor.SSuggestSysAlertSetting) (*jsonutils.JSONArray, error) {
//scaleEvalMatchs := make([]*monitor.EvalMatch, 0)
firing, evalMatchMap, err := rule.getScaleEvalResult(*instance.ScaleRule)
if err != nil {
return jsonutils.NewArray(), errors.Wrap(err, "rule getScaleEvalResult happen error")
}
if firing {
serverArr, err := rule.getResourcesByEvalMatchsMap(evalMatchMap, instance)
if err != nil {
return jsonutils.NewArray(), errors.Wrap(err, "rule getResource error")
}
return serverArr, nil
}
return jsonutils.NewArray(), nil
}
func (rule *ScaleDown) getScaleEvalResult(scales []monitor.Scale) (bool, map[string][]*monitor.EvalMatch, error) {
firing := false
scaleEvalMatchs := make(map[string][]*monitor.EvalMatch, 0)
for index, scale := range scales {
condition := monitor.AlertCondition{
Type: "query",
Query: rule.newAlertQuery(scale),
Evaluator: monitor.Condition{Type: getQueryEvalType(scale), Params: []float64{scale.Threshold}},
Reducer: monitor.Condition{Type: "avg"},
Operator: scale.Operator,
}
factory := alerting.GetConditionFactories()[condition.Type]
queryCondition, err := factory(&condition, index)
if err != nil {
return firing, scaleEvalMatchs, errors.Wrapf(err, "construct query condition %s",
jsonutils.Marshal(condition))
}
//evalContext := alerting.NewEvalContext(context.Background(), auth.AdminCredential(), nil)
evalContext := alerting.EvalContext{
Ctx: context.Background(),
UserCred: auth.AdminCredential(),
IsDebug: true,
IsTestRun: true,
}
conditionResult, err := queryCondition.Eval(&evalContext)
if err != nil {
return firing, scaleEvalMatchs, errors.Wrap(err, "condition eval error")
}
if index == 0 {
firing = conditionResult.Firing
}
// calculating Firing based on operator
if conditionResult.Operator == "or" {
firing = firing || conditionResult.Firing
} else {
firing = firing && conditionResult.Firing
}
if firing {
evalMatchs := conditionResult.EvalMatches
if conditionResult.Operator == "and" {
if index != 0 {
evalMatchs = getAndEvalMatches(scaleEvalMatchs, evalMatchs)
if len(evalMatchs) == 0 {
return false, scaleEvalMatchs, nil
}
}
}
key := fmt.Sprintf("%s--%d", scale.Field, index)
scaleEvalMatchs[key] = evalMatchs
}
}
return firing, scaleEvalMatchs, nil
}
func (rule *ScaleDown) getResourcesByEvalMatchsMap(evalMatchsMap map[string][]*monitor.EvalMatch, instance *monitor.SSuggestSysAlertSetting) (*jsonutils.JSONArray, error) {
matchLength := 0
var maxEvalMatch []*monitor.EvalMatch
for _, evalMatchs := range evalMatchsMap {
if len(evalMatchs) > matchLength {
matchLength = len(evalMatchs)
maxEvalMatch = evalMatchs
}
}
serverArr := jsonutils.NewArray()
for _, evalMatch := range maxEvalMatch {
server, mappingId, mappingVal := getServerFromEvalMatch(evalMatch)
if mappingId == "" {
continue
}
suggestSysAlert, err := getSuggestSysAlertFromJson(server, rule)
if err != nil {
return serverArr, errors.Wrap(err, "Scale getSuggestSysAlertFromJson error")
}
suggestSysAlert.Action = monitor.SCALE_DOWN_DRIVER_ACTION
suggestSysAlert.MonitorConfig = jsonutils.Marshal(instance)
suggestSysAlert.Problem = describeEvalResultTojson(evalMatchsMap, mappingId, mappingVal)
serverArr.Add(jsonutils.Marshal(suggestSysAlert))
}
return serverArr, nil
}
func getServerFromEvalMatch(evalMatch *monitor.EvalMatch) (jsonutils.JSONObject, string, string) {
idTag := getMetricIdTag(evalMatch.Tags)
var server jsonutils.JSONObject
mappingId := ""
mappingVal := ""
for id, val := range idTag {
serverobj, err := getVm(val)
if err != nil {
continue
}
server = serverobj
mappingId = id
mappingVal = val
break
}
return server, mappingId, mappingVal
}
func describeEvalResultTojson(evalMatchsMap map[string][]*monitor.EvalMatch, mappingId, mappingVal string) jsonutils.JSONObject {
problem := jsonutils.NewDict()
for _, evalMatchs := range evalMatchsMap {
for _, evalMatch := range evalMatchs {
idTag := getMetricIdTag(evalMatch.Tags)
if val, ok := idTag[mappingId]; ok {
if val == mappingVal {
problem.Add(jsonutils.NewFloat(*evalMatch.Value), evalMatch.Metric)
}
}
}
}
return problem
}
func getVm(id string) (jsonutils.JSONObject, error) {
session := auth.GetAdminSession(context.Background(), "", "")
query := jsonutils.NewDict()
query.Add(jsonutils.NewString("0"), "limit")
query.Add(jsonutils.NewString("system"), "scope")
server, err := modules.Servers.GetById(session, id, query)
if err != nil {
return nil, err
}
return server, nil
}
func getMetricIdTag(tags map[string]string) map[string]string {
idTags := make(map[string]string, 0)
for key, val := range tags {
if strings.HasSuffix(key, "_id") {
idTags[key] = val
}
}
return idTags
}
func getAndEvalMatches(scaleEvalMatchs map[string][]*monitor.EvalMatch, andscaleEvalMatchs []*monitor.EvalMatch) []*monitor.EvalMatch {
for key, evalMatchs := range scaleEvalMatchs {
andscaleEvalMatchs = getAndEvalMatches_(evalMatchs, andscaleEvalMatchs)
if len(andscaleEvalMatchs) == 0 {
return andscaleEvalMatchs
}
scaleEvalMatchs[key] = getAndEvalMatches_(andscaleEvalMatchs, evalMatchs)
}
return andscaleEvalMatchs
}
//by first param to scale other param's length
func getAndEvalMatches_(scaleEvalMatchs, andscaleEvalMatchs []*monitor.EvalMatch) []*monitor.EvalMatch {
resEvalMatchs := make([]*monitor.EvalMatch, 0)
for _, evalMatch := range scaleEvalMatchs {
idTags := getMetricIdTag(evalMatch.Tags)
for _, andEvalMatch := range andscaleEvalMatchs {
andIdTags := getMetricIdTag(evalMatch.Tags)
for key, val := range idTags {
if andVal, ok := andIdTags[key]; ok {
if val == andVal {
resEvalMatchs = append(resEvalMatchs, andEvalMatch)
break
}
}
}
}
}
return resEvalMatchs
}
func (rule *ScaleDown) newAlertQuery(scale monitor.Scale) monitor.AlertQuery {
suggestSysRules, _ := models.SuggestSysRuleManager.GetRules(rule.GetType())
datasource, _ := models.DataSourceManager.GetDefaultSource()
return monitor.AlertQuery{
Model: newMetricQuery(scale),
DataSourceId: datasource.Id,
From: suggestSysRules[0].TimeFrom,
To: "now",
}
}
func newMetricQuery(scale monitor.Scale) monitor.MetricQuery {
sels := make([]monitor.MetricQuerySelect, 0)
sels = append(sels, monitor.NewMetricQuerySelect(monitor.MetricQueryPart{Type: "field", Params: []string{scale.Field}}))
return monitor.MetricQuery{
Database: scale.Database,
Measurement: scale.Measurement,
Selects: sels,
GroupBy: []monitor.MetricQueryPart{
{
Type: "field",
Params: []string{"*"},
},
},
}
}
func (rule *ScaleDown) StartResolveTask(ctx context.Context, userCred mcclient.TokenCredential,
suggestSysAlert *models.SSuggestSysAlert,
params *jsonutils.JSONDict) error {
log.Println("scaleDown StartResolveTask do nothing")
return nil
}
func (rule *ScaleDown) Resolve(data *models.SSuggestSysAlert) error {
log.Println("scaleDown Resolve do nothing")
return nil
}