Merge pull request #14493 from zexi/monitor-lb-migrate

feat(monitor): VMs of host loadbalance
This commit is contained in:
Zexi Li
2022-06-23 17:00:43 +08:00
committed by GitHub
43 changed files with 2461 additions and 71 deletions
+9
View File
@@ -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)
})
}
+10 -4
View File
@@ -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
})
}
+1 -2
View File
@@ -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))
+1 -2
View File
@@ -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))
+1 -2
View File
@@ -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))
+1 -2
View File
@@ -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))
+1 -1
View File
@@ -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))
+7
View File
@@ -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
}
+1 -2
View File
@@ -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{})
+1 -2
View File
@@ -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))
@@ -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))
@@ -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))
}
+1 -2
View File
@@ -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))
}
@@ -21,6 +21,6 @@ import (
)
func init() {
cmd := shell.NewJointCmd(modules.MonitorResourceAlertManager)
cmd := shell.NewJointCmd(modules.MonitorResourceAlertManager).SetPrefix("monitor")
cmd.List(new(options.MonitorResourceAlertListOptions))
}
+264
View File
@@ -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"`
}
+8 -3
View File
@@ -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,
+9 -4
View File
@@ -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"`
}
+4
View File
@@ -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)
+8 -3
View File
@@ -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
}
+13
View File
@@ -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,
@@ -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,
}
}
@@ -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
}
+5 -6
View File
@@ -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
@@ -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
}
-1
View File
@@ -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)
+671
View File
@@ -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 := &note
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
}
@@ -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)
}
})
}
}
+187
View File
@@ -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
}
+1
View File
@@ -0,0 +1 @@
package balancer // import "yunion.io/x/onecloud/pkg/monitor/controller/balancer"
+171
View File
@@ -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
}
+169
View File
@@ -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
}
+109
View File
@@ -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
}
+24 -1
View File
@@ -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()
+4 -3
View File
@@ -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 {
+18
View File
@@ -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 {
}
+10 -21
View File
@@ -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 {
+327
View File
@@ -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)
}
+1 -1
View File
@@ -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
+21
View File
@@ -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))
+1
View File
@@ -64,6 +64,7 @@ func InitHandlers(app *appsrv.Application) {
models.AlertPanelManager,
models.MonitorResourceManager,
models.AlertRecordShieldManager,
models.GetMigrationAlertManager(),
} {
db.RegisterModelManager(manager)
handler := db.NewModelHandler(manager)
+8
View File
@@ -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() {
+5 -5
View File
@@ -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 {
@@ -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