diff --git a/pkg/hostman/hostinfo/hostconsts/hostconsts.go b/pkg/hostman/hostinfo/hostconsts/hostconsts.go index 798c312dcd..d9ec0c1a18 100644 --- a/pkg/hostman/hostinfo/hostconsts/hostconsts.go +++ b/pkg/hostman/hostinfo/hostconsts/hostconsts.go @@ -1,9 +1,11 @@ package hostconsts const ( - TELEGRAF_TAG_KEY_BRAND = "brand" - TELEGRAF_TAG_KEY_RES_TYPE = "res_type" - TELEGRAF_TAG_KEY_HOST_TYPE = "host_type" + TELEGRAF_TAG_KEY_BRAND = "brand" + TELEGRAF_TAG_KEY_PLATFORM = "platform" + TELEGRAF_TAG_KEY_HYPERVISOR = "hypervisor" + TELEGRAF_TAG_KEY_RES_TYPE = "res_type" + TELEGRAF_TAG_KEY_HOST_TYPE = "host_type" TELEGRAF_TAG_ONECLOUD_BRAND = "OneCloud" TELEGRAF_TAG_ONECLOUD_RES_TYPE = "host" diff --git a/pkg/monitor/alerting/notifier.go b/pkg/monitor/alerting/notifier.go index 987bab66a2..d3753b5cea 100644 --- a/pkg/monitor/alerting/notifier.go +++ b/pkg/monitor/alerting/notifier.go @@ -221,5 +221,9 @@ func newAlertRecordRule(evalCtx *EvalContext) monitor.AlertRecordRule { alertRule.MeasurementDesc = evalCtx.EvalMatches[0].MeasurementDesc alertRule.FieldDesc = evalCtx.EvalMatches[0].FieldDesc } + if len(evalCtx.AlertOkEvalMatches) != 0 { + alertRule.MeasurementDesc = evalCtx.AlertOkEvalMatches[0].MeasurementDesc + alertRule.FieldDesc = evalCtx.AlertOkEvalMatches[0].FieldDesc + } return alertRule } diff --git a/pkg/monitor/alertresourcedrivers/cloudaccount.go b/pkg/monitor/alertresourcedrivers/cloudaccount.go index 233824af76..3b83e23fda 100644 --- a/pkg/monitor/alertresourcedrivers/cloudaccount.go +++ b/pkg/monitor/alertresourcedrivers/cloudaccount.go @@ -36,10 +36,10 @@ func (drvF cloudaccountDriverF) GetType() monitor.AlertResourceType { func (drvF cloudaccountDriverF) IsEvalMatched(input monitor.EvalMatch) bool { tags := input.Tags - _, hasId := tags[CLOUDACCOUNT_TAG_ID_KEY] - if !hasId { - return false - } + //_, hasId := tags[CLOUDACCOUNT_TAG_ID_KEY] + //if !hasId { + // return false + //} _, hasName := tags[CLOUDACCOUNT_TAG_NAME_KEY] if !hasName { return false diff --git a/pkg/monitor/models/alertrecord.go b/pkg/monitor/models/alertrecord.go index ca7ab50488..fa8cabb323 100644 --- a/pkg/monitor/models/alertrecord.go +++ b/pkg/monitor/models/alertrecord.go @@ -171,6 +171,11 @@ func (record *SAlertRecord) PostCreate(ctx context.Context, userCred mcclient.To log.Errorf("Reconcile from alert record error: %v", err) return } + err := GetAlertResourceManager().NotifyAlertResourceCount(ctx) + if err != nil { + log.Errorf("NotifyAlertResourceCount error: %v", err) + return + } } func (record *SAlertRecord) GetState() monitor.AlertStateType { diff --git a/pkg/monitor/models/alertresource.go b/pkg/monitor/models/alertresource.go index ab7896cfcb..48e5ef110c 100644 --- a/pkg/monitor/models/alertresource.go +++ b/pkg/monitor/models/alertresource.go @@ -17,6 +17,7 @@ package models import ( "context" "fmt" + "sync" "yunion.io/x/jsonutils" "yunion.io/x/log" @@ -27,11 +28,15 @@ import ( "yunion.io/x/onecloud/pkg/apis/monitor" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/mcclient/auth" + mc_modules "yunion.io/x/onecloud/pkg/mcclient/modules" + npk "yunion.io/x/onecloud/pkg/mcclient/modules/notify" "yunion.io/x/onecloud/pkg/util/stringutils2" ) var ( alertResourceManager *SAlertResourceManager + adminUsers *sync.Map ) func init() { @@ -373,3 +378,110 @@ func (res *SAlertResource) CustomizeDelete( } return nil } + +func (manager *SAlertResourceManager) NotifyAlertResourceCount(ctx context.Context) error { + log.Errorln("exec NotifyAlertResourceCount func") + cn, err := manager.getResourceCount() + if err != nil { + return err + } + alertResourceCount := resourceCount{ + AlertResourceCount: cn, + } + if adminUsers == nil { + manager.GetAdminRoleUsers(ctx, nil, true) + } + adminUsersTmp := *adminUsers + ids := make([]string, 0) + adminUsersTmp.Range(func(key, value interface{}) bool { + ids = append(ids, key.(string)) + return true + }) + if len(ids) == 0 { + return fmt.Errorf("no find users in receivers has admin role") + } + //if len(ids) != 0 { + // notifyclient.RawNotifyWithCtx(ctx, ids, false, npk.NotifyByWebConsole, npk.NotifyPriorityCritical, + // "alertResourceCount", jsonutils.Marshal(&alertResourceCount)) + // return nil + //} else { + // return fmt.Errorf("no find users in receivers has admin role") + //} + manager.sendWebsocketInfo(ids, alertResourceCount) + return nil +} + +type resourceCount struct { + AlertResourceCount int `json:"alert_resource_count"` +} + +func (manager *SAlertResourceManager) getResourceCount() (int, error) { + query := manager.Query("id") + cn, err := query.CountWithError() + if err != nil { + return cn, errors.Wrap(err, "SAlertResourceManager get resource count error") + } + + return cn, nil +} + +func (manager *SAlertResourceManager) GetAdminRoleUsers(ctx context.Context, userCred mcclient.TokenCredential, + isStart bool) { + if adminUsers == nil { + adminUsers = new(sync.Map) + } + offset := 0 + query := jsonutils.NewDict() + session := auth.GetAdminSession(ctx, "", "") + rid, err := mc_modules.RolesV3.GetId(session, "admin", jsonutils.NewDict()) + if err != nil { + errors.Errorf("get role id error:%v", err) + return + } + query.Add(jsonutils.NewString(rid), "role", "id") + for { + query.Set("offset", jsonutils.NewInt(int64(offset))) + result, err := mc_modules.RoleAssignments.List(session, query) + if err != nil { + errors.Errorf("get admin role list error:%v", err) + return + } + for _, roleAssign := range result.Data { + userId, err := roleAssign.GetString("user", "id") + if err != nil { + log.Errorf("roleAssign:%v", roleAssign) + return + } + //_, err = mc_modules.NotifyReceiver.GetById(session, userId, jsonutils.NewDict()) + //if err != nil { + // log.Errorf("Recipients GetById err:%v", err) + // continue + //} + adminUsers.Store(userId, roleAssign) + } + offset = result.Offset + len(result.Data) + if offset >= result.Total { + break + } + } +} + +func (manager *SAlertResourceManager) sendWebsocketInfo(uids []string, alertResourceCount resourceCount) { + session := auth.GetAdminSession(context.Background(), "", "") + params := jsonutils.NewDict() + params.Set("obj_type", jsonutils.NewString("monitor")) + params.Set("obj_id", jsonutils.NewString("")) + params.Set("obj_name", jsonutils.NewString("")) + params.Set("success", jsonutils.JSONTrue) + params.Set("action", jsonutils.NewString("alertResourceCount")) + params.Set("notes", jsonutils.NewString(fmt.Sprintf("priority=%s; content=%s", string(npk.NotifyPriorityCritical), + jsonutils.Marshal(&alertResourceCount).String()))) + for _, uid := range uids { + params.Set("user_id", jsonutils.NewString(uid)) + params.Set("user", jsonutils.NewString(uid)) + _, err := mc_modules.Websockets.Create(session, params) + if err != nil { + log.Errorf("websocket send info err:%v", err) + } + } +} diff --git a/pkg/monitor/models/datasource.go b/pkg/monitor/models/datasource.go index f8fdedaa2e..06b83e6d24 100644 --- a/pkg/monitor/models/datasource.go +++ b/pkg/monitor/models/datasource.go @@ -37,6 +37,7 @@ import ( identityapi "yunion.io/x/onecloud/pkg/apis/identity" "yunion.io/x/onecloud/pkg/apis/monitor" "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/hostman/hostinfo/hostconsts" "yunion.io/x/onecloud/pkg/httperrors" "yunion.io/x/onecloud/pkg/mcclient/auth" merrors "yunion.io/x/onecloud/pkg/monitor/errors" @@ -608,10 +609,34 @@ func (self *SDataSourceManager) GetMetricMeasurement(query jsonutils.JSONObject, if err != nil { return jsonutils.JSONNull, errors.Wrap(err, "getTagValue error") } + self.filterRtnTags(output) return jsonutils.Marshal(output), nil } +func (self *SDataSourceManager) filterRtnTags(output *monitor.InfluxMeasurement) { + for _, tag := range []string{hostconsts.TELEGRAF_TAG_KEY_BRAND, hostconsts.TELEGRAF_TAG_KEY_PLATFORM, + hostconsts.TELEGRAF_TAG_KEY_HYPERVISOR} { + if val, ok := output.TagValue[tag]; ok { + output.TagValue[hostconsts.TELEGRAF_TAG_KEY_BRAND] = val + break + } + } + for _, tag := range []string{"source", "status", hostconsts.TELEGRAF_TAG_KEY_HOST_TYPE, + hostconsts.TELEGRAF_TAG_KEY_RES_TYPE, "is_vm", "os_type", hostconsts.TELEGRAF_TAG_KEY_PLATFORM, + hostconsts.TELEGRAF_TAG_KEY_HYPERVISOR, "domain_name", "region"} { + if _, ok := output.TagValue[tag]; ok { + delete(output.TagValue, tag) + } + } + + repTag := make([]string, 0) + for tag, _ := range output.TagValue { + repTag = append(repTag, tag) + } + output.TagKey = repTag +} + func (self *SDataSourceManager) filterTagValue(measurement monitor.InfluxMeasurement, timeF timeFilter, db *influxdb.SInfluxdb, tagValChan *influxdbTagValueChan, tagFilter string) error { ctx, _ := context.WithTimeout(context.Background(), time.Second*5) diff --git a/pkg/monitor/models/unifiedmonitor.go b/pkg/monitor/models/unifiedmonitor.go index 6843311854..f524909876 100644 --- a/pkg/monitor/models/unifiedmonitor.go +++ b/pkg/monitor/models/unifiedmonitor.go @@ -11,6 +11,7 @@ import ( "yunion.io/x/onecloud/pkg/apis/monitor" "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/hostman/hostinfo/hostconsts" "yunion.io/x/onecloud/pkg/httperrors" "yunion.io/x/onecloud/pkg/mcclient" merrors "yunion.io/x/onecloud/pkg/monitor/errors" @@ -392,6 +393,12 @@ func fillSerieTags(series *tsdb.TimeSeriesSlice) { break } } + for _, tag := range []string{"source", "status", hostconsts.TELEGRAF_TAG_KEY_HOST_TYPE, + hostconsts.TELEGRAF_TAG_KEY_RES_TYPE, "cpu", "is_vm", "os_type", "domain_name", "region"} { + if _, ok := serie.Tags[tag]; ok { + delete(serie.Tags, tag) + } + } (*series)[i] = serie } } diff --git a/pkg/monitor/options/options.go b/pkg/monitor/options/options.go index 32001cb2ec..e0eabe20d1 100644 --- a/pkg/monitor/options/options.go +++ b/pkg/monitor/options/options.go @@ -22,12 +22,13 @@ type AlerterOptions struct { common_options.CommonOptions common_options.DBOptions - DataProxyTimeout int `help:"query data source proxy timeout" default:"30"` - AlertingMinIntervalSeconds int64 `help:"alerting min schedule frequency" default:"10"` - AlertingMaxAttempts int `help:"alerting engine max attempt" default:"3"` - AlertingEvaluationTimeoutSeconds int64 `help:"alerting evaluation timeout" default:"5"` - AlertingNotificationTimeoutSeconds int64 `help:"alerting notification timeout" default:"30"` - InitScopeSuggestConfigIntervalSeconds int `help:"internal to init scope suggest configs" default:"900"` + DataProxyTimeout int `help:"query data source proxy timeout" default:"30"` + AlertingMinIntervalSeconds int64 `help:"alerting min schedule frequency" default:"10"` + AlertingMaxAttempts int `help:"alerting engine max attempt" default:"3"` + AlertingEvaluationTimeoutSeconds int64 `help:"alerting evaluation timeout" default:"5"` + AlertingNotificationTimeoutSeconds int64 `help:"alerting notification timeout" default:"30"` + InitScopeSuggestConfigIntervalSeconds int `help:"internal to init scope suggest configs" default:"900"` + InitAlertResourceAdminRoleUsersIntervalSeconds int `help:"internal to init alert resource admin role users " default:"3600"` } var ( diff --git a/pkg/monitor/service/service.go b/pkg/monitor/service/service.go index 75d32f823b..23917f7c32 100644 --- a/pkg/monitor/service/service.go +++ b/pkg/monitor/service/service.go @@ -66,6 +66,7 @@ func StartService() { cron := cronman.InitCronJobManager(true, opts.CronJobWorkerCount) suggestsysdrivers.InitSuggestSysRuleCronjob() cron.AddJobAtIntervalsWithStartRun("InitScopeSuggestConfigs", time.Duration(opts.InitScopeSuggestConfigIntervalSeconds)*time.Second, models.SuggestSysRuleConfigManager.InitScopeConfigs, true) + cron.AddJobAtIntervalsWithStartRun("InitAlertResourceAdminRoleUsers", time.Duration(opts.InitAlertResourceAdminRoleUsersIntervalSeconds)*time.Second, models.GetAlertResourceManager().GetAdminRoleUsers, true) cron.Start() defer cron.Stop()