diff --git a/pkg/apis/monitor/suggestsys_const.go b/pkg/apis/monitor/suggestsys_const.go index cfcc4527e9..15fbdfe6d3 100644 --- a/pkg/apis/monitor/suggestsys_const.go +++ b/pkg/apis/monitor/suggestsys_const.go @@ -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("缩减机器配置") ) diff --git a/pkg/apis/monitor/suggestsysrule.go b/pkg/apis/monitor/suggestsysrule.go index 219dfebc79..921d7201cd 100644 --- a/pkg/apis/monitor/suggestsysrule.go +++ b/pkg/apis/monitor/suggestsysrule.go @@ -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"` +} diff --git a/pkg/mcclient/modules/mod_suggestsysrule.go b/pkg/mcclient/modules/mod_suggestsysrule.go index 593b5ba759..0082c94b44 100644 --- a/pkg/mcclient/modules/mod_suggestsysrule.go +++ b/pkg/mcclient/modules/mod_suggestsysrule.go @@ -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, + } +} diff --git a/pkg/mcclient/modules/monitor/suggestsysrule.go b/pkg/mcclient/modules/monitor/suggestsysrule.go index fefe76b683..aaea27585d 100644 --- a/pkg/mcclient/modules/monitor/suggestsysrule.go +++ b/pkg/mcclient/modules/monitor/suggestsysrule.go @@ -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, + } +} diff --git a/pkg/mcclient/options/monitor/influxdbshema.go b/pkg/mcclient/options/monitor/influxdbshema.go new file mode 100644 index 0000000000..fc46638b6f --- /dev/null +++ b/pkg/mcclient/options/monitor/influxdbshema.go @@ -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 +} diff --git a/pkg/mcclient/options/monitor/suggestsysrule.go b/pkg/mcclient/options/monitor/suggestsysrule.go index b100743e52..1f2f549011 100644 --- a/pkg/mcclient/options/monitor/suggestsysrule.go +++ b/pkg/mcclient/options/monitor/suggestsysrule.go @@ -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 } diff --git a/pkg/monitor/alerting/conditions/query.go b/pkg/monitor/alerting/conditions/query.go index 2abb9c06d7..11e2c47841 100644 --- a/pkg/monitor/alerting/conditions/query.go +++ b/pkg/monitor/alerting/conditions/query.go @@ -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`. diff --git a/pkg/monitor/alerting/rule.go b/pkg/monitor/alerting/rule.go index 08341c5728..ff8eb67216 100644 --- a/pkg/monitor/alerting/rule.go +++ b/pkg/monitor/alerting/rule.go @@ -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 +} diff --git a/pkg/monitor/models/datasource.go b/pkg/monitor/models/datasource.go index 76ad9d2f44..c38bbd6a95 100644 --- a/pkg/monitor/models/datasource.go +++ b/pkg/monitor/models/datasource.go @@ -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 +} diff --git a/pkg/monitor/models/suggestsysalart.go b/pkg/monitor/models/suggestsysalart.go index f87eb478ce..f502d35685 100644 --- a/pkg/monitor/models/suggestsysalart.go +++ b/pkg/monitor/models/suggestsysalart.go @@ -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( diff --git a/pkg/monitor/models/suggestsysrule.go b/pkg/monitor/models/suggestsysrule.go index 7049914d31..01501ea986 100644 --- a/pkg/monitor/models/suggestsysrule.go +++ b/pkg/monitor/models/suggestsysrule.go @@ -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() diff --git a/pkg/monitor/suggestsysdrivers/driverconfig.go b/pkg/monitor/suggestsysdrivers/driverconfig.go index 9e37441ae0..ceab887c71 100644 --- a/pkg/monitor/suggestsysdrivers/driverconfig.go +++ b/pkg/monitor/suggestsysdrivers/driverconfig.go @@ -25,7 +25,7 @@ import ( ) func init() { - models.RegisterSuggestSysRuleDrivers(NewEIPUsedDriver(), NewDiskUnusedDriver(), NewLBUnusedDriver()) + models.RegisterSuggestSysRuleDrivers(NewEIPUsedDriver(), NewDiskUnusedDriver(), NewLBUnusedDriver(), NewScaleDownDriver()) } func InitSuggestSysRuleCronjob() { diff --git a/pkg/monitor/suggestsysdrivers/scaledown.go b/pkg/monitor/suggestsysdrivers/scaledown.go new file mode 100644 index 0000000000..d474cc7e25 --- /dev/null +++ b/pkg/monitor/suggestsysdrivers/scaledown.go @@ -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 +}