From e5f01ad535336f3eb4dd75ef0c05fa56934d77e4 Mon Sep 17 00:00:00 2001 From: Zexi Li Date: Mon, 16 May 2022 10:52:29 +0800 Subject: [PATCH] feat(monitor): VMs of host loadbalance --- cmd/climc/shell/events/events.go | 9 + cmd/climc/shell/monitor/alert.go | 14 +- cmd/climc/shell/monitor/alertdashboard.go | 3 +- cmd/climc/shell/monitor/alertpanel.go | 3 +- cmd/climc/shell/monitor/alertrecord.go | 3 +- cmd/climc/shell/monitor/alertrecordshield.go | 3 +- cmd/climc/shell/monitor/alertresource.go | 2 +- cmd/climc/shell/monitor/common.go | 7 + cmd/climc/shell/monitor/commonalert.go | 3 +- cmd/climc/shell/monitor/commonalertmetric.go | 3 +- .../shell/monitor/commonalertmetricfield.go | 3 +- cmd/climc/shell/monitor/migration_alert.go | 28 + cmd/climc/shell/monitor/monitor_resource.go | 3 +- .../shell/monitor/monitor_resourece_alert.go | 2 +- pkg/apis/monitor/migration_alert.go | 264 +++++++ pkg/apis/monitor/nodealert.go | 11 +- pkg/apis/monitor/notification.go | 13 +- pkg/mcclient/modules/mod_logs.go | 4 + pkg/mcclient/modules/monitor/alert.go | 11 +- pkg/mcclient/modules/monitor/helper.go | 13 + .../modules/monitor/mod_migration_alert.go | 45 ++ .../options/monitor/migration_alert.go | 92 +++ pkg/monitor/alerting/notifier.go | 11 +- .../alerting/notifiers/auto_migration.go | 102 +++ pkg/monitor/alerting/scheduler.go | 1 - pkg/monitor/controller/balancer/balancer.go | 671 ++++++++++++++++++ .../controller/balancer/balancer_test.go | 114 +++ pkg/monitor/controller/balancer/cpu.go | 187 +++++ pkg/monitor/controller/balancer/doc.go | 1 + pkg/monitor/controller/balancer/mem.go | 171 +++++ pkg/monitor/controller/balancer/recorder.go | 169 +++++ pkg/monitor/controller/balancer/utils.go | 109 +++ pkg/monitor/models/alert.go | 25 +- pkg/monitor/models/alertnotification.go | 7 +- pkg/monitor/models/balancerule.go | 18 + pkg/monitor/models/commonalert.go | 31 +- pkg/monitor/models/migration_alert.go | 327 +++++++++ pkg/monitor/models/nodealert.go | 2 +- pkg/monitor/models/notification.go | 21 + pkg/monitor/service/handlers.go | 1 + pkg/monitor/service/service.go | 8 + pkg/monitor/tsdb/driver/influxdb/models.go | 10 +- .../tsdb/driver/influxdb/response_parser.go | 7 +- 43 files changed, 2461 insertions(+), 71 deletions(-) create mode 100644 cmd/climc/shell/monitor/migration_alert.go create mode 100644 pkg/apis/monitor/migration_alert.go create mode 100644 pkg/mcclient/modules/monitor/mod_migration_alert.go create mode 100644 pkg/mcclient/options/monitor/migration_alert.go create mode 100644 pkg/monitor/alerting/notifiers/auto_migration.go create mode 100644 pkg/monitor/controller/balancer/balancer.go create mode 100644 pkg/monitor/controller/balancer/balancer_test.go create mode 100644 pkg/monitor/controller/balancer/cpu.go create mode 100644 pkg/monitor/controller/balancer/doc.go create mode 100644 pkg/monitor/controller/balancer/mem.go create mode 100644 pkg/monitor/controller/balancer/recorder.go create mode 100644 pkg/monitor/controller/balancer/utils.go create mode 100644 pkg/monitor/models/balancerule.go create mode 100644 pkg/monitor/models/migration_alert.go diff --git a/cmd/climc/shell/events/events.go b/cmd/climc/shell/events/events.go index e0e4b2000b..43642b612c 100644 --- a/cmd/climc/shell/events/events.go +++ b/cmd/climc/shell/events/events.go @@ -70,6 +70,10 @@ func doIdentityEventList(s *mcclient.ClientSession, args *EventListOptions) erro return DoEventList(modules.IdentityLogs, s, args) } +func doMonitorEventList(s *mcclient.ClientSession, args *EventListOptions) error { + return DoEventList(modules.MonitorLogs, s, args) +} + func DoEventList(man modulebase.ResourceManager, s *mcclient.ClientSession, args *EventListOptions) error { params := jsonutils.NewDict() if len(args.Type) > 0 { @@ -253,4 +257,9 @@ func init() { nargs := EventListOptions{BaseEventListOptions: args.BaseEventListOptions, Id: args.ID, Type: []string{"credential"}} return doIdentityEventList(s, &nargs) }) + + R(&TypeEventListOptions{}, "monitor-migrationalert-event", "Show operation event logs of monitor auto migrations", func(s *mcclient.ClientSession, args *TypeEventListOptions) error { + nargs := EventListOptions{BaseEventListOptions: args.BaseEventListOptions, Id: args.ID, Type: []string{"migrationalert"}} + return doMonitorEventList(s, &nargs) + }) } diff --git a/cmd/climc/shell/monitor/alert.go b/cmd/climc/shell/monitor/alert.go index 5a8a4331cc..98676158c3 100644 --- a/cmd/climc/shell/monitor/alert.go +++ b/cmd/climc/shell/monitor/alert.go @@ -15,9 +15,11 @@ package monitor import ( + "encoding/json" "fmt" - "yunion.io/x/onecloud/cmd/climc/shell" + "yunion.io/x/pkg/errors" + monitorapi "yunion.io/x/onecloud/pkg/apis/monitor" "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/onecloud/pkg/mcclient/modules/monitor" @@ -25,7 +27,7 @@ import ( ) func init() { - cmd := shell.NewResourceCmd(monitor.Alerts) + cmd := NewResourceCmd(monitor.Alerts) cmd.List(new(options.AlertListOptions)) cmd.Create(new(options.AlertCreateOptions)) cmd.Show(new(options.AlertShowOptions)) @@ -42,9 +44,13 @@ func init() { } ret, err := monitor.Alerts.DoTestRun(s, args.ID, data) if err != nil { - return err + return errors.Wrap(err, "DoTestRun") } - fmt.Println(ret.JSON(ret).YAMLString()) + jsonBytes, err := json.MarshalIndent(ret, "", " ") + if err != nil { + return errors.Wrap(err, "MarshalIndent") + } + fmt.Printf("%s\n", jsonBytes) return nil }) } diff --git a/cmd/climc/shell/monitor/alertdashboard.go b/cmd/climc/shell/monitor/alertdashboard.go index 5293964108..ad83536475 100644 --- a/cmd/climc/shell/monitor/alertdashboard.go +++ b/cmd/climc/shell/monitor/alertdashboard.go @@ -15,13 +15,12 @@ package monitor import ( - "yunion.io/x/onecloud/cmd/climc/shell" modules "yunion.io/x/onecloud/pkg/mcclient/modules/monitor" options "yunion.io/x/onecloud/pkg/mcclient/options/monitor" ) func init() { - cmd := shell.NewResourceCmd(modules.AlertDashBoardManager) + cmd := NewResourceCmd(modules.AlertDashBoardManager) cmd.Create(new(options.AlertDashBoardCreateOptions)) cmd.List(new(options.AlertDashBoardListOptions)) cmd.Show(new(options.AlertDashBoardShowOptions)) diff --git a/cmd/climc/shell/monitor/alertpanel.go b/cmd/climc/shell/monitor/alertpanel.go index 3eb0dcae4b..169f42644d 100644 --- a/cmd/climc/shell/monitor/alertpanel.go +++ b/cmd/climc/shell/monitor/alertpanel.go @@ -15,13 +15,12 @@ package monitor import ( - "yunion.io/x/onecloud/cmd/climc/shell" modules "yunion.io/x/onecloud/pkg/mcclient/modules/monitor" options "yunion.io/x/onecloud/pkg/mcclient/options/monitor" ) func init() { - cmd := shell.NewResourceCmd(modules.AlertPanelManager) + cmd := NewResourceCmd(modules.AlertPanelManager) cmd.Create(new(options.AlertPanelCreateOptions)) cmd.List(new(options.AlertPanelListOptions)) cmd.Show(new(options.AlertPanelShowOptions)) diff --git a/cmd/climc/shell/monitor/alertrecord.go b/cmd/climc/shell/monitor/alertrecord.go index eade48caf9..07d364ab51 100644 --- a/cmd/climc/shell/monitor/alertrecord.go +++ b/cmd/climc/shell/monitor/alertrecord.go @@ -15,13 +15,12 @@ package monitor import ( - "yunion.io/x/onecloud/cmd/climc/shell" modules "yunion.io/x/onecloud/pkg/mcclient/modules/monitor" options "yunion.io/x/onecloud/pkg/mcclient/options/monitor" ) func init() { - cmd := shell.NewResourceCmd(modules.AlertRecordManager) + cmd := NewResourceCmd(modules.AlertRecordManager) cmd.List(new(options.AlertRecordListOptions)) cmd.Show(new(options.AlertRecordShowOptions)) cmd.Get("", new(options.AlertRecordTotalOptions)) diff --git a/cmd/climc/shell/monitor/alertrecordshield.go b/cmd/climc/shell/monitor/alertrecordshield.go index cb9235c5c8..70627c1c0d 100644 --- a/cmd/climc/shell/monitor/alertrecordshield.go +++ b/cmd/climc/shell/monitor/alertrecordshield.go @@ -15,13 +15,12 @@ package monitor import ( - "yunion.io/x/onecloud/cmd/climc/shell" modules "yunion.io/x/onecloud/pkg/mcclient/modules/monitor" options "yunion.io/x/onecloud/pkg/mcclient/options/monitor" ) func init() { - cmd := shell.NewResourceCmd(modules.AlertRecordShieldManager) + cmd := NewResourceCmd(modules.AlertRecordShieldManager) cmd.Create(new(options.AlertRecordShieldCreateOptions)) cmd.List(new(options.AlertRecordShieldListOptions)) cmd.Show(new(options.AlertRecordShieldShowOptions)) diff --git a/cmd/climc/shell/monitor/alertresource.go b/cmd/climc/shell/monitor/alertresource.go index dc3adf6779..45e3b23860 100644 --- a/cmd/climc/shell/monitor/alertresource.go +++ b/cmd/climc/shell/monitor/alertresource.go @@ -21,7 +21,7 @@ import ( ) func init() { - cmd := shell.NewResourceCmd(monitor.AlertResources) + cmd := NewResourceCmd(monitor.AlertResources) cmd.List(new(options.AlertResourceListOptions)) cmd.Show(new(options.AlertResourceShowOptions)) cmd.BatchDelete(new(options.AlertResourceDeleteOptions)) diff --git a/cmd/climc/shell/monitor/common.go b/cmd/climc/shell/monitor/common.go index 9299db260d..9894da898d 100644 --- a/cmd/climc/shell/monitor/common.go +++ b/cmd/climc/shell/monitor/common.go @@ -16,6 +16,7 @@ package monitor import ( "yunion.io/x/onecloud/cmd/climc/shell" + "yunion.io/x/onecloud/pkg/mcclient/modulebase" "yunion.io/x/onecloud/pkg/util/printutils" ) @@ -25,3 +26,9 @@ var ( printObject = printutils.PrintJSONObject printBatchResults = printutils.PrintJSONBatchResults ) + +func NewResourceCmd(manager modulebase.IBaseManager) *shell.ResourceCmd { + cmd := shell.NewResourceCmd(manager) + cmd.SetPrefix("monitor") + return cmd +} diff --git a/cmd/climc/shell/monitor/commonalert.go b/cmd/climc/shell/monitor/commonalert.go index 9df940b5bb..32b3b1f56a 100644 --- a/cmd/climc/shell/monitor/commonalert.go +++ b/cmd/climc/shell/monitor/commonalert.go @@ -15,13 +15,12 @@ package monitor import ( - "yunion.io/x/onecloud/cmd/climc/shell" modules "yunion.io/x/onecloud/pkg/mcclient/modules/monitor" options "yunion.io/x/onecloud/pkg/mcclient/options/monitor" ) func init() { - cmd := shell.NewResourceCmd(modules.CommonAlertManager) + cmd := NewResourceCmd(modules.CommonAlertManager) cmd.List(new(options.CommonAlertListOptions)) cmd.Show(new(options.CommonAlertShowOptions)) cmd.Perform("enable", &options.CommonAlertShowOptions{}) diff --git a/cmd/climc/shell/monitor/commonalertmetric.go b/cmd/climc/shell/monitor/commonalertmetric.go index d6dc5634dc..27cde5e0b2 100644 --- a/cmd/climc/shell/monitor/commonalertmetric.go +++ b/cmd/climc/shell/monitor/commonalertmetric.go @@ -15,13 +15,12 @@ package monitor import ( - "yunion.io/x/onecloud/cmd/climc/shell" modules "yunion.io/x/onecloud/pkg/mcclient/modules/monitor" options "yunion.io/x/onecloud/pkg/mcclient/options/monitor" ) func init() { - cmd := shell.NewResourceCmd(modules.MetricManager) + cmd := NewResourceCmd(modules.MetricManager) cmd.List(new(options.MonitorMetricListOptions)) cmd.Update(new(options.MetricUpdateOptions)) cmd.Show(new(options.MetricShowOptions)) diff --git a/cmd/climc/shell/monitor/commonalertmetricfield.go b/cmd/climc/shell/monitor/commonalertmetricfield.go index 4a8024b7aa..8634bd4886 100644 --- a/cmd/climc/shell/monitor/commonalertmetricfield.go +++ b/cmd/climc/shell/monitor/commonalertmetricfield.go @@ -15,13 +15,12 @@ package monitor import ( - "yunion.io/x/onecloud/cmd/climc/shell" modules "yunion.io/x/onecloud/pkg/mcclient/modules/monitor" options "yunion.io/x/onecloud/pkg/mcclient/options/monitor" ) func init() { - cmd := shell.NewResourceCmd(modules.MetricFieldManager) + cmd := NewResourceCmd(modules.MetricFieldManager) cmd.List(new(options.MonitorMetricFieldListOptions)) cmd.Update(new(options.MetricFieldUpdateOptions)) cmd.Show(new(options.MetricFieldShowOptions)) diff --git a/cmd/climc/shell/monitor/migration_alert.go b/cmd/climc/shell/monitor/migration_alert.go new file mode 100644 index 0000000000..3e267269ab --- /dev/null +++ b/cmd/climc/shell/monitor/migration_alert.go @@ -0,0 +1,28 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package monitor + +import ( + modules "yunion.io/x/onecloud/pkg/mcclient/modules/monitor" + options "yunion.io/x/onecloud/pkg/mcclient/options/monitor" +) + +func init() { + cmd := NewResourceCmd(modules.MigrationAlertManager) + cmd.List(new(options.MigrationAlertListOptions)) + cmd.Show(new(options.MigrationAlertShowOptions)) + cmd.Create(new(options.MigrationAlertCreateOptions)) + cmd.Delete(new(options.MigrationAlertShowOptions)) +} diff --git a/cmd/climc/shell/monitor/monitor_resource.go b/cmd/climc/shell/monitor/monitor_resource.go index 566eda8da9..0dbbd09743 100644 --- a/cmd/climc/shell/monitor/monitor_resource.go +++ b/cmd/climc/shell/monitor/monitor_resource.go @@ -15,12 +15,11 @@ package monitor import ( - "yunion.io/x/onecloud/cmd/climc/shell" modules "yunion.io/x/onecloud/pkg/mcclient/modules/monitor" options "yunion.io/x/onecloud/pkg/mcclient/options/monitor" ) func init() { - cmd := shell.NewResourceCmd(modules.MonitorResourceManager) + cmd := NewResourceCmd(modules.MonitorResourceManager) cmd.Get("", new(options.MonitorResourceJointAlertOptions)) } diff --git a/cmd/climc/shell/monitor/monitor_resourece_alert.go b/cmd/climc/shell/monitor/monitor_resourece_alert.go index 9f5ab5b7f2..ff42266d94 100644 --- a/cmd/climc/shell/monitor/monitor_resourece_alert.go +++ b/cmd/climc/shell/monitor/monitor_resourece_alert.go @@ -21,6 +21,6 @@ import ( ) func init() { - cmd := shell.NewJointCmd(modules.MonitorResourceAlertManager) + cmd := shell.NewJointCmd(modules.MonitorResourceAlertManager).SetPrefix("monitor") cmd.List(new(options.MonitorResourceAlertListOptions)) } diff --git a/pkg/apis/monitor/migration_alert.go b/pkg/apis/monitor/migration_alert.go new file mode 100644 index 0000000000..b4afab6480 --- /dev/null +++ b/pkg/apis/monitor/migration_alert.go @@ -0,0 +1,264 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package monitor + +import ( + "fmt" + "strings" + "sync" + "time" + + "yunion.io/x/jsonutils" +) + +func init() { + for _, drv := range []IMigrationAlertMetric{ + newMMemAailable(), + newMCPUUsageActive(), + } { + GetMigrationAlertMetricDrivers().register(drv) + } +} + +var ( + migMetricDrivers *MigrationMetricDrivers +) + +func GetMigrationAlertMetricDrivers() *MigrationMetricDrivers { + if migMetricDrivers == nil { + migMetricDrivers = newMigrationMetricDrivers() + } + return migMetricDrivers +} + +type MigrationMetricDrivers struct { + *sync.Map +} + +func newMigrationMetricDrivers() *MigrationMetricDrivers { + return &MigrationMetricDrivers{ + Map: new(sync.Map), + } +} + +func (d *MigrationMetricDrivers) register(di IMigrationAlertMetric) *MigrationMetricDrivers { + d.Store(di.GetType(), di) + return d +} + +func (d *MigrationMetricDrivers) Get(t MigrationAlertMetricType) (IMigrationAlertMetric, error) { + obj, ok := d.Load(t) + if !ok { + return nil, fmt.Errorf("driver type %q not found", t) + } + return obj.(IMigrationAlertMetric), nil +} + +type MigrationAlertMetricType string + +const ( + MigrationAlertMetricTypeCPUUsageActive = "cpu.usage_active" + MigrationAlertMetricTypeMemAvailable = "mem.available" +) + +type MetricQueryFields struct { + ResourceType MigrationAlertResourceType + Database string + Measurement string + Field string + Comparator string +} + +type IMigrationAlertMetric interface { + GetType() MigrationAlertMetricType + GetQueryFields() *MetricQueryFields +} + +type mMemAvailable struct{} + +func newMMemAailable() IMigrationAlertMetric { + return new(mMemAvailable) +} + +func (_ mMemAvailable) GetType() MigrationAlertMetricType { + return MigrationAlertMetricTypeMemAvailable +} + +func (_ mMemAvailable) GetQueryFields() *MetricQueryFields { + return &MetricQueryFields{ + ResourceType: MigrationAlertResourceTypeHost, + Database: METRIC_DATABASE_TELE, + Measurement: "mem", + Field: "available", + Comparator: ConditionLessThan, // < + } +} + +type mCPUUsageActive struct{} + +func newMCPUUsageActive() IMigrationAlertMetric { + return new(mCPUUsageActive) +} + +func (_ mCPUUsageActive) GetType() MigrationAlertMetricType { + return MigrationAlertMetricTypeCPUUsageActive +} + +func (_ mCPUUsageActive) GetQueryFields() *MetricQueryFields { + return &MetricQueryFields{ + ResourceType: MigrationAlertResourceTypeHost, + Database: METRIC_DATABASE_TELE, + Measurement: "cpu", + Field: "usage_active", + Comparator: ConditionGreaterThan, // > + } +} + +func IsValidMigrationAlertMetricType(t MigrationAlertMetricType) error { + _, err := GetMigrationAlertMetricDrivers().Get(t) + return err +} + +type MigrationAlertResourceType string + +const ( + MigrationAlertResourceTypeHost = "host" +) + +type MigrationAlertCreateInput struct { + AlertCreateInput + // Threshold is the value to trigger migration + Threshold float64 `json:"threshold"` + // Period of querying metrics + Period string `json:"period"` + // MetricType is supported metric type by auto migration + MetricType MigrationAlertMetricType `json:"metric_type"` + // MigrationAlertSettings contain migration configuration + MigrationSettings *MigrationAlertSettings `json:"migration_settings"` +} + +func (m MigrationAlertCreateInput) GetMetricDriver() IMigrationAlertMetric { + d, err := GetMigrationAlertMetricDrivers().Get(m.MetricType) + if err != nil { + panic(err) + } + return d +} + +func (m MigrationAlertCreateInput) ToAlertCreateInput() *AlertCreateInput { + freq, _ := time.ParseDuration(m.Period) + ret := new(AlertCreateInput) + ret.Name = m.Name + ret.Frequency = int64(freq / time.Second) + ret.Level = m.Level + ret.CustomizeConfig = jsonutils.Marshal(m.MigrationSettings) + + drv := m.GetMetricDriver() + fs := drv.GetQueryFields() + + ret.Settings = AlertSetting{ + Conditions: []AlertCondition{ + { + Type: "query", + Operator: "and", + Query: AlertQuery{ + Model: m.getQuery(fs), + From: m.Period, + To: "now", + }, + Evaluator: m.GetEvaluator(fs), + Reducer: Condition{Type: "avg"}, + }, + }, + } + + return ret +} + +func (m MigrationAlertCreateInput) GetEvaluator(fs *MetricQueryFields) Condition { + return Condition{ + Type: fs.Comparator, + Operators: nil, + Params: []float64{m.Threshold}, + } +} + +func (m MigrationAlertCreateInput) getQuery(fs *MetricQueryFields) MetricQuery { + sels := make([]MetricQuerySelect, 0) + sels = append(sels, NewMetricQuerySelect( + MetricQueryPart{ + Type: "field", + Params: []string{fs.Field}, + }, + MetricQueryPart{ + Type: "mean", + Params: nil, + }, + )) + q := MetricQuery{ + Selects: sels, + GroupBy: []MetricQueryPart{ + { + Type: "field", + Params: []string{"*"}, + }, + { + Type: "fill", + Params: []string{"null"}, + }, + }, + Measurement: fs.Measurement, + Database: fs.Database, + } + q.Tags = []MetricQueryTag{ + { + Condition: "and", + Key: "res_type", + Operator: "=", + Value: "host", + }, + } + if m.MigrationSettings != nil { + if m.MigrationSettings.Source != nil && len(m.MigrationSettings.Source.HostIds) > 0 { + ids := strings.Join(m.MigrationSettings.Source.HostIds, "|") + q.Tags = append(q.Tags, MetricQueryTag{ + Key: "host_id", + Operator: "=~", + Value: fmt.Sprintf("/%s/", ids), + }) + } + } + return q +} + +type MigrationAlertSettings struct { + Source *MigrationAlertSettingsSource `json:"source"` + Target *MigrationAlertSettingsTarget `json:"target"` +} + +type MigrationAlertSettingsSource struct { + GuestIds []string `json:"guest_ids"` + HostIds []string `json:"host_ids"` +} + +type MigrationAlertSettingsTarget struct { + HostIds []string `json:"host_ids"` +} + +type MigrationAlertListInput struct { + AlertListInput + + MetricType string `json:"metric_type"` +} diff --git a/pkg/apis/monitor/nodealert.go b/pkg/apis/monitor/nodealert.go index a1ea1a93ba..a4901d14eb 100644 --- a/pkg/apis/monitor/nodealert.go +++ b/pkg/apis/monitor/nodealert.go @@ -121,13 +121,18 @@ func (input NodeAlertCreateInput) GetEvaluator() Condition { return GetNodeAlertEvaluator(input.Comparator, input.Threshold) } +const ( + ConditionGreaterThan = "gt" + ConditionLessThan = "lt" +) + func GetNodeAlertEvaluator(comparator string, threshold float64) Condition { - typ := "gt" + typ := ConditionGreaterThan switch comparator { case ">=", ">": - typ = "gt" + typ = ConditionGreaterThan case "<=", "<": - typ = "lt" + typ = ConditionLessThan } return Condition{ Type: typ, diff --git a/pkg/apis/monitor/notification.go b/pkg/apis/monitor/notification.go index db6cb494ff..aa8f9f3112 100644 --- a/pkg/apis/monitor/notification.go +++ b/pkg/apis/monitor/notification.go @@ -31,10 +31,11 @@ var ( ) const ( - AlertNotificationTypeOneCloud = "onecloud" - AlertNotificationTypeDingding = "dingding" - AlertNotificationTypeFeishu = "feishu" - AlertNotificationTypeAutoScaling = "autoscaling" + AlertNotificationTypeOneCloud = "onecloud" + AlertNotificationTypeDingding = "dingding" + AlertNotificationTypeFeishu = "feishu" + AlertNotificationTypeAutoScaling = "autoscaling" + AlertNotificationTypeAutoMigration = "automigration" ) type NotificationCreateInput struct { @@ -103,3 +104,7 @@ type NotificationSettingFeishu struct { AppId string `json:"app_id"` AppSecret string `json:"app_secret"` } + +type NotificationSettingAutoMigration struct { + AlertId string `json:"alert_id"` +} diff --git a/pkg/mcclient/modules/mod_logs.go b/pkg/mcclient/modules/mod_logs.go index 8fe5c3def2..c8fbabde3d 100644 --- a/pkg/mcclient/modules/mod_logs.go +++ b/pkg/mcclient/modules/mod_logs.go @@ -33,6 +33,7 @@ var ( ActionLogs modulebase.ResourceManager CloudeventLogs modulebase.ResourceManager ComputeLogs modulebase.ResourceManager + MonitorLogs modulebase.ResourceManager ) func (this *LogsManager) Get(session *mcclient.ClientSession, id string, params jsonutils.JSONObject) (jsonutils.JSONObject, error) { @@ -84,6 +85,9 @@ func init() { []string{"id", "ops_time", "obj_id", "obj_type", "obj_name", "user", "user_id", "tenant", "tenant_id", "owner_tenant_id", "action", "notes"}, []string{}) ComputeLogs.SetApiVersion(mcclient.V2_API_VERSION) + MonitorLogs = NewMonitorV2Manager("event", "events", + []string{"id", "ops_time", "obj_id", "obj_type", "obj_name", "user", "tenant", "action", "notes"}, + []string{}) Logs = LogsManager{ComputeLogs} register(&Logs) diff --git a/pkg/mcclient/modules/monitor/alert.go b/pkg/mcclient/modules/monitor/alert.go index bc82582d17..1cb0b6b5c4 100644 --- a/pkg/mcclient/modules/monitor/alert.go +++ b/pkg/mcclient/modules/monitor/alert.go @@ -15,7 +15,10 @@ package monitor import ( + "encoding/json" + "yunion.io/x/jsonutils" + "yunion.io/x/pkg/errors" "yunion.io/x/onecloud/pkg/apis/monitor" "yunion.io/x/onecloud/pkg/mcclient" @@ -125,9 +128,11 @@ func (m *SAlertManager) DoCreate(s *mcclient.ClientSession, config *AlertConfig) func (m *SAlertManager) DoTestRun(s *mcclient.ClientSession, id string, input *monitor.AlertTestRunInput) (*monitor.AlertTestRunOutput, error) { ret, err := m.PerformAction(s, id, "test-run", input.JSON(input)) if err != nil { - return nil, err + return nil, errors.Wrap(err, "call test-run") } out := new(monitor.AlertTestRunOutput) - err = ret.Unmarshal(out) - return out, err + if err := json.Unmarshal([]byte(ret.String()), out); err != nil { + return nil, errors.Wrapf(err, "Unmarshal %s", ret.String()) + } + return out, nil } diff --git a/pkg/mcclient/modules/monitor/helper.go b/pkg/mcclient/modules/monitor/helper.go index ed9944437d..8d6a7fa577 100644 --- a/pkg/mcclient/modules/monitor/helper.go +++ b/pkg/mcclient/modules/monitor/helper.go @@ -16,6 +16,7 @@ package monitor import ( "fmt" + "strings" "time" "yunion.io/x/onecloud/pkg/apis/monitor" @@ -467,6 +468,18 @@ func (w *AlertQueryWhere) OR() *AlertQueryWhere { return w } +func (w *AlertQueryWhere) REGEX(key, val string) *AlertQueryWhere { + return w.filter("=~", key, fmt.Sprintf("/%s/", val)) +} + +func (w *AlertQueryWhere) IN(key string, vals []string) *AlertQueryWhere { + if len(vals) == 0 { + return w + } + valStr := strings.Join(vals, "|") + return w.REGEX(key, valStr) +} + func (w *AlertQueryWhere) newTag(op string, key string, value string) monitor.MetricQueryTag { return monitor.MetricQueryTag{ Key: key, diff --git a/pkg/mcclient/modules/monitor/mod_migration_alert.go b/pkg/mcclient/modules/monitor/mod_migration_alert.go new file mode 100644 index 0000000000..b962446efe --- /dev/null +++ b/pkg/mcclient/modules/monitor/mod_migration_alert.go @@ -0,0 +1,45 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package monitor + +import ( + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/mcclient/modulebase" + "yunion.io/x/onecloud/pkg/mcclient/modules" +) + +var ( + MigrationAlertManager *SMigrationAlertManager +) + +type SMigrationAlertManager struct { + *modulebase.ResourceManager +} + +func init() { + MigrationAlertManager = NewMigrationAlertManager() + modules.Register(MigrationAlertManager) + MigrationAlertManager.SetApiVersion(mcclient.V2_API_VERSION) + modules.RegisterV2(MigrationAlertManager) +} + +func NewMigrationAlertManager() *SMigrationAlertManager { + m := modules.NewMonitorV2Manager("migrationalert", "migrationalerts", + []string{"id", "name", "metric_type"}, + []string{}) + return &SMigrationAlertManager{ + ResourceManager: &m, + } +} diff --git a/pkg/mcclient/options/monitor/migration_alert.go b/pkg/mcclient/options/monitor/migration_alert.go new file mode 100644 index 0000000000..3506080da0 --- /dev/null +++ b/pkg/mcclient/options/monitor/migration_alert.go @@ -0,0 +1,92 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package monitor + +import ( + "strings" + + "yunion.io/x/jsonutils" + "yunion.io/x/pkg/errors" + + "yunion.io/x/onecloud/pkg/apis/monitor" + "yunion.io/x/onecloud/pkg/mcclient/options" +) + +type MigrationAlertListOptions struct { + options.BaseListOptions + MetricType string `help:"Migration alert metric type" choices:"cpu.usage_active|mem.available"` +} + +func (o *MigrationAlertListOptions) Params() (jsonutils.JSONObject, error) { + return options.ListStructToParams(o) +} + +type MigrationAlertShowOptions struct { + ID string `help:"ID of alart " json:"-"` +} + +func (o *MigrationAlertShowOptions) Params() (jsonutils.JSONObject, error) { + return options.StructToParams(o) +} + +func (o *MigrationAlertShowOptions) GetId() string { + return o.ID +} + +type MigrationAlertCreateOptions struct { + NAME string `help:"Name of the migration alert"` + METRIC string `help:"Metric type" choices:"cpu.usage_active.gt|mem.available.lt"` + THRESHOLD float64 `help:"Metric threshold"` + Period string `help:"Period of execution, e.g. '5m', '1h'" default:"5m"` + SourceHost []string `help:"Source hosts' id or name"` + SourceGuest []string `help:"Source guests's id or name"` + TargetHost []string `help:"Target hosts' id or name"` +} + +func (o *MigrationAlertCreateOptions) parseMetric(m string) (monitor.MigrationAlertMetricType, error) { + parts := strings.Split(m, ".") + if len(parts) != 3 { + return "", errors.Errorf("Invalid metric %q", m) + } + return monitor.MigrationAlertMetricType(strings.Join([]string{parts[0], parts[1]}, ".")), nil +} + +func (o *MigrationAlertCreateOptions) Params() (jsonutils.JSONObject, error) { + input := new(monitor.MigrationAlertCreateInput) + input.Name = o.NAME + input.Threshold = o.THRESHOLD + input.Period = o.Period + mt, err := o.parseMetric(o.METRIC) + if err != nil { + return nil, errors.Wrap(err, "metric_type") + } + input.MetricType = mt + + input.MigrationSettings = &monitor.MigrationAlertSettings{ + Source: &monitor.MigrationAlertSettingsSource{}, + Target: &monitor.MigrationAlertSettingsTarget{}, + } + if len(o.SourceHost) != 0 { + input.MigrationSettings.Source.HostIds = o.SourceHost + } + if len(o.SourceGuest) != 0 { + input.MigrationSettings.Source.GuestIds = o.SourceGuest + } + if len(o.TargetHost) != 0 { + input.MigrationSettings.Target.HostIds = o.TargetHost + } + + return input.JSON(input), nil +} diff --git a/pkg/monitor/alerting/notifier.go b/pkg/monitor/alerting/notifier.go index 99b5403d14..f651937ce1 100644 --- a/pkg/monitor/alerting/notifier.go +++ b/pkg/monitor/alerting/notifier.go @@ -61,7 +61,7 @@ type notifierStateSlice []*notifierState func (n *notificationService) sendNotification(evalCtx *EvalContext, state *notifierState) error { if !evalCtx.IsTestRun { if err := state.state.SetToPending(); err != nil { - return err + return errors.Wrap(err, "SetToPending") } } return n.sendAndMarkAsComplete(evalCtx, state) @@ -73,8 +73,7 @@ func (n *notificationService) sendAndMarkAsComplete(evalCtx *EvalContext, state log.Debugf("Sending notification, type %s, id %s", notifier.GetType(), notifier.GetNotifierId()) if err := notifier.Notify(evalCtx, state.state.GetParams()); err != nil { - log.Errorf("failed to send notification %s: %v", notifier.GetNotifierId(), err) - return err + return errors.Wrapf(err, "notify driver %s(%s)", notifier.GetType(), notifier.GetNotifierId()) } if evalCtx.IsTestRun { @@ -82,7 +81,7 @@ func (n *notificationService) sendAndMarkAsComplete(evalCtx *EvalContext, state } err := state.state.UpdateSendTime() if err != nil { - return errors.Wrap(err, "notifierState UpdateSendTime err") + return errors.Wrap(err, "notifierState UpdateSendTime") } return state.state.SetToCompleted() } @@ -90,7 +89,7 @@ func (n *notificationService) sendAndMarkAsComplete(evalCtx *EvalContext, state func (n *notificationService) sendNotifications(evalCtx *EvalContext, states notifierStateSlice) error { for _, state := range states { if err := n.sendNotification(evalCtx, state); err != nil { - log.Errorf("failed to send %s notification: %v", state.notifier.GetNotifierId(), err) + log.Errorf("failed to send %s(%s) notification: %v", state.notifier.GetType(), state.notifier.GetNotifierId(), err) if evalCtx.IsTestRun { return err } @@ -102,7 +101,7 @@ func (n *notificationService) sendNotifications(evalCtx *EvalContext, states not func (n *notificationService) getNeededNotifiers(nIds []string, evalCtx *EvalContext) (notifierStateSlice, error) { notis, err := models.NotificationManager.GetNotificationsWithDefault(nIds) if err != nil { - return nil, err + return nil, errors.Wrapf(err, "GetNotificationsWithDefault with %v", nIds) } var result notifierStateSlice diff --git a/pkg/monitor/alerting/notifiers/auto_migration.go b/pkg/monitor/alerting/notifiers/auto_migration.go new file mode 100644 index 0000000000..c9190b5649 --- /dev/null +++ b/pkg/monitor/alerting/notifiers/auto_migration.go @@ -0,0 +1,102 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package notifiers + +import ( + "yunion.io/x/jsonutils" + "yunion.io/x/log" + "yunion.io/x/pkg/errors" + + "yunion.io/x/onecloud/pkg/apis/monitor" + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/mcclient/auth" + "yunion.io/x/onecloud/pkg/monitor/alerting" + "yunion.io/x/onecloud/pkg/monitor/controller/balancer" + "yunion.io/x/onecloud/pkg/monitor/models" + "yunion.io/x/onecloud/pkg/monitor/options" +) + +func init() { + alerting.RegisterNotifier(&alerting.NotifierPlugin{ + Type: monitor.AlertNotificationTypeAutoMigration, + Factory: newAutoMigratingNotifier, + ValidateCreateData: func(cred mcclient.IIdentityProvider, input monitor.NotificationCreateInput) (monitor.NotificationCreateInput, error) { + settings := new(monitor.NotificationSettingAutoMigration) + if err := input.Settings.Unmarshal(settings); err != nil { + return input, errors.Wrap(err, "Unmarshal setting") + } + input.Settings = jsonutils.Marshal(settings) + return input, nil + }, + }) +} + +type autoMigrationNotifier struct { + NotifierBase + + Settings *monitor.NotificationSettingAutoMigration +} + +func newAutoMigratingNotifier(conf alerting.NotificationConfig) (alerting.Notifier, error) { + settings := new(monitor.NotificationSettingAutoMigration) + if err := conf.Settings.Unmarshal(settings); err != nil { + return nil, errors.Wrap(err, "unmarshal setting") + } + return &autoMigrationNotifier{ + NotifierBase: NewNotifierBase(conf), + Settings: settings, + }, nil +} + +func (am *autoMigrationNotifier) getBalancerRules(ctx *alerting.EvalContext, alert *models.SMigrationAlert) (*balancer.Rules, error) { + if len(ctx.EvalMatches) >= 1 { + log.Warningf("EvalMatches great than 1, use first one") + } else { + return nil, errors.Errorf("Matches not >= 1 %d", len(ctx.EvalMatches)) + } + match := ctx.EvalMatches[0] + drv, err := balancer.GetMetricDrivers().Get(alert.GetMetricType()) + if err != nil { + return nil, errors.Wrap(err, "Get metric driver") + } + log.Infof("autoMigrationNotifier for evalMatch: %s", jsonutils.Marshal(match)) + return balancer.NewRules(ctx, match, alert, drv) +} + +func (am *autoMigrationNotifier) Notify(ctx *alerting.EvalContext, data jsonutils.JSONObject) error { + if !ctx.Firing { + // do nothing + return nil + } + + alertId := ctx.Rule.Id + man := models.GetMigrationAlertManager() + obj, err := man.FetchById(alertId) + if err != nil { + return errors.Wrapf(err, "Fetch alert by Id: %q", alertId) + } + alert := obj.(*models.SMigrationAlert) + + rules, err := am.getBalancerRules(ctx, alert) + if err != nil { + return errors.Wrapf(err, "get balancer rules") + } + + if err := balancer.DoBalance(ctx.Ctx, auth.GetAdminSession(ctx.Ctx, options.Options.Region, ""), rules, balancer.NewRecorder()); err != nil { + return errors.Wrap(err, "DoBalance") + } + + return nil +} diff --git a/pkg/monitor/alerting/scheduler.go b/pkg/monitor/alerting/scheduler.go index d0ad8d9ed5..30e886ba69 100644 --- a/pkg/monitor/alerting/scheduler.go +++ b/pkg/monitor/alerting/scheduler.go @@ -35,7 +35,6 @@ func newScheduler() scheduler { func (s *schedulerImpl) Update(rules []*Rule) { log.Debugf("Scheduling update, rule count %d", len(rules)) - log.Errorf("Scheduling update, rule count %d", len(rules)) jobs := make(map[string]*Job) diff --git a/pkg/monitor/controller/balancer/balancer.go b/pkg/monitor/controller/balancer/balancer.go new file mode 100644 index 0000000000..7193526673 --- /dev/null +++ b/pkg/monitor/controller/balancer/balancer.go @@ -0,0 +1,671 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package balancer + +import ( + "context" + "fmt" + "math" + "sort" + "sync" + + "yunion.io/x/jsonutils" + "yunion.io/x/log" + "yunion.io/x/pkg/errors" + "yunion.io/x/pkg/utils" + + computeapi "yunion.io/x/onecloud/pkg/apis/compute" + "yunion.io/x/onecloud/pkg/apis/monitor" + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/mcclient/modules/compute" + compute_options "yunion.io/x/onecloud/pkg/mcclient/options/compute" + "yunion.io/x/onecloud/pkg/monitor/alerting" + "yunion.io/x/onecloud/pkg/monitor/models" + "yunion.io/x/onecloud/pkg/monitor/tsdb" +) + +func init() { + for _, drv := range []IMetricDriver{ + newMemAvailable(), + newCPUUsageActive(), + } { + GetMetricDrivers().register(drv) + } +} + +var ( + drivers *MetricDrivers +) + +type MetricDrivers struct { + *sync.Map +} + +func NewMetricDrivers() *MetricDrivers { + return &MetricDrivers{ + Map: new(sync.Map), + } +} + +func (d *MetricDrivers) register(id IMetricDriver) *MetricDrivers { + d.Store(id.GetType(), id) + return d +} + +func (d *MetricDrivers) Get(mT monitor.MigrationAlertMetricType) (IMetricDriver, error) { + drv, ok := d.Load(mT) + if !ok { + return nil, errors.Errorf("Not found driver by %q", mT) + } + return drv.(IMetricDriver), nil +} + +func GetMetricDrivers() *MetricDrivers { + if drivers == nil { + drivers = NewMetricDrivers() + } + return drivers +} + +type IMetricDriver interface { + GetType() monitor.MigrationAlertMetricType + GetTsdbQuery() *TsdbQuery + GetCandidate(gst jsonutils.JSONObject, host IHost, ds *tsdb.DataSource) (ICandidate, error) + SetHostCurrent(h IHost, values map[string]float64) error + GetTarget(host jsonutils.JSONObject) (ITarget, error) + GetCondition(s *monitor.AlertSetting) (ICondition, error) +} + +type ICondition interface { + GetThreshold() float64 + // GetSourceThresholdDelta must > 0 + GetSourceThresholdDelta(threshold float64, srcHost IHost) float64 + IsFitTarget(t ITarget, c ICandidate) error +} + +type Rules struct { + Alert *models.SMigrationAlert + Condtion ICondition + Source *SourceRule + Target *TargetRule +} + +func (r *Rules) GetAlert() *models.SMigrationAlert { + return r.Alert +} + +func NewRules(_ *alerting.EvalContext, m *monitor.EvalMatch, alert *models.SMigrationAlert, drv IMetricDriver) (*Rules, error) { + hostId, ok := m.Tags["host_id"] + if !ok { + return nil, errors.Errorf("Not found host_id in tags: %#v", m.Tags) + } + ok, hObjs := models.MonitorResourceManager.GetResourceObjByResType(monitor.METRIC_RES_TYPE_HOST) + if !ok { + return nil, errors.Errorf("GetResourceObjByResType host returns false") + } + var srcHostObj jsonutils.JSONObject = nil + for _, obj := range hObjs { + id, err := obj.GetString("id") + if err != nil { + return nil, errors.Wrapf(err, "get host obj id: %s", obj) + } + if id == hostId { + srcHostObj = obj + break + } + } + if srcHostObj == nil { + return nil, errors.Errorf("Not found source host object by id: %q, %q", hostId, srcHostObj) + } + srcHost, err := drv.GetTarget(srcHostObj) + if err != nil { + return nil, errors.Wrap(err, "new host") + } + + allHosts := []IResource{srcHost} + msettings, _ := alert.GetMigrationSettings() + targetHosts, err := filterTargetHosts(drv, srcHost, hObjs, msettings) + if err != nil { + return nil, errors.Wrap(err, "filterTargetHosts") + } + for _, oh := range targetHosts { + allHosts = append(allHosts, oh) + } + + dsObj, err := models.DataSourceManager.GetDefaultSource() + if err != nil { + return nil, errors.Wrapf(err, "Get default DataSource") + } + ds := dsObj.ToTSDBDataSource("") + cds, err := findGuestsOfHost(drv, srcHost, ds, msettings) + if err != nil { + return nil, errors.Wrapf(err, "findGuestsOfHost %s", srcHost.GetName()) + } + metrics, err := InfluxdbQuery(ds, "host_id", allHosts, drv.GetTsdbQuery()) + if err != nil { + return nil, errors.Wrapf(err, "InfluxdbQuery all hosts metrics") + } + for _, host := range allHosts { + m := metrics.Get(host.GetId()) + if m == nil { + return nil, errors.Errorf("Influxdb metrics of %s(%s) not found", host.GetName(), host.GetId()) + } + if err := drv.SetHostCurrent(host.(IHost), m.Values); err != nil { + return nil, errors.Wrapf(err, "SetHostCurrent %q", host.GetName()) + } + } + settings, err := alert.GetSettings() + if err != nil { + return nil, errors.Wrapf(err, "Get alert settings") + } + cond, err := drv.GetCondition(settings) + if err != nil { + return nil, errors.Wrapf(err, "Get Condtion") + } + rs := &Rules{ + Alert: alert, + Condtion: cond, + } + rs.Source = NewSourceRule(srcHost, cds) + rs.Target = NewTargetRule(targetHosts) + return rs, nil +} + +func filterTargetHosts(drv IMetricDriver, srcHost IHost, allHost []jsonutils.JSONObject, ms *monitor.MigrationAlertSettings) ([]ITarget, error) { + specifyHostIds := []string{} + if ms != nil && ms.Target != nil { + specifyHostIds = ms.Target.HostIds + } + srcHostId := srcHost.GetId() + srcHostObj := srcHost.GetObject() + srcHostType, err := srcHostObj.GetString("host_type") + if err != nil { + return nil, errors.Wrap(err, "get source host_type") + } + srcArch, err := srcHostObj.GetString("cpu_architecture") + if err != nil { + return nil, errors.Wrapf(err, "get source cpu_architecture") + } + targets := make([]ITarget, 0) + for _, obj := range allHost { + id, err := obj.GetString("id") + if err != nil { + return nil, errors.Wrapf(err, "get host obj id: %s", obj) + } + if id == srcHostId { + continue + } + if len(specifyHostIds) != 0 { + if !utils.IsInStringArray(id, specifyHostIds) { + continue + } + } + hostType, _ := obj.GetString("host_type") + if hostType != srcHostType { + continue + } + arch, _ := obj.GetString("cpu_architecture") + if arch != srcArch { + continue + } + th, err := drv.GetTarget(obj) + if err != nil { + return nil, errors.Wrapf(err, "drv.GetTarget %s", obj) + } + targets = append(targets, th) + } + return targets, nil +} + +func findGuestsOfHost(drv IMetricDriver, host IHost, ds *tsdb.DataSource, ms *monitor.MigrationAlertSettings) ([]ICandidate, error) { + ok, objs := models.MonitorResourceManager.GetResourceObjByResType(monitor.METRIC_RES_TYPE_GUEST) + if !ok { + return nil, errors.Errorf("GetResourceObjByResType by guest return false") + } + + specifyGuestIds := []string{} + if ms != nil && ms.Source != nil { + specifyGuestIds = ms.Source.GuestIds + } + + ret := make([]ICandidate, 0) + for _, obj := range objs { + gHostId, err := obj.GetString("host_id") + if err != nil { + return nil, errors.Wrapf(err, "get host_id from cache guest %s", obj) + } + if gHostId == host.GetId() { + status, err := obj.GetString("status") + if err != nil { + return nil, errors.Wrapf(err, "get status of guest: %s", obj) + } + name, _ := obj.GetString("name") + // filter running guest + if status != computeapi.VM_RUNNING { + log.Debugf("ignore guest %s cause status is %s", name, status) + continue + } + gId, _ := obj.GetString("id") + if len(specifyGuestIds) != 0 { + if !utils.IsInStringArray(gId, specifyGuestIds) { + log.Debugf("ignore guest %s(%s) cause not in specified ids %v", name, gId, specifyGuestIds) + continue + } + } + c, err := drv.GetCandidate(obj, host, ds) + if err != nil { + return nil, errors.Wrapf(err, "drv.GetCandidate of guest %s", obj) + } + if c.GetScore() == 0 { + log.Debugf("ignore guest %s cause %s score is 0", c.GetName(), drv.GetType()) + continue + } + ret = append(ret, c) + } + } + return ret, nil +} + +// SourceRule 定义触发了报警的宿主机和上面可以迁移的虚拟机 +type SourceRule struct { + Host IHost + Candidates []ICandidate +} + +func NewSourceRule(host IHost, cds []ICandidate) *SourceRule { + return &SourceRule{ + Host: host, + Candidates: cds, + } +} + +type IHost interface { + IResource + GetHostResource() *HostResource + GetCurrent() float64 + SetCurrent(float64) IHost + Compare(oh IHost) bool +} + +// TargetRule 定义可以选择迁移的宿主机 +type TargetRule struct { + Items []ITarget +} + +func NewTargetRule(hosts []ITarget) *TargetRule { + return &TargetRule{ + Items: hosts, + } +} + +type ItemType string + +const ( + ItemTypeHost = "host" + ItemTypeGuest = "guest" +) + +type ITarget interface { + IHost + Selected(c ICandidate) ITarget +} + +type iTargets []ITarget + +func (i iTargets) Len() int { + return len(i) +} + +func (ts iTargets) Less(i, j int) bool { + a, b := ts[i], ts[j] + return a.Compare(b) +} + +func (ts iTargets) Swap(i, j int) { + ts[i], ts[j] = ts[j], ts[i] +} + +func RecoverInProcessAlerts(ctx context.Context, s *mcclient.ClientSession) error { + alerts, err := models.GetMigrationAlertManager().GetInMigrationAlerts() + if err != nil { + return errors.Wrap(err, "GetInMigrationAlerts") + } + recorder := NewRecorder() + errs := make([]error, 0) + for _, alert := range alerts { + notes, err := alert.GetMigrateNotes() + if err != nil { + errs = append(errs, errors.Wrapf(err, "GetMigrateNotes for %s", alert.GetId())) + continue + } + for _, note := range notes { + log.Infof("Start recover alert %s(%s) note %s", alert.GetName(), alert.GetId(), jsonutils.Marshal(note)) + notePrt := ¬e + recorder.StartWatchMigratingProcess(ctx, s, alert, notePrt) + } + } + return errors.NewAggregate(errs) +} + +func DoBalance(ctx context.Context, s *mcclient.ClientSession, rules *Rules, recorder IRecorder) error { + // check whether having migration in process + alerts, err := models.GetMigrationAlertManager().GetInMigrationAlerts() + if err != nil { + return errors.Wrap(err, "GetInMigrationAlerts") + } + if len(alerts) != 0 { + ids := make([]string, len(alerts)) + for i := range alerts { + alert := alerts[i] + ids[i] = fmt.Sprintf("%s(%s)", alert.GetName(), alert.GetId()) + } + return errors.Errorf("Others migration alerts in process: %v", ids) + } + rst, err := findResult(rules) + if err != nil { + err = errors.Wrapf(err, "find result to migrate") + recorder.RecordError(s.GetToken(), rules.GetAlert(), err, EventActionFindResultFail) + return err + } + + if err := doMigrate(ctx, s, rules, rst, recorder); err != nil { + return errors.Wrapf(err, "do migrate for result %#v", rst) + } + + return nil +} + +type result struct { + pairs []*resultPair +} + +type resultPair struct { + source ICandidate + target ITarget +} + +func findResult(rules *Rules) (*result, error) { + // 找到 rules.Source 里面可以迁移的虚拟机 + guests, err := findCandidates(rules.Source, rules.Condtion) + if err != nil { + return nil, errors.Wrap(err, "find source candidates to migrate") + } + + // TODO + // 将找到的 guests 进行配对,调用 scheduler-forecast 接口判断能否迁移到宿主机 + // 如果不能迁就提出这些 guests,重新 findCandidates + + // 将找到的虚拟机分配到对应的宿主机,形成 1-1 配对 + return pairMigratResult(guests, rules.Target, rules.Condtion) +} + +type IResource interface { + GetId() string + GetName() string + GetObject() jsonutils.JSONObject +} + +type ICandidate interface { + IResource + GetHostName() string + GetScore() float64 +} + +func getCScores(css []ICandidate) []float64 { + ret := make([]float64, len(css)) + for i := range css { + ret[i] = css[i].GetScore() + } + return ret +} + +func findFitCandidates(css []ICandidate, delta float64) ([]ICandidate, error) { + if len(css) == 0 { + return nil, errors.Errorf("Not found fit input for delta %f", delta) + } + first := css[0] + rest := css[1:] + if first.GetScore() >= delta { + return []ICandidate{first}, nil + } + rRests, err := findFitCandidates(rest, delta-first.GetScore()) + if err != nil { + return nil, errors.Wrapf(err, "Found in rest %v", rest) + } + ret := []ICandidate{first} + ret = append(ret, rRests...) + return ret, nil +} + +type candidatesThresholdSort struct { + cds []ICandidate + threshold float64 +} + +func newCandidatesThresholdSort(cds []ICandidate, th float64) *candidatesThresholdSort { + return &candidatesThresholdSort{ + cds: cds, + threshold: th, + } +} + +func (s *candidatesThresholdSort) Len() int { + return len(s.cds) +} + +func (s *candidatesThresholdSort) Swap(i, j int) { + s.cds[i], s.cds[j] = s.cds[j], s.cds[i] +} + +func (s *candidatesThresholdSort) Less(i, j int) bool { + d1 := math.Abs(s.cds[i].GetScore() - s.threshold) + d2 := math.Abs(s.cds[j].GetScore() - s.threshold) + return d1 < d2 +} + +func (s *candidatesThresholdSort) CandidatesString() string { + str := fmt.Sprintf("delta: %f\n", s.threshold) + for _, c := range s.cds { + str += fmt.Sprintf("%s: %f\n", c.GetName(), c.GetScore()) + } + return str +} + +func (s *candidatesThresholdSort) Debug(prefix string) { + log.Infof("%s:\n%s", prefix, s.CandidatesString()) +} + +func sortCandidatesByThreshold(cds []ICandidate, th float64) []ICandidate { + ss := newCandidatesThresholdSort(cds, th) + // ss.Debug("Pre sort") + sort.Sort(ss) + // ss.Debug("After sort") + return ss.cds +} + +func findCandidates(src *SourceRule, cond ICondition) ([]ICandidate, error) { + threshold := cond.GetThreshold() + delta := cond.GetSourceThresholdDelta(threshold, src.Host) + cds := sortCandidatesByThreshold(src.Candidates, threshold) + return findFitCandidates(cds, delta) +} + +func findFitTarget(c ICandidate, targets iTargets, cond ICondition) (ITarget, error) { + // sort targets + sort.Sort(targets) + var errs []error + for i := range targets { + target := targets[i] + if err := cond.IsFitTarget(target, c); err == nil { + return target, nil + } else { + errs = append(errs, err) + } + } + return nil, errors.NewAggregate(errs) +} + +func pairMigratResult(gsts []ICandidate, target *TargetRule, cond ICondition) (*result, error) { + pairs := make([]*resultPair, 0) + hosts := target.Items + for _, gst := range gsts { + host, err := findFitTarget(gst, hosts, cond) + if err != nil { + return nil, errors.Wrapf(err, "not found target for guest %#v", gst) + } + host.Selected(gst) + pairs = append(pairs, &resultPair{ + source: gst, + target: host, + }) + } + if len(gsts) != len(pairs) { + return nil, errors.Errorf("Paired: %d candidates != %d hosts", len(gsts), len(pairs)) + } + return &result{ + pairs: pairs, + }, nil +} + +func doMigrate(ctx context.Context, s *mcclient.ClientSession, rules *Rules, rst *result, recorder IRecorder) error { + // migrate must be executed on by one + alert := rules.GetAlert() + for _, pair := range rst.pairs { + note, err := NewMigrateNote(pair, nil) + if err != nil { + return errors.Wrap(err, "NewMigrateNotes") + } + if obj, err := doMigrateByPair(s, pair); err != nil { + err = errors.Wrapf(err, "doMigrateByPair %s to %s", pair.source.GetName(), pair.target.GetName()) + if rErr := recorder.RecordMigrateError(s.GetToken(), alert, note, err); rErr != nil { + log.Errorf("RecordMigrate %s to %s error: %v", obj, pair.target.GetId(), rErr) + } + return err + } else { + if err := recorder.RecordMigrate(ctx, s, alert, note); err != nil { + log.Errorf("RecordMigrate %s to %s error: %v", obj, pair.target.GetId(), err) + } + } + } + return nil +} + +func doMigrateByPair(s *mcclient.ClientSession, pair *resultPair) (jsonutils.JSONObject, error) { + gst := pair.source + trueObj := true + input := &compute_options.ServerLiveMigrateOptions{ + ID: gst.GetId(), + PreferHost: pair.target.GetId(), + SkipCpuCheck: &trueObj, + SkipKernelCheck: &trueObj, + } + params, err := input.Params() + if err != nil { + return nil, errors.Wrapf(err, "live migrate input %#v", input) + } + obj, err := compute.Servers.PerformAction(s, input.GetId(), "live-migrate", params) + if err != nil { + return nil, errors.Wrapf(err, "live migrate with params: %s", params) + } + return obj, nil +} + +type resource struct { + id string + name string + obj jsonutils.JSONObject +} + +func newResource(obj jsonutils.JSONObject) (IResource, error) { + id, err := obj.GetString("id") + if err != nil { + return nil, errors.Wrap(err, "get id") + } + name, err := obj.GetString("name") + if err != nil { + return nil, errors.Wrap(err, "get name") + } + return &resource{ + id: id, + name: name, + obj: obj, + }, nil +} + +func (r *resource) GetId() string { + return r.id +} + +func (r *resource) GetName() string { + return r.name +} + +func (r *resource) GetObject() jsonutils.JSONObject { + return r.obj +} + +type guestResource struct { + IResource + hostName string + guest jsonutils.JSONObject +} + +func newGuestResource(gst jsonutils.JSONObject, hostName string) (*guestResource, error) { + res, err := newResource(gst) + if err != nil { + return nil, errors.Wrap(err, "newResource") + } + return &guestResource{ + IResource: res, + hostName: hostName, + guest: gst, + }, nil +} + +func (s *guestResource) GetHostName() string { + return s.hostName +} + +type HostResource struct { + IResource + host jsonutils.JSONObject + totalMemSize float64 + cpuCount int64 +} + +func newHostResource(host jsonutils.JSONObject) (*HostResource, error) { + res, err := newResource(host) + if err != nil { + return nil, errors.Wrap(err, "newResource") + } + memSize, err := host.Int("mem_size") + if err != nil { + return nil, errors.Wrap(err, "get mem_size") + } + cpuCount, err := host.Int("cpu_count") + if err != nil { + return nil, errors.Wrap(err, "get cpu_count") + } + return &HostResource{ + IResource: res, + host: host, + totalMemSize: float64(memSize * 1024 * 1024), + cpuCount: cpuCount, + }, nil +} + +func (h *HostResource) GetHostResource() *HostResource { + return h +} diff --git a/pkg/monitor/controller/balancer/balancer_test.go b/pkg/monitor/controller/balancer/balancer_test.go new file mode 100644 index 0000000000..1efddbbf63 --- /dev/null +++ b/pkg/monitor/controller/balancer/balancer_test.go @@ -0,0 +1,114 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package balancer + +import ( + "reflect" + "testing" + + "yunion.io/x/jsonutils" +) + +type floatC float64 + +func (f floatC) GetId() string { + return "" +} + +func (f floatC) GetName() string { + return "" +} + +func (f floatC) GetObject() jsonutils.JSONObject { + return nil +} + +func (f floatC) GetHostName() string { + return "" +} + +func (f floatC) GetScore() float64 { + return float64(f) +} + +func newFCs(n ...float64) []ICandidate { + ret := make([]ICandidate, len(n)) + for i := range n { + ret[i] = floatC(n[i]) + } + return ret +} + +func Test_findFitCandidates(t *testing.T) { + type args struct { + input []ICandidate + delta float64 + } + tests := []struct { + name string + args args + want []ICandidate + wantErr bool + }{ + { + name: "{}", + args: args{ + input: newFCs(), + delta: 3.0, + }, + want: nil, + wantErr: true, + }, + { + name: "{1, 2, 3}, 3", + args: args{ + input: newFCs(1, 2, 3), + delta: 3.0, + }, + want: newFCs(1, 2), + wantErr: false, + }, + { + name: "{1, 2, 3}, 0.5", + args: args{ + input: newFCs(1, 2, 3), + delta: 0.5, + }, + want: newFCs(1), + wantErr: false, + }, + { + name: "{1, 2, 3}, 4", + args: args{ + input: newFCs(1, 2, 3), + delta: 4, + }, + want: newFCs(1, 2, 3), + wantErr: false, + }, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got, err := findFitCandidates(tt.args.input, tt.args.delta) + if (err != nil) != tt.wantErr { + t.Errorf("findN() got = %v, error = %v, wantErr %v, err = %v", got, err, tt.wantErr, err) + return + } + if !reflect.DeepEqual(got, tt.want) { + t.Errorf("findN() = %v, want %v", got, tt.want) + } + }) + } +} diff --git a/pkg/monitor/controller/balancer/cpu.go b/pkg/monitor/controller/balancer/cpu.go new file mode 100644 index 0000000000..9539dbcef5 --- /dev/null +++ b/pkg/monitor/controller/balancer/cpu.go @@ -0,0 +1,187 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package balancer + +import ( + "yunion.io/x/jsonutils" + "yunion.io/x/log" + "yunion.io/x/pkg/errors" + + "yunion.io/x/onecloud/pkg/apis/monitor" + "yunion.io/x/onecloud/pkg/monitor/tsdb" +) + +func setHostCurrent(host IHost, vals map[string]float64, key string) error { + val, ok := vals[key] + if !ok { + return errors.Errorf("not found %q in vals %#v", key, vals) + } + host.SetCurrent(val) + return nil +} + +type cpuUsageActive struct{} + +func newCPUUsageActive() IMetricDriver { + return &cpuUsageActive{} +} + +func (c *cpuUsageActive) GetType() monitor.MigrationAlertMetricType { + return monitor.MigrationAlertMetricTypeCPUUsageActive +} + +func (c *cpuUsageActive) GetTsdbQuery() *TsdbQuery { + return &TsdbQuery{ + Database: monitor.METRIC_DATABASE_TELE, + Measurement: "cpu", + Fields: []string{"usage_active"}, + } +} + +func (c *cpuUsageActive) GetCandidate(obj jsonutils.JSONObject, host IHost, ds *tsdb.DataSource) (ICandidate, error) { + return newCPUCandidate(obj, host.GetHostResource(), ds) +} + +func (c *cpuUsageActive) SetHostCurrent(host IHost, vals map[string]float64) error { + return setHostCurrent(host, vals, "usage_active") +} + +func (c *cpuUsageActive) GetTarget(host jsonutils.JSONObject) (ITarget, error) { + return newTargetCPUHost(host) +} + +func (ma *cpuUsageActive) GetCondition(s *monitor.AlertSetting) (ICondition, error) { + t, err := GetAlertSettingThreshold(s) + if err != nil { + return nil, errors.Wrap(err, "GetAlertSettingThreshold") + } + return newCPUCond(t), nil +} + +// cpuCondition implements ICondition +type cpuCondition struct { + value float64 +} + +func newCPUCond(val float64) ICondition { + return &cpuCondition{ + value: val, + } +} + +func (c *cpuCondition) GetThreshold() float64 { + return c.value +} + +func (c *cpuCondition) GetSourceThresholdDelta(threshold float64, host IHost) float64 { + // cpu.usage_active + log.Errorf("--host %s current: %f, threshold: %f", host.GetName(), host.GetCurrent(), threshold) + return host.GetCurrent() - threshold +} + +func (m *cpuCondition) IsFitTarget(t ITarget, c ICandidate) error { + if t.GetCurrent()+c.GetScore() < m.GetThreshold() { + return nil + } + return errors.Errorf("host:%s:current(%f) + guest:%s:score(%f) >= threshold(%f)", t.GetName(), t.GetCurrent(), c.GetName(), c.GetScore(), m.GetThreshold()) +} + +// cpuCandidate implements ICandidate +type cpuCandidate struct { + *guestResource + score float64 +} + +func newCPUCandidate(gst jsonutils.JSONObject, host *HostResource, ds *tsdb.DataSource) (ICandidate, error) { + res, err := newGuestResource(gst, host.GetName()) + if err != nil { + return nil, errors.Wrap(err, "newGuestResource") + } + + gstCPUCount, err := res.guest.Int("vcpu_count") + if err != nil { + return nil, errors.Wrap(err, "get vcpu_count") + } + + // fetch metric from influxdb + metrics, err := InfluxdbQuery(ds, "vm_id", []IResource{res}, &TsdbQuery{ + Database: monitor.METRIC_DATABASE_TELE, + Measurement: "vm_cpu", + Fields: []string{"usage_active"}, + }) + if err != nil { + return nil, errors.Wrapf(err, "InfluxdbQuery guest %q(%q)", res.GetName(), res.GetId()) + } + + metric := metrics.Get(res.GetId()) + usage := metric.Values["usage_active"] + score := usage * (float64(gstCPUCount) / float64(host.cpuCount)) + return &cpuCandidate{ + guestResource: res, + score: score, + }, nil +} + +func (c cpuCandidate) GetScore() float64 { + return c.score +} + +type cpuHost struct { + *HostResource + usageActive float64 +} + +func newCPUHost(obj jsonutils.JSONObject) (IHost, error) { + host, err := newHostResource(obj) + if err != nil { + return nil, errors.Wrap(err, "newHostResource") + } + return &cpuHost{ + HostResource: host, + }, nil +} + +func (h *cpuHost) GetCurrent() float64 { + return h.usageActive +} + +func (h *cpuHost) SetCurrent(val float64) IHost { + h.usageActive = val + return h +} + +func (h *cpuHost) Compare(oh IHost) bool { + return h.GetCurrent() < oh.GetCurrent() +} + +type targetCPUHost struct { + IHost +} + +func newTargetCPUHost(obj jsonutils.JSONObject) (ITarget, error) { + host, err := newCPUHost(obj) + if err != nil { + return nil, errors.Wrap(err, "newCPUHost") + } + ts := &targetCPUHost{ + IHost: host, + } + return ts, nil +} + +func (ts *targetCPUHost) Selected(c ICandidate) ITarget { + ts.SetCurrent(ts.GetCurrent() + c.GetScore()) + return ts +} diff --git a/pkg/monitor/controller/balancer/doc.go b/pkg/monitor/controller/balancer/doc.go new file mode 100644 index 0000000000..ff3d6a1c02 --- /dev/null +++ b/pkg/monitor/controller/balancer/doc.go @@ -0,0 +1 @@ +package balancer // import "yunion.io/x/onecloud/pkg/monitor/controller/balancer" diff --git a/pkg/monitor/controller/balancer/mem.go b/pkg/monitor/controller/balancer/mem.go new file mode 100644 index 0000000000..379a2d247d --- /dev/null +++ b/pkg/monitor/controller/balancer/mem.go @@ -0,0 +1,171 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package balancer + +import ( + "yunion.io/x/jsonutils" + "yunion.io/x/pkg/errors" + + "yunion.io/x/onecloud/pkg/apis/monitor" + "yunion.io/x/onecloud/pkg/monitor/tsdb" +) + +type memAvailable struct{} + +func newMemAvailable() IMetricDriver { + return &memAvailable{} +} + +func (m *memAvailable) GetType() monitor.MigrationAlertMetricType { + return monitor.MigrationAlertMetricTypeMemAvailable +} + +func (ma *memAvailable) GetTsdbQuery() *TsdbQuery { + return &TsdbQuery{ + Database: monitor.METRIC_DATABASE_TELE, + Measurement: "mem", + Fields: []string{"total", "free", "available"}, + } +} + +func (ma *memAvailable) GetCandidate(obj jsonutils.JSONObject, host IHost, _ *tsdb.DataSource) (ICandidate, error) { + return newMemCandidate(obj, host) +} + +func (ma *memAvailable) SetHostCurrent(host IHost, vals map[string]float64) error { + return setHostCurrent(host, vals, "free") +} + +func (ma *memAvailable) GetTarget(host jsonutils.JSONObject) (ITarget, error) { + return newTargetMemHost(host) +} + +func (ma *memAvailable) GetCondition(s *monitor.AlertSetting) (ICondition, error) { + t, err := GetAlertSettingThreshold(s) + if err != nil { + return nil, errors.Wrap(err, "GetAlertSettingThreshold") + } + return newMemoryCond(t), nil +} + +// memCondition implements ICondition +type memCondition struct { + value float64 +} + +func newMemoryCond(value float64) ICondition { + return &memCondition{ + value: value, + } +} + +func (m *memCondition) GetThreshold() float64 { + return m.value +} + +func (m *memCondition) GetSourceThresholdDelta(threshold float64, host IHost) float64 { + // mem.available + return threshold - host.GetCurrent() +} + +func (m *memCondition) IsFitTarget(t ITarget, c ICandidate) error { + if t.GetCurrent()-c.GetScore() > m.GetThreshold() { + return nil + } + return errors.Errorf("host:%s:current(%f) - guest:%s:score(%f) <= threshold(%f)", t.GetName(), t.GetCurrent(), c.GetName(), c.GetScore(), m.GetThreshold()) +} + +// memCandidate implements ICandidate +type memCandidate struct { + *guestResource + score float64 +} + +func newMemCandidate(gst jsonutils.JSONObject, host IHost) (ICandidate, error) { + res, err := newGuestResource(gst, host.GetName()) + if err != nil { + return nil, errors.Wrap(err, "newGuestResource") + } + + memSizeMB, err := gst.Int("vmem_size") + if err != nil { + return nil, errors.Wrap(err, "get vmem_size") + } + + /* unit of influxdb query is byte + > select free, available, total from mem where host_id = 'eda7c6f5-f714-4d59-8d6a-16b658712b07' limit 1; + name: mem + time free available total + ---- ---- --------- ----- + 2022-05-02T00:00:00Z 15399550976 94193070080 270276599808 + */ + + return &memCandidate{ + guestResource: res, + score: float64(memSizeMB * 1024 * 1024), + }, nil +} + +func (m *memCandidate) GetScore() float64 { + return m.score +} + +type memHost struct { + *HostResource + availableMemSize float64 +} + +func newMemHost(obj jsonutils.JSONObject) (IHost, error) { + host, err := newHostResource(obj) + if err != nil { + return nil, errors.Wrap(err, "newHostResource") + } + return &memHost{ + HostResource: host, + }, nil +} + +func (ts *memHost) GetCurrent() float64 { + return ts.availableMemSize +} + +func (ts *memHost) SetCurrent(val float64) IHost { + ts.availableMemSize = val + return ts +} + +func (ts *memHost) Compare(oh IHost) bool { + return ts.GetCurrent() > oh.GetCurrent() +} + +type targetMemHost struct { + IHost +} + +func newTargetMemHost(obj jsonutils.JSONObject) (ITarget, error) { + host, err := newMemHost(obj) + if err != nil { + return nil, errors.Wrap(err, "newMemHost") + } + ts := &targetMemHost{ + IHost: host, + } + return ts, nil +} + +func (ts *targetMemHost) Selected(c ICandidate) ITarget { + ts.SetCurrent(ts.GetCurrent() - c.GetScore()) + return ts +} diff --git a/pkg/monitor/controller/balancer/recorder.go b/pkg/monitor/controller/balancer/recorder.go new file mode 100644 index 0000000000..bf936156a8 --- /dev/null +++ b/pkg/monitor/controller/balancer/recorder.go @@ -0,0 +1,169 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package balancer + +import ( + "context" + "strings" + "time" + + "yunion.io/x/jsonutils" + "yunion.io/x/log" + "yunion.io/x/pkg/errors" + "yunion.io/x/pkg/util/wait" + + "yunion.io/x/onecloud/pkg/apis/compute" + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/mcclient" + computemod "yunion.io/x/onecloud/pkg/mcclient/modules/compute" + "yunion.io/x/onecloud/pkg/monitor/models" +) + +type EventAction string + +const ( + EventActionFindResultFail = "find_result_fail" + EventActionMigrating = "migrating" + EventActionMigrateSuccess = "migrate_success" + EventActionMigrateFail = "migrate_fail" + EventActionMigrateError = "migrate_error" +) + +type IRecorder interface { + Record(userCred mcclient.TokenCredential, alert *models.SMigrationAlert, notes interface{}, act EventAction) + RecordError(userCred mcclient.TokenCredential, alert *models.SMigrationAlert, err error, act EventAction) + + RecordMigrate(ctx context.Context, s *mcclient.ClientSession, alert *models.SMigrationAlert, note *models.MigrateNote) error + RecordMigrateError(userCred mcclient.TokenCredential, alert *models.SMigrationAlert, note *models.MigrateNote, err error) error + + StartWatchMigratingProcess(ctx context.Context, s *mcclient.ClientSession, alert *models.SMigrationAlert, note *models.MigrateNote) +} + +type sRecorder struct{} + +func NewRecorder() IRecorder { + return new(sRecorder) +} + +func (r *sRecorder) Record(userCred mcclient.TokenCredential, alert *models.SMigrationAlert, notes interface{}, act EventAction) { + db.OpsLog.LogEvent(alert, string(act), notes, userCred) +} + +func (r *sRecorder) RecordError(userCred mcclient.TokenCredential, alert *models.SMigrationAlert, err error, act EventAction) { + db.OpsLog.LogEvent(alert, string(act), err, userCred) +} + +func NewMigrateNote(pair *resultPair, err error) (*models.MigrateNote, error) { + gst := new(models.MigrateNoteGuest) + src := pair.source + if err := src.GetObject().Unmarshal(gst); err != nil { + return nil, errors.Wrap(err, "Unmarshal source") + } + gst.Host = src.GetHostName() + gst.Score = pair.source.GetScore() + + target := pair.target + host := &models.MigrateNoteTarget{ + Id: target.GetId(), + Name: target.GetName(), + Score: target.GetCurrent(), + } + + note := &models.MigrateNote{ + Guest: gst, + Target: host, + } + if err != nil { + note.Error = err.Error() + } + return note, nil +} + +func (r *sRecorder) RecordMigrate(ctx context.Context, s *mcclient.ClientSession, alert *models.SMigrationAlert, note *models.MigrateNote) error { + r.Record(s.GetToken(), alert, note, EventActionMigrating) + if err := alert.SetMigrateNote(ctx, note, false); err != nil { + return errors.Wrap(err, "SetMigrateNote") + } + // start watcher to trace migrating process + r.StartWatchMigratingProcess(ctx, s, alert, note) + return nil +} + +func (r *sRecorder) StartWatchMigratingProcess(ctx context.Context, s *mcclient.ClientSession, alert *models.SMigrationAlert, note *models.MigrateNote) { + go r.startWatchMigratingProcess(ctx, s, alert, note) +} + +func (r *sRecorder) startWatchMigratingProcess(ctx context.Context, s *mcclient.ClientSession, alert *models.SMigrationAlert, note *models.MigrateNote) { + interval := time.Second * 30 + noteStr := jsonutils.Marshal(note).String() + serverId := note.Guest.Id + serverName := note.Guest.Name + sourceHostId := note.Guest.HostId + targetHostId := note.Target.Id + if err := wait.PollImmediateInfinite(interval, func() (bool, error) { + log.Infof("start to watch migrating process %s", noteStr) + srvObj, err := computemod.Servers.Get(s, serverId, jsonutils.NewDict()) + if err != nil { + return false, errors.Wrapf(err, "Get server %s(%s) from cloud", serverName, serverId) + } + status, err := srvObj.GetString("status") + if err != nil { + return false, errors.Wrapf(err, "Get server status %s(%s)", serverName, serverId) + } + curHostId, err := srvObj.GetString("host_id") + if err != nil { + return false, errors.Wrapf(err, "Get server host_id %s(%s)", serverName, serverId) + } + if status == compute.VM_RUNNING { + if curHostId != sourceHostId { + if curHostId == targetHostId { + return true, nil + } else { + return false, errors.Errorf("Server expected in target host %s, current %s", targetHostId, curHostId) + } + } else { + return false, errors.Errorf("Server still in source host %s", sourceHostId) + } + } + if strings.HasSuffix(status, "_fail") || strings.HasSuffix(status, "_failed") { + if status == compute.VM_MIGRATE_FAILED { + // try sync status to orignal + if _, err := computemod.Servers.PerformAction(s, serverId, "syncstatus", nil); err != nil { + log.Errorf("Sync server %s(%s) migrate_failed status error: %v", serverName, serverId, err) + } + } + return false, errors.Errorf("Server fail status %s", status) + } + log.Errorf("%s: server status %q, continue watching", noteStr, status) + return false, nil + }); err != nil { + note.Error = err.Error() + r.Record(s.GetToken(), alert, note, EventActionMigrateFail) + if err := alert.SetMigrateNote(ctx, note, true); err != nil { + log.Errorf("Delete alert %s(%s) migrate note on failure: %s", alert.GetName(), alert.GetId(), noteStr) + } + } else { + r.Record(s.GetToken(), alert, note, EventActionMigrateSuccess) + if err := alert.SetMigrateNote(ctx, note, true); err != nil { + log.Errorf("Delete alert %s(%s) migrate note on success: %s", alert.GetName(), alert.GetId(), noteStr) + } + } +} + +func (r *sRecorder) RecordMigrateError(userCred mcclient.TokenCredential, alert *models.SMigrationAlert, note *models.MigrateNote, err error) error { + note.Error = err.Error() + r.Record(userCred, alert, note, EventActionMigrateError) + return nil +} diff --git a/pkg/monitor/controller/balancer/utils.go b/pkg/monitor/controller/balancer/utils.go new file mode 100644 index 0000000000..42fa5c1c5d --- /dev/null +++ b/pkg/monitor/controller/balancer/utils.go @@ -0,0 +1,109 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package balancer + +import ( + "context" + + "yunion.io/x/jsonutils" + "yunion.io/x/pkg/errors" + + api "yunion.io/x/onecloud/pkg/apis/monitor" + "yunion.io/x/onecloud/pkg/mcclient/modules/monitor" + "yunion.io/x/onecloud/pkg/monitor/tsdb" + "yunion.io/x/onecloud/pkg/monitor/tsdb/driver/influxdb" +) + +type HostMetric struct { + Id string + Values map[string]float64 +} + +type HostMetrics struct { + metrics []*HostMetric + indexes map[string]*HostMetric +} + +func NewHostMetrics(ms []*HostMetric) *HostMetrics { + h := &HostMetrics{ + metrics: ms, + indexes: make(map[string]*HostMetric), + } + for _, m := range ms { + h.indexes[m.Id] = m + } + return h +} + +func (hs HostMetrics) JSONString() string { + return jsonutils.Marshal(hs.metrics).String() +} + +func (hs HostMetrics) Get(id string) *HostMetric { + return hs.indexes[id] +} + +type TsdbQuery struct { + Database string + Measurement string + Fields []string +} + +func InfluxdbQuery( + ds *tsdb.DataSource, + idKey string, + hosts []IResource, + query *TsdbQuery) (*HostMetrics, error) { + q := monitor.NewAlertQuery(query.Database, query.Measurement).From("5m").To("now") + sels := q.Selects() + for _, field := range query.Fields { + sels.Select(field).MEAN().AS(field) + } + ids := []string{} + for _, h := range hosts { + ids = append(ids, h.GetId()) + } + q.Where().IN(idKey, ids) + q.GroupBy().TAG(idKey).FILL_NULL() + qCtx := q.ToTsdbQuery() + endpoint, err := influxdb.NewInfluxdbExecutor(nil) + if err != nil { + return nil, errors.Wrap(err, "influxdb.NewInfluxdbExecutor") + } + resp, err := endpoint.Query(context.TODO(), ds, qCtx) + if err != nil { + return nil, errors.Wrap(err, "influxdb endpoint Query") + } + ss := resp.Results[""].Series + ms := make([]*HostMetric, len(ss)) + for i, s := range ss { + m := &HostMetric{ + Id: s.Tags[idKey], + Values: make(map[string]float64), + } + for j, f := range query.Fields { + m.Values[f] = *(s.Points[0][j].(*float64)) + } + ms[i] = m + } + return NewHostMetrics(ms), nil +} + +func GetAlertSettingThreshold(s *api.AlertSetting) (float64, error) { + if len(s.Conditions) != 1 { + return 0, errors.Errorf("AlertSetting conditions %d != 1", len(s.Conditions)) + } + return s.Conditions[0].Evaluator.Params[0], nil +} diff --git a/pkg/monitor/models/alert.go b/pkg/monitor/models/alert.go index 0de2a8d0eb..7d32875ff3 100644 --- a/pkg/monitor/models/alert.go +++ b/pkg/monitor/models/alert.go @@ -131,7 +131,7 @@ type SAlert struct { // change to `Alerting` and send alert notifications. For int64 `nullable:"false" list:"user" update:"user"` - EvalData jsonutils.JSONObject `list:"user" list:"user"` + EvalData jsonutils.JSONObject `list:"user"` State string `width:"36" charset:"ascii" nullable:"false" default:"unknown" list:"user" update:"user"` NoDataState string `width:"36" charset:"ascii" nullable:"false" default:"no_data" create:"optional" list:"user" update:"user"` ExecutionErrorState string `width:"36" charset:"ascii" nullable:"false" default:"alerting" create:"optional" list:"user" update:"user"` @@ -501,6 +501,29 @@ func (alert *SAlert) GetNotifications() ([]SAlertnotification, error) { return notis, nil } +func (alert *SAlert) deleteNotifications(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) error { + notis, err := alert.GetNotifications() + if err != nil { + return err + } + for _, noti := range notis { + conf, err := noti.GetNotification() + if err != nil { + return errors.Wrap(err, "GetNotification") + } + if err := conf.CustomizeDelete(ctx, userCred, query, data); err != nil { + return errors.Wrapf(err, "notification %s(%s) CustomizeDelete", conf.GetName(), conf.GetId()) + } + if err := noti.Detach(ctx, userCred); err != nil { + return errors.Wrapf(err, "notification %s(%s) Detach ", conf.GetName(), conf.GetId()) + } + if err := conf.Delete(ctx, userCred); err != nil { + return errors.Wrapf(err, "notification %s(%s) Delete", conf.GetName(), conf.GetId()) + } + } + return nil +} + func (alert *SAlert) GetAlertResources() ([]*SAlertResource, error) { jRess := make([]SAlertResourceAlert, 0) jm := GetAlertResourceAlertManager() diff --git a/pkg/monitor/models/alertnotification.go b/pkg/monitor/models/alertnotification.go index f31780df63..deea431bd4 100644 --- a/pkg/monitor/models/alertnotification.go +++ b/pkg/monitor/models/alertnotification.go @@ -29,9 +29,10 @@ import ( ) const ( - AlertNotificationUsedByMeterAlert = "meter_alert" - AlertNotificationUsedByNodeAlert = "node_alert" - AlertNotificationUsedByCommonAlert = "common_alert" + AlertNotificationUsedByMeterAlert = "meter_alert" + AlertNotificationUsedByNodeAlert = "node_alert" + AlertNotificationUsedByCommonAlert = "common_alert" + AlertNotificationUsedByMigrationAlert = "migration_alert" ) type SAlertNotificationManager struct { diff --git a/pkg/monitor/models/balancerule.go b/pkg/monitor/models/balancerule.go new file mode 100644 index 0000000000..e48d5da859 --- /dev/null +++ b/pkg/monitor/models/balancerule.go @@ -0,0 +1,18 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package models + +type SBalanceRule struct { +} diff --git a/pkg/monitor/models/commonalert.go b/pkg/monitor/models/commonalert.go index e7b9024880..66f753dfab 100644 --- a/pkg/monitor/models/commonalert.go +++ b/pkg/monitor/models/commonalert.go @@ -66,11 +66,19 @@ var ( ) func init() { + GetCommonAlertManager() + //registry.RegisterService(CommonAlertManager) +} + +func GetCommonAlertManager() *SCommonAlertManager { + if CommonAlertManager != nil { + return CommonAlertManager + } CommonAlertManager = &SCommonAlertManager{ SAlertManager: *NewAlertManager(SCommonAlert{}, "commonalert", "commonalerts"), } CommonAlertManager.SetVirtualObject(CommonAlertManager) - //registry.RegisterService(CommonAlertManager) + return CommonAlertManager } type ISubscriptionManager interface { @@ -1089,26 +1097,7 @@ func (self *SCommonAlert) DeleteAttachAlertRecords(ctx context.Context, userCred func (alert *SCommonAlert) customizeDeleteNotis( ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) error { - notis, err := alert.GetNotifications() - if err != nil { - return err - } - for _, noti := range notis { - conf, err := noti.GetNotification() - if err != nil { - return err - } - if err := conf.CustomizeDelete(ctx, userCred, query, data); err != nil { - return err - } - if err := noti.Detach(ctx, userCred); err != nil { - return err - } - if err := conf.Delete(ctx, userCred); err != nil { - return err - } - } - return nil + return alert.deleteNotifications(ctx, userCred, query, data) } func (alert *SCommonAlert) Delete(ctx context.Context, userCred mcclient.TokenCredential) error { diff --git a/pkg/monitor/models/migration_alert.go b/pkg/monitor/models/migration_alert.go new file mode 100644 index 0000000000..a59638eba7 --- /dev/null +++ b/pkg/monitor/models/migration_alert.go @@ -0,0 +1,327 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package models + +import ( + "context" + "database/sql" + "time" + + "yunion.io/x/jsonutils" + "yunion.io/x/pkg/errors" + "yunion.io/x/sqlchemy" + + "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" +) + +var ( + migrationAlertMan *SMigrationAlertManager +) + +func init() { + migrationAlertMan = GetMigrationAlertManager() +} + +func GetMigrationAlertManager() *SMigrationAlertManager { + if migrationAlertMan != nil { + return migrationAlertMan + } + migrationAlertMan = &SMigrationAlertManager{ + SAlertManager: *NewAlertManager(SMigrationAlert{}, "migrationalert", "migrationalerts"), + } + migrationAlertMan.SetVirtualObject(migrationAlertMan) + return migrationAlertMan +} + +type SMigrationAlertManager struct { + SAlertManager +} + +type SMigrationAlert struct { + SAlert + + MetricType string `create:"admin_required" list:"admin" get:"admin"` + MigrateNotes jsonutils.JSONObject `nullable:"true" list:"admin" get:"admin" update:"admin" create:"admin_optional"` +} + +type MigrateNoteGuest struct { + Id string `json:"id"` + Name string `json:"name"` + HostId string `json:"host_id"` + Host string `json:"host"` + VCPUCount int `json:"vcpu_count"` + VMemSize int `json:"vmem_size"` + Score float64 `json:"score"` +} + +type MigrateNoteTarget struct { + Id string `json:"id"` + Name string `json:"name"` + Score float64 `json:"score"` +} + +type MigrateNote struct { + Guest *MigrateNoteGuest `json:"guest"` + Target *MigrateNoteTarget `json:"target_host"` + Error string `json:"error"` +} + +func (m *SMigrationAlertManager) ListItemFilter(ctx context.Context, q *sqlchemy.SQuery, userCred mcclient.TokenCredential, query monitor.MigrationAlertListInput) (*sqlchemy.SQuery, error) { + q, err := m.SAlertManager.ListItemFilter(ctx, q, userCred, query.AlertListInput) + if err != nil { + return nil, err + } + if len(query.MetricType) > 0 { + q = q.Equals("metric_type", query.MetricType) + } + q = q.Equals("used_by", AlertNotificationUsedByMigrationAlert) + return q, nil +} + +func (m *SMigrationAlertManager) ValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, input *monitor.MigrationAlertCreateInput) (*monitor.MigrationAlertCreateInput, error) { + if input.Period == "" { + input.Period = "5m" + } + if _, err := time.ParseDuration(input.Period); err != nil { + return nil, httperrors.NewInputParameterError("Invalid period format: %s", input.Period) + } + + if err := monitor.IsValidMigrationAlertMetricType(input.MetricType); err != nil { + return nil, httperrors.NewInputParameterError("Invalid metric_type %v", err) + } + + if input.MigrationSettings == nil { + input.MigrationSettings = &monitor.MigrationAlertSettings{ + Source: new(monitor.MigrationAlertSettingsSource), + Target: new(monitor.MigrationAlertSettingsTarget), + } + } + if err := m.ValidateMigrationSettings(input.MigrationSettings); err != nil { + return nil, errors.Wrap(err, "validate migration settings") + } + + aInput := *(input.ToAlertCreateInput()) + aInput, err := AlertManager.ValidateCreateData(ctx, userCred, nil, query, aInput) + if err != nil { + return input, errors.Wrap(err, "AlertManager.ValidateCreateData") + } + input.AlertCreateInput = aInput + + return input, nil +} + +func (m *SMigrationAlertManager) ValidateMigrationSettings(s *monitor.MigrationAlertSettings) error { + if s.Source != nil { + if err := m.ValidateMigrationSettingsSource(s.Source); err != nil { + return errors.Wrap(err, "validate source") + } + } + if s.Target != nil { + if err := m.ValidateMigrationSettingsTarget(s.Target); err != nil { + return errors.Wrap(err, "validate target") + } + } + return nil +} + +func (m *SMigrationAlertManager) GetResourceByIdOrName(rType string, id string) (jsonutils.JSONObject, error) { + ok, objs := MonitorResourceManager.GetResourceObjByResType(rType) + if !ok { + return nil, errors.Errorf("Get by %q", rType) + } + for _, obj := range objs { + name, _ := obj.GetString("name") + if name == id { + return obj, nil + } + objId, _ := obj.GetString("id") + if objId == id { + return obj, nil + } + } + return nil, errors.Errorf("Not found resource %q by %q", rType, id) +} + +func (m *SMigrationAlertManager) GetHostByIdOrName(id string) (jsonutils.JSONObject, error) { + return m.GetResourceByIdOrName(monitor.METRIC_RES_TYPE_HOST, id) +} + +func (m *SMigrationAlertManager) GetGuestByIdOrName(id string) (jsonutils.JSONObject, error) { + return m.GetResourceByIdOrName(monitor.METRIC_RES_TYPE_GUEST, id) +} + +func (m *SMigrationAlertManager) validateResource(vf func(idOrName string) (jsonutils.JSONObject, error), ids []string) ([]string, error) { + nIds := make([]string, len(ids)) + for idx, idName := range ids { + obj, err := vf(idName) + if err != nil { + return nil, errors.Wrapf(err, "find by %s", idName) + } + id, err := obj.GetString("id") + if err != nil { + return nil, errors.Wrap(err, "get id") + } + nIds[idx] = id + } + return nIds, nil +} + +func (m *SMigrationAlertManager) ValidateMigrationSettingsSource(input *monitor.MigrationAlertSettingsSource) error { + if len(input.HostIds) != 0 { + hIds, err := m.validateResource(m.GetHostByIdOrName, input.HostIds) + if err != nil { + return errors.Wrap(err, "validate host") + } + input.HostIds = hIds + } + if len(input.GuestIds) != 0 { + gIds, err := m.validateResource(m.GetGuestByIdOrName, input.GuestIds) + if err != nil { + return errors.Wrap(err, "validate guest") + } + input.GuestIds = gIds + } + return nil +} + +func (m *SMigrationAlertManager) ValidateMigrationSettingsTarget(input *monitor.MigrationAlertSettingsTarget) error { + if len(input.HostIds) != 0 { + hIds, err := m.validateResource(m.GetHostByIdOrName, input.HostIds) + if err != nil { + return errors.Wrap(err, "validate host") + } + input.HostIds = hIds + } + return nil +} + +func (alert *SMigrationAlert) CustomizeCreate(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, data jsonutils.JSONObject) error { + out := new(monitor.MigrationAlertCreateInput) + if err := data.Unmarshal(out); err != nil { + return errors.Wrap(err, "Unmarshal to MigrationAlertCreateInput") + } + fs := out.GetMetricDriver().GetQueryFields() + alert.ResType = string(fs.ResourceType) + alert.UsedBy = AlertNotificationUsedByMigrationAlert + + if err := alert.SAlert.CustomizeCreate(ctx, userCred, ownerId, query, data); err != nil { + return errors.Wrap(err, "SAlert.CustomizeCreate") + } + return alert.CreateNotification(ctx, userCred) +} + +func (m *SMigrationAlertManager) FetchAllMigrationAlerts() ([]SMigrationAlert, error) { + objs := make([]SMigrationAlert, 0) + q := m.Query() + q = q.IsTrue("enabled") + err := db.FetchModelObjects(m, q, &objs) + if err != nil && err != sql.ErrNoRows { + return nil, errors.Wrap(err, "db.FetchModelObjects") + } + return objs, nil +} + +func (m *SMigrationAlertManager) GetInMigrationAlerts() ([]*SMigrationAlert, error) { + alerts, err := m.FetchAllMigrationAlerts() + if err != nil { + return nil, errors.Wrap(err, "FetchAllMigrationAlerts") + } + objs := make([]*SMigrationAlert, 0) + for _, a := range alerts { + if ok, _, _ := a.IsInMigrationProcess(); ok { + tmp := a + objs = append(objs, &tmp) + } + } + return objs, nil +} + +func (alert *SMigrationAlert) GetMigrateNotes() (map[string]MigrateNote, error) { + if alert.MigrateNotes == nil { + return make(map[string]MigrateNote, 0), nil + } + objs := make(map[string]MigrateNote, 0) + if err := alert.MigrateNotes.Unmarshal(&objs); err != nil { + return nil, errors.Wrap(err, "Unmarshal") + } + return objs, nil +} + +func (alert *SMigrationAlert) SetMigrateNote(ctx context.Context, ns *MigrateNote, isDelete bool) error { + _, err := db.UpdateWithLock(ctx, alert, func() error { + curNotes, err := alert.GetMigrateNotes() + if err != nil { + return errors.Wrap(err, "GetMigrateNotes") + } + if isDelete { + delete(curNotes, ns.Guest.Id) + } else { + curNotes[ns.Guest.Id] = *ns + } + alert.MigrateNotes = jsonutils.Marshal(curNotes) + return nil + }) + return err +} + +func (alert *SMigrationAlert) IsInMigrationProcess() (bool, map[string]MigrateNote, error) { + notes, err := alert.GetMigrateNotes() + if err != nil { + return false, nil, errors.Wrap(err, "GetMigrateNotes") + } + if len(notes) == 0 { + return false, nil, nil + } + return true, nil, nil +} + +func (alert *SMigrationAlert) GetMigrationSettings() (*monitor.MigrationAlertSettings, error) { + if alert.CustomizeConfig == nil { + return nil, errors.Errorf("CustomizeConfig is nil") + } + out := new(monitor.MigrationAlertSettings) + if err := alert.CustomizeConfig.Unmarshal(out); err != nil { + return nil, err + } + return out, nil +} + +func (alert *SMigrationAlert) CreateNotification(ctx context.Context, userCred mcclient.TokenCredential) error { + if alert.Id == "" { + alert.Id = db.DefaultUUIDGenerator() + } + noti, err := NotificationManager.CreateAutoMigrationNotification(ctx, userCred, alert) + if err != nil { + return errors.Wrap(err, "CreateAutoMigrationNotification") + } + if _, err := alert.AttachNotification(ctx, userCred, noti, monitor.AlertNotificationStateUnknown, ""); err != nil { + return errors.Wrap(err, "alert.AttachNotification") + } + return nil +} + +func (alert *SMigrationAlert) GetMetricType() monitor.MigrationAlertMetricType { + return monitor.MigrationAlertMetricType(alert.MetricType) +} + +func (alert *SMigrationAlert) CustomizeDelete(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) error { + if err := alert.deleteNotifications(ctx, userCred, query, data); err != nil { + return errors.Wrap(err, "delete related notification") + } + return alert.SAlert.CustomizeDelete(ctx, userCred, query, data) +} diff --git a/pkg/monitor/models/nodealert.go b/pkg/monitor/models/nodealert.go index 86ecebbeca..0f04d814c4 100644 --- a/pkg/monitor/models/nodealert.go +++ b/pkg/monitor/models/nodealert.go @@ -123,7 +123,7 @@ func (man *SNodeAlertManager) ValidateCreateData( if err != nil { return nil, err } - alertInput := data.ToAlertCreateInput(name, field, measurement, "telegraf") + alertInput := data.ToAlertCreateInput(name, field, measurement, monitor.METRIC_DATABASE_TELE) alertInput, err = AlertManager.ValidateCreateData(ctx, userCred, ownerId, query, alertInput) if err != nil { return nil, err diff --git a/pkg/monitor/models/notification.go b/pkg/monitor/models/notification.go index 6d5a1d40cd..f747348db7 100644 --- a/pkg/monitor/models/notification.go +++ b/pkg/monitor/models/notification.go @@ -180,6 +180,27 @@ func (man *SNotificationManager) CreateOneCloudNotification(ctx context.Context, duration, _ := time.ParseDuration(silentPeriod) input.Frequency = duration / time.Second } + return man.createNotification(ctx, userCred, input) +} + +func (man *SNotificationManager) CreateAutoMigrationNotification(ctx context.Context, userCred mcclient.TokenCredential, alert *SMigrationAlert) (*SNotification, error) { + alertName := alert.GetName() + newName, err := db.GenerateName(ctx, man, userCred, alertName) + if err != nil { + return nil, errors.Wrapf(err, "generate name: %s", alertName) + } + settings := monitor.NotificationSettingAutoMigration{ + AlertId: alert.GetId(), + } + input := &monitor.NotificationCreateInput{ + Name: newName, + Type: monitor.AlertNotificationTypeAutoMigration, + Settings: jsonutils.Marshal(settings), + } + return man.createNotification(ctx, userCred, input) +} + +func (man *SNotificationManager) createNotification(ctx context.Context, userCred mcclient.TokenCredential, input *monitor.NotificationCreateInput) (*SNotification, error) { obj, err := db.DoCreate(man, ctx, userCred, nil, input.JSON(input), userCred) if err != nil { return nil, errors.Wrapf(err, "create notification input: %s", input.JSON(input)) diff --git a/pkg/monitor/service/handlers.go b/pkg/monitor/service/handlers.go index 3ea3329334..f098c42169 100644 --- a/pkg/monitor/service/handlers.go +++ b/pkg/monitor/service/handlers.go @@ -64,6 +64,7 @@ func InitHandlers(app *appsrv.Application) { models.AlertPanelManager, models.MonitorResourceManager, models.AlertRecordShieldManager, + models.GetMigrationAlertManager(), } { db.RegisterModelManager(manager) handler := db.NewModelHandler(manager) diff --git a/pkg/monitor/service/service.go b/pkg/monitor/service/service.go index 0600a29c6c..a2e8afb1f2 100644 --- a/pkg/monitor/service/service.go +++ b/pkg/monitor/service/service.go @@ -29,10 +29,12 @@ import ( "yunion.io/x/onecloud/pkg/cloudcommon/cronman" "yunion.io/x/onecloud/pkg/cloudcommon/db" common_options "yunion.io/x/onecloud/pkg/cloudcommon/options" + "yunion.io/x/onecloud/pkg/mcclient/auth" _ "yunion.io/x/onecloud/pkg/monitor/alerting" _ "yunion.io/x/onecloud/pkg/monitor/alerting/conditions" _ "yunion.io/x/onecloud/pkg/monitor/alerting/notifiers" _ "yunion.io/x/onecloud/pkg/monitor/alertresourcedrivers" + "yunion.io/x/onecloud/pkg/monitor/controller/balancer" "yunion.io/x/onecloud/pkg/monitor/models" _ "yunion.io/x/onecloud/pkg/monitor/notifydrivers" "yunion.io/x/onecloud/pkg/monitor/options" @@ -86,6 +88,12 @@ func StartService() { //common_app.ServeForever(app, baseOpts) InitInfluxDBSubscriptionHandlers(app, baseOpts) + // start migration recover routine + go func() { + if err := balancer.RecoverInProcessAlerts(app.GetContext(), auth.GetAdminSession(app.GetContext(), options.Options.Region, "")); err != nil { + log.Errorf("RecoverInProcessAlerts error: %v", err) + } + }() } func startServices() { diff --git a/pkg/monitor/tsdb/driver/influxdb/models.go b/pkg/monitor/tsdb/driver/influxdb/models.go index 833c6aec92..ca90ffacee 100644 --- a/pkg/monitor/tsdb/driver/influxdb/models.go +++ b/pkg/monitor/tsdb/driver/influxdb/models.go @@ -35,14 +35,14 @@ type Query struct { type Select []QueryPart type Response struct { - Results []Result - Err error + Results []Result `json:"results"` + Err error `json:"err"` } type Result struct { - Series []Row - Message []*Message - Err error + Series []Row `json:"series"` + Message []*Message `json:"message"` + Err error `json:"err"` } type Message struct { diff --git a/pkg/monitor/tsdb/driver/influxdb/response_parser.go b/pkg/monitor/tsdb/driver/influxdb/response_parser.go index f78bc0df8f..d5a68dec56 100644 --- a/pkg/monitor/tsdb/driver/influxdb/response_parser.go +++ b/pkg/monitor/tsdb/driver/influxdb/response_parser.go @@ -21,6 +21,9 @@ import ( "strconv" "strings" + "yunion.io/x/log" + "yunion.io/x/pkg/errors" + "yunion.io/x/onecloud/pkg/monitor/tsdb" ) @@ -92,6 +95,8 @@ func (rp *ResponseParser) transformRowsV2(rows []Row, queryResult *tsdb.QueryRes point, err := rp.parseTimepointV2(valuePair) if err == nil { points = append(points, point) + } else { + log.Errorf("rp.parseTimepointV2 error: %v", err) } } tags := make(map[string]string) @@ -200,7 +205,7 @@ func (rp *ResponseParser) parseTimepointV2(valuePair []interface{}) (tsdb.TimePo timestampNumber, _ := valuePair[0].(json.Number) timestamp, err := timestampNumber.Float64() if err != nil { - return tsdb.TimePoint{}, err + return tsdb.TimePoint{}, errors.Wrapf(err, "timestampNumber.Float64 of %#v", timestampNumber) } timepoint = append(timepoint, timestamp) return timepoint, nil