From 792bfe34649c10cf455f6f1c2e2ff40a1b2cc604 Mon Sep 17 00:00:00 2001 From: Zexi Li Date: Wed, 1 Mar 2023 00:03:29 +0800 Subject: [PATCH] fix(monitor): query resources according by listing region resources (#16074) --- pkg/apis/monitor/unifiedmonitor_const.go | 15 +- .../alerting/conditions/metricquery.go | 22 ++- .../alerting/conditions/nodataquery.go | 3 +- pkg/monitor/alerting/conditions/query.go | 78 ++++---- pkg/monitor/metricquery/interfaces.go | 3 +- pkg/monitor/models/datasource.go | 180 +++++++++++++----- pkg/monitor/models/unifiedmonitor.go | 40 ++-- 7 files changed, 225 insertions(+), 116 deletions(-) diff --git a/pkg/apis/monitor/unifiedmonitor_const.go b/pkg/apis/monitor/unifiedmonitor_const.go index f529c269f4..2903eed7f2 100644 --- a/pkg/apis/monitor/unifiedmonitor_const.go +++ b/pkg/apis/monitor/unifiedmonitor_const.go @@ -77,13 +77,14 @@ type MetricInputQuery struct { Slimit string `json:"slimit"` Soffset string `json:"soffset"` //default group by - Unit bool `json:"unit"` - Interval string `json:"interval"` - DomainId string `json:"domain_id"` - ProjectId string `json:"project_id"` - MetricQuery []*AlertQuery `json:"metric_query"` - Signature string `json:"signature"` - ShowMeta bool `json:"show_meta"` + Unit bool `json:"unit"` + Interval string `json:"interval"` + DomainId string `json:"domain_id"` + ProjectId string `json:"project_id"` + MetricQuery []*AlertQuery `json:"metric_query"` + Signature string `json:"signature"` + ShowMeta bool `json:"show_meta"` + ForceCheckSeries bool `json:"force_check_series"` } type SimpleQueryInput struct { diff --git a/pkg/monitor/alerting/conditions/metricquery.go b/pkg/monitor/alerting/conditions/metricquery.go index e550aa55ca..2c86c05d1c 100644 --- a/pkg/monitor/alerting/conditions/metricquery.go +++ b/pkg/monitor/alerting/conditions/metricquery.go @@ -25,6 +25,8 @@ import ( "yunion.io/x/pkg/utils" "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" mq "yunion.io/x/onecloud/pkg/monitor/metricquery" "yunion.io/x/onecloud/pkg/monitor/models" @@ -69,10 +71,12 @@ func NewMetricQueryCondition(models []*monitor.AlertCondition) (*MetricQueryCond return cond, nil } -func (query *MetricQueryCondition) ExecuteQuery() (*mq.Metrics, error) { - timeRange := tsdb.NewTimeRange(query.QueryCons[0].Query.From, query.QueryCons[0].Query.To) +func (query *MetricQueryCondition) ExecuteQuery(userCred mcclient.TokenCredential, forceCheckSeries bool) (*mq.Metrics, error) { + firstCond := query.QueryCons[0] + timeRange := tsdb.NewTimeRange(firstCond.Query.From, firstCond.Query.To) + ctx := gocontext.Background() evalContext := alerting.EvalContext{ - Ctx: gocontext.Background(), + Ctx: ctx, IsDebug: true, IsTestRun: false, } @@ -84,17 +88,23 @@ func (query *MetricQueryCondition) ExecuteQuery() (*mq.Metrics, error) { Series: make(tsdb.TimeSeriesSlice, 0), Metas: queryResult.metas, } - if query.noCheckSeries() { + if query.noCheckSeries() && !forceCheckSeries { metrics.Series = queryResult.series return &metrics, nil } + s := auth.GetSession(ctx, userCred, "") + ress, err := firstCond.getOnecloudResources(s, false) + if err != nil { + return nil, errors.Wrap(err, "get resources from region") + } + for _, serie := range queryResult.series { - isLatestOfSerie, resource := query.QueryCons[0].serieIsLatestResource(nil, serie) + isLatestOfSerie, resource := firstCond.serieIsLatestResource(ress, serie) if !isLatestOfSerie { continue } - query.QueryCons[0].FillSerieByResourceField(resource, serie) + firstCond.FillSerieByResourceField(resource, serie) metrics.Series = append(metrics.Series, serie) } return &metrics, nil diff --git a/pkg/monitor/alerting/conditions/nodataquery.go b/pkg/monitor/alerting/conditions/nodataquery.go index ae12c8aae9..cc6db4d0f1 100644 --- a/pkg/monitor/alerting/conditions/nodataquery.go +++ b/pkg/monitor/alerting/conditions/nodataquery.go @@ -22,6 +22,7 @@ import ( "yunion.io/x/onecloud/pkg/apis/monitor" "yunion.io/x/onecloud/pkg/hostman/hostinfo/hostconsts" + "yunion.io/x/onecloud/pkg/mcclient/auth" "yunion.io/x/onecloud/pkg/monitor/alerting" "yunion.io/x/onecloud/pkg/monitor/models" "yunion.io/x/onecloud/pkg/monitor/tsdb" @@ -87,7 +88,7 @@ serLoop: } } } - allResources, err := c.GetQueryResources() + allResources, err := c.GetQueryResources(auth.GetAdminSession(context.Ctx, ""), true) if err != nil { return nil, errors.Wrap(err, "GetQueryResources err") } diff --git a/pkg/monitor/alerting/conditions/query.go b/pkg/monitor/alerting/conditions/query.go index ae95d23697..bdd3aa1b89 100644 --- a/pkg/monitor/alerting/conditions/query.go +++ b/pkg/monitor/alerting/conditions/query.go @@ -26,12 +26,13 @@ import ( "yunion.io/x/onecloud/pkg/apis/monitor" "yunion.io/x/onecloud/pkg/hostman/hostinfo/hostconsts" - "yunion.io/x/onecloud/pkg/mcclient/auth" + "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/onecloud/pkg/mcclient/modulebase" mc_mds "yunion.io/x/onecloud/pkg/mcclient/modules/compute" "yunion.io/x/onecloud/pkg/mcclient/modules/identity" "yunion.io/x/onecloud/pkg/monitor/alerting" "yunion.io/x/onecloud/pkg/monitor/models" + "yunion.io/x/onecloud/pkg/monitor/options" "yunion.io/x/onecloud/pkg/monitor/tsdb" "yunion.io/x/onecloud/pkg/monitor/validators" ) @@ -244,14 +245,20 @@ func (c *QueryCondition) serieIsLatestResource(resources []jsonutils.JSONObject, if len(tagId) == 0 { tagId = "host_id" } - seriId := series.Tags[tagId] - //for _, resource := range resources { - // id, _ := resource.GetString("id") - // if seriId == id { - // return true, resource - // } - //} - return models.MonitorResourceManager.GetResourceObj(seriId) + resId := series.Tags[tagId] + if len(resources) != 0 { + for _, resource := range resources { + id, _ := resource.GetString("id") + if resId == id { + // return true, resource + return models.MonitorResourceManager.GetResourceObj(resId) + } else { + continue + } + } + return false, nil + } + return models.MonitorResourceManager.GetResourceObj(resId) } func (c *QueryCondition) FillSerieByResourceField(resource jsonutils.JSONObject, @@ -541,22 +548,22 @@ func (c *QueryCondition) setResType() { } } -func (c *QueryCondition) GetQueryResources() ([]jsonutils.JSONObject, error) { - allHosts, err := c.getOnecloudResources() +func (c *QueryCondition) GetQueryResources(s *mcclient.ClientSession, showDetails bool) ([]jsonutils.JSONObject, error) { + allRes, err := c.getOnecloudResources(s, showDetails) if err != nil { - return nil, errors.Wrap(err, "getOnecloudHosts error") + return nil, errors.Wrap(err, "getOnecloudResources error") } - allHosts = c.filterAllResources(allHosts) - return allHosts, nil + allRes = c.filterAllResources(allRes) + return allRes, nil } -func (c *QueryCondition) getOnecloudResources() ([]jsonutils.JSONObject, error) { +func (c *QueryCondition) getOnecloudResources(s *mcclient.ClientSession, showDetails bool) ([]jsonutils.JSONObject, error) { var err error allResources := make([]jsonutils.JSONObject, 0) query := jsonutils.NewDict() query.Add(jsonutils.NewStringArray([]string{"running", "ready"}), "status") - query.Add(jsonutils.NewString("true"), "admin") + // query.Add(jsonutils.NewString("true"), "admin") //if len(c.Query.Model.Tags) != 0 { // query, err = c.convertTagsQuery(evalContext, query) // if err != nil { @@ -568,33 +575,33 @@ func (c *QueryCondition) getOnecloudResources() ([]jsonutils.JSONObject, error) query := jsonutils.NewDict() query.Set("host_type", jsonutils.NewString(hostconsts.TELEGRAF_TAG_KEY_HYPERVISOR)) query.Set("enabled", jsonutils.NewInt(1)) - allResources, err = ListAllResources(&mc_mds.Hosts, query) + allResources, err = ListAllResources(s, &mc_mds.Hosts, query, showDetails) case monitor.METRIC_RES_TYPE_GUEST: - allResources, err = ListAllResources(&mc_mds.Servers, query) + allResources, err = ListAllResources(s, &mc_mds.Servers, query, showDetails) case monitor.METRIC_RES_TYPE_AGENT: - allResources, err = ListAllResources(&mc_mds.Servers, query) + allResources, err = ListAllResources(s, &mc_mds.Servers, query, showDetails) case monitor.METRIC_RES_TYPE_RDS: - allResources, err = ListAllResources(&mc_mds.DBInstance, query) + allResources, err = ListAllResources(s, &mc_mds.DBInstance, query, showDetails) case monitor.METRIC_RES_TYPE_REDIS: - allResources, err = ListAllResources(&mc_mds.ElasticCache, query) + allResources, err = ListAllResources(s, &mc_mds.ElasticCache, query, showDetails) case monitor.METRIC_RES_TYPE_OSS: - allResources, err = ListAllResources(&mc_mds.Buckets, query) + allResources, err = ListAllResources(s, &mc_mds.Buckets, query, showDetails) case monitor.METRIC_RES_TYPE_CLOUDACCOUNT: query.Remove("status") query.Add(jsonutils.NewBool(true), "enabled") - allResources, err = ListAllResources(&mc_mds.Cloudaccounts, query) + allResources, err = ListAllResources(s, &mc_mds.Cloudaccounts, query, showDetails) case monitor.METRIC_RES_TYPE_TENANT: - allResources, err = ListAllResources(&identity.Projects, query) + allResources, err = ListAllResources(s, &identity.Projects, query, showDetails) case monitor.METRIC_RES_TYPE_DOMAIN: - allResources, err = ListAllResources(&identity.Domains, query) + allResources, err = ListAllResources(s, &identity.Domains, query, showDetails) case monitor.METRIC_RES_TYPE_STORAGE: query.Remove("status") - allResources, err = ListAllResources(&mc_mds.Storages, query) + allResources, err = ListAllResources(s, &mc_mds.Storages, query, showDetails) default: query := jsonutils.NewDict() query.Set("brand", jsonutils.NewString(hostconsts.TELEGRAF_TAG_ONECLOUD_BRAND)) query.Set("host-type", jsonutils.NewString(hostconsts.TELEGRAF_TAG_KEY_HYPERVISOR)) - allResources, err = ListAllResources(&mc_mds.Hosts, query) + allResources, err = ListAllResources(s, &mc_mds.Hosts, query, showDetails) } if err != nil { @@ -603,25 +610,24 @@ func (c *QueryCondition) getOnecloudResources() ([]jsonutils.JSONObject, error) return allResources, nil } -func ListAllResources(manager modulebase.Manager, params *jsonutils.JSONDict) ([]jsonutils.JSONObject, error) { +func ListAllResources(s *mcclient.ClientSession, manager modulebase.Manager, params *jsonutils.JSONDict, showDetails bool) ([]jsonutils.JSONObject, error) { if params == nil { params = jsonutils.NewDict() } - params.Add(jsonutils.NewString("system"), "scope") - params.Add(jsonutils.NewInt(20), "limit") - params.Add(jsonutils.NewBool(true), "details") + if s.GetToken().HasSystemAdminPrivilege() { + params.Add(jsonutils.NewString("system"), "scope") + } + params.Add(jsonutils.NewInt(int64(options.Options.APIListBatchSize)), "limit") + params.Add(jsonutils.NewBool(showDetails), "details") var count int - session := auth.GetAdminSession(context.Background(), "") objs := make([]jsonutils.JSONObject, 0) for { params.Set("offset", jsonutils.NewInt(int64(count))) - result, err := manager.List(session, params) + result, err := manager.List(s, params) if err != nil { return nil, errors.Wrapf(err, "list %s resources with params %s", manager.KeyString(), params.String()) } - for _, data := range result.Data { - objs = append(objs, data) - } + objs = append(objs, result.Data...) total := result.Total count = count + len(result.Data) if count >= total { diff --git a/pkg/monitor/metricquery/interfaces.go b/pkg/monitor/metricquery/interfaces.go index 263bb9d256..c649f31cd9 100644 --- a/pkg/monitor/metricquery/interfaces.go +++ b/pkg/monitor/metricquery/interfaces.go @@ -16,6 +16,7 @@ package metricquery import ( "yunion.io/x/onecloud/pkg/apis/monitor" + "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/onecloud/pkg/monitor/tsdb" ) @@ -26,7 +27,7 @@ type Metrics struct { } type MetricQuery interface { - ExecuteQuery() (*Metrics, error) + ExecuteQuery(userCred mcclient.TokenCredential, forceCheckSeries bool) (*Metrics, error) } type QueryFactory func(model []*monitor.AlertCondition) (MetricQuery, error) diff --git a/pkg/monitor/models/datasource.go b/pkg/monitor/models/datasource.go index 65da5dc65d..b0314c77cd 100644 --- a/pkg/monitor/models/datasource.go +++ b/pkg/monitor/models/datasource.go @@ -38,6 +38,7 @@ import ( "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" "yunion.io/x/onecloud/pkg/mcclient/auth" merrors "yunion.io/x/onecloud/pkg/monitor/errors" "yunion.io/x/onecloud/pkg/monitor/options" @@ -497,7 +498,7 @@ func (self *SDataSourceManager) getFilterMeasurement(queryChan *influxdbQueryCha rtnMeasurement := new(monitor.InfluxMeasurement) var buffer bytes.Buffer buffer.WriteString(fmt.Sprintf(`SELECT last(*) FROM %s WHERE %s`, measurement.Measurement, - self.renderTimeFilter(from, to))) + renderTimeFilter(from, to))) if len(tagFilter) != 0 { buffer.WriteString(" AND ") buffer.WriteString(fmt.Sprintf(" %s", tagFilter)) @@ -548,7 +549,7 @@ func (self *SDataSourceManager) getFilterMeasurement(queryChan *influxdbQueryCha return nil } -func (self *SDataSourceManager) renderTimeFilter(from, to string) string { +func renderTimeFilter(from, to string) string { if strings.Contains(from, "now-") { from = "now() - " + strings.Replace(from, "now-", "", 1) } else { @@ -564,7 +565,7 @@ func (self *SDataSourceManager) renderTimeFilter(from, to string) string { } -func (self *SDataSourceManager) GetMetricMeasurement(query jsonutils.JSONObject, tagFilter string) (jsonutils.JSONObject, error) { +func (self *SDataSourceManager) GetMetricMeasurement(userCred mcclient.TokenCredential, query jsonutils.JSONObject, tagFilter string) (jsonutils.JSONObject, error) { database, _ := query.GetString("database") if database == "" { return jsonutils.JSONNull, merrors.NewArgIsEmptyErr("database") @@ -588,7 +589,7 @@ func (self *SDataSourceManager) GetMetricMeasurement(query jsonutils.JSONObject, timeF, err := self.getFromAndToFromParam(query) if err != nil { - return nil, err + return nil, errors.Wrap(err, "getFromAndToFromParam") } db := influxdb.NewInfluxdb(dataSource.Url) @@ -598,44 +599,48 @@ func (self *SDataSourceManager) GetMetricMeasurement(query jsonutils.JSONObject, output.Measurement = measurement output.Database = database output.TagValue = make(map[string][]string, 0) - for _, val := range monitor.METRIC_ATTRI { - err = getAttributesOnMeasurement(database, val, output, db) - if err != nil { - return jsonutils.JSONNull, errors.Wrap(err, "getAttributesOnMeasurement error") - } - } + // for _, val := range monitor.METRIC_ATTRI { + // if err := getAttributesOnMeasurement(database, val, output, db); err != nil { + // return jsonutils.JSONNull, errors.Wrap(err, "getAttributesOnMeasurement error") + // } + // } output.FieldKey = []string{field} //err = getTagValue(database, output, db) - tagValChan := influxdbTagValueChan{ - rtnChan: make(chan map[string][]string, len(output.FieldKey)), - count: len(output.FieldKey), - //count: 1, - } - ctx, cancel := context.WithTimeout(context.Background(), time.Second*30) - tagValGroup, _ := errgroup.WithContext(ctx) - defer cancel() - tagValGroup.Go(func() error { - return self.filterTagValue(*output, timeF, db, &tagValChan, tagFilter) - }) - tagValGroup.Go(func() error { - for i := 0; i < tagValChan.count; i++ { - select { - case tagVal := <-tagValChan.rtnChan: - if len(tagVal) != 0 { - tagValUnion(output, tagVal) - } - case <-ctx.Done(): - return fmt.Errorf("filter Union TagValue time out") - } - } - return nil - }) - err = tagValGroup.Wait() - if err != nil { - return jsonutils.JSONNull, errors.Wrap(err, "getTagValue error") + // tagValChan := influxdbTagValueChan{ + // rtnChan: make(chan map[string][]string, len(output.FieldKey)), + // count: len(output.FieldKey), + // //count: 1, + // } + + //** ctx, cancel := context.WithTimeout(context.Background(), time.Second*30) + // tagValGroup, _ := errgroup.WithContext(ctx) + // defer cancel() + // tagValGroup.Go(func() error { + // return self.filterTagValue(*output, timeF, db, &tagValChan, tagFilter) + // }) + // tagValGroup.Go(func() error { + // for i := 0; i < tagValChan.count; i++ { + // select { + // case tagVal := <-tagValChan.rtnChan: + // if len(tagVal) != 0 { + // tagValUnion(output, tagVal) + // } + // case <-ctx.Done(): + // return fmt.Errorf("filter Union TagValue time out") + // } + // } + // return nil + // }) + // err = tagValGroup.Wait() + // if err != nil { + // return jsonutils.JSONNull, errors.Wrap(err, "getTagValue error") + //** } + if err := getTagValues(userCred, output, timeF, dataSource.GetId(), tagFilter); err != nil { + return jsonutils.JSONNull, errors.Wrap(err, "getTagValues error") } + self.filterRtnTags(output) return jsonutils.Marshal(output), nil @@ -666,7 +671,7 @@ func (self *SDataSourceManager) filterRtnTags(output *monitor.InfluxMeasurement) 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) + ctx, _ := context.WithTimeout(context.Background(), time.Second*15) tagValGroup2, _ := errgroup.WithContext(ctx) tagValChan2 := influxdbTagValueChan{ rtnChan: make(chan map[string][]string, len(measurement.TagKey)), @@ -795,9 +800,10 @@ func (self *SDataSourceManager) DropSubscription(subscription InfluxdbSubscripti } func getAttributesOnMeasurement(database, tp string, output *monitor.InfluxMeasurement, db *influxdb.SInfluxdb) error { - dbRtn, err := db.Query(fmt.Sprintf("SHOW %s KEYS ON %s FROM %s", tp, database, output.Measurement)) + query := fmt.Sprintf("SHOW %s KEYS ON %s FROM %s", tp, database, output.Measurement) + dbRtn, err := db.Query(query) if err != nil { - return errors.Wrap(err, "SHOW MEASUREMENTS") + return errors.Wrapf(err, "SHOW MEASUREMENTS: %s", query) } if len(dbRtn) == 0 || len(dbRtn[0]) == 0 { return nil @@ -820,6 +826,94 @@ func getAttributesOnMeasurement(database, tp string, output *monitor.InfluxMeasu return nil } +func getTagValues(userCred mcclient.TokenCredential, output *monitor.InfluxMeasurement, timeF timeFilter, dsId string, tagFilter string) error { + mq := monitor.MetricQuery{ + Database: output.Database, + Measurement: output.Measurement, + Selects: []monitor.MetricQuerySelect{ + { + { + Type: "field", + Params: []string{output.FieldKey[0]}, + }, + { + Type: "last", + }, + }, + }, + GroupBy: []monitor.MetricQueryPart{ + { + Type: "field", + Params: []string{"*"}, + }, + }, + } + if tagFilter != "" { + parts := strings.Split(tagFilter, " ") + mq.Tags = []monitor.MetricQueryTag{ + { + Key: parts[0], + Operator: parts[1], + Value: parts[2], + }, + } + } + + aq := &monitor.AlertQuery{ + Model: mq, + From: timeF.From, + To: timeF.To, + DataSourceId: dsId, + } + + q := monitor.MetricInputQuery{ + From: timeF.From, + To: timeF.To, + MetricQuery: []*monitor.AlertQuery{ + aq, + }, + ForceCheckSeries: true, + } + + ret, err := doQuery(userCred, q) + if err != nil { + return errors.Wrapf(err, "getTagValues query error %s", jsonutils.Marshal(q)) + } + + // 2. group tag and values + tagValMap := make(map[string][]string) + tagKeys := make([]string, 0) + if len(ret.Series) == 0 { + return nil + } + + for _, s := range ret.Series { + tagMap := s.Tags + for key, valStr := range tagMap { + valStr = renderTagVal(valStr) + if len(valStr) == 0 || valStr == "null" || filterTagValue(valStr) { + continue + } + if filterTagKey(key) { + continue + } + if valArr, ok := tagValMap[key]; ok { + if !utils.IsInStringArray(valStr, valArr) { + tagValMap[key] = append(valArr, valStr) + } + continue + } + tagValMap[key] = []string{valStr} + tagKeys = append(tagKeys, key) + } + } + output.TagValue = tagValMap + sort.Strings(tagKeys) + output.TagKey = tagKeys + + return nil +} + func getTagValue(database string, output *monitor.InfluxMeasurement, db *influxdb.SInfluxdb) error { if len(output.TagKey) == 0 { return nil @@ -868,7 +962,7 @@ func (self *SDataSourceManager) getFilterMeasurementTagValue(tagValueChan *influ measurement monitor.InfluxMeasurement, db *influxdb.SInfluxdb, tagFilter string) error { var buffer bytes.Buffer buffer.WriteString(fmt.Sprintf(`SELECT last("%s") FROM "%s" WHERE %s `, field, measurement.Measurement, - self.renderTimeFilter(from, to))) + renderTimeFilter(from, to))) if len(tagFilter) != 0 { buffer.WriteString(fmt.Sprintf(` AND %s `, tagFilter)) } @@ -885,7 +979,7 @@ func (self *SDataSourceManager) getFilterMeasurementTagValue(tagValueChan *influ tagMap, _ := rtn[rtnIndex][serieIndex].Tags.GetMap() for key, valObj := range tagMap { valStr, _ := valObj.GetString() - valStr = self._renderTagVal(valStr) + valStr = renderTagVal(valStr) if len(valStr) == 0 || valStr == "null" || filterTagValue(valStr) { continue } @@ -909,7 +1003,7 @@ func (self *SDataSourceManager) getFilterMeasurementTagValue(tagValueChan *influ return nil } -func (self *SDataSourceManager) _renderTagVal(val string) string { +func renderTagVal(val string) string { return strings.ReplaceAll(val, "+", " ") } func floatEquals(a, b float64) bool { diff --git a/pkg/monitor/models/unifiedmonitor.go b/pkg/monitor/models/unifiedmonitor.go index e063625749..ac2cfc2276 100644 --- a/pkg/monitor/models/unifiedmonitor.go +++ b/pkg/monitor/models/unifiedmonitor.go @@ -166,8 +166,7 @@ func getProjectIdFilterByProject(projectId string) (string, error) { return fmt.Sprintf(`"%s" =~ /%s/`, "tenant_id", projectId), nil } -func (self *SUnifiedMonitorManager) GetPropertyMetricMeasurement(ctx context.Context, userCred mcclient.TokenCredential, - query jsonutils.JSONObject) (jsonutils.JSONObject, error) { +func (self *SUnifiedMonitorManager) GetPropertyMetricMeasurement(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) (jsonutils.JSONObject, error) { metricFunc := monitor.MetricFunc{ FieldOptType: monitor.UNIFIED_MONITOR_FIELD_OPT_TYPE, FieldOptValue: monitor.UNIFIED_MONITOR_FIELD_OPT_VALUE, @@ -176,11 +175,11 @@ func (self *SUnifiedMonitorManager) GetPropertyMetricMeasurement(ctx context.Con } filter, err := getTagFilterByRequestQuery(ctx, userCred, query) if err != nil { - return nil, err + return nil, errors.Wrapf(err, "getTagFilterByRequestQuery %s", query.String()) } - rtn, err := DataSourceManager.GetMetricMeasurement(query, filter) + rtn, err := DataSourceManager.GetMetricMeasurement(userCred, query, filter) if err != nil { - return nil, err + return nil, errors.Wrapf(err, "GetMetricMeasurement by query %s, filter %s", query.String(), filter) } rtn.(*jsonutils.JSONDict).Add(jsonutils.Marshal(&metricFunc), "func") return rtn, nil @@ -209,9 +208,8 @@ func (self *SUnifiedMonitorManager) PerformQuery(ctx context.Context, userCred m ownId = userCred } setDefaultValue(q, inputQuery, scope, ownId) - err = self.ValidateInputQuery(q) - if err != nil { - return jsonutils.NewDict(), err + if err := self.ValidateInputQuery(q); err != nil { + return jsonutils.NewDict(), errors.Wrapf(err, "ValidateInputQuery") } } @@ -224,13 +222,13 @@ func (self *SUnifiedMonitorManager) PerformQuery(ctx context.Context, userCred m } } - rtn, err := doQuery(*inputQuery) + rtn, err := doQuery(userCred, *inputQuery) if err != nil { - return jsonutils.NewDict(), err + return jsonutils.NewDict(), errors.Wrapf(err, "doQuery with input %s", data) } if len(inputQuery.Soffset) != 0 && len(inputQuery.Slimit) != 0 { - seriesTotal := self.fillSearchSeriesTotalQuery(*inputQuery.MetricQuery[0]) + seriesTotal := self.fillSearchSeriesTotalQuery(userCred, *inputQuery.MetricQuery[0]) rtn.SeriesTotal = seriesTotal } @@ -239,18 +237,17 @@ func (self *SUnifiedMonitorManager) PerformQuery(ctx context.Context, userCred m return jsonutils.Marshal(rtn), nil } -func (self *SUnifiedMonitorManager) fillSearchSeriesTotalQuery(fork monitor.AlertQuery) int64 { +func (self *SUnifiedMonitorManager) fillSearchSeriesTotalQuery(userCred mcclient.TokenCredential, fork monitor.AlertQuery) int64 { newGroupByPart := make([]monitor.MetricQueryPart, 0) newGroupByPart = append(newGroupByPart, fork.Model.GroupBy[0]) fork.Model.GroupBy = newGroupByPart forkInputQury := new(monitor.MetricInputQuery) forkInputQury.MetricQuery = []*monitor.AlertQuery{&fork} - rtn, err := doQuery(*forkInputQury) + rtn, err := doQuery(userCred, *forkInputQury) if err != nil { log.Errorf("exec forkInputQury err:%v", err) return 0 } - log.Errorf("series len:%d", len(rtn.Series)) return int64(len(rtn.Series)) } @@ -283,7 +280,7 @@ func (self *SUnifiedMonitorManager) handleDataPreSignature(ctx context.Context, } } -func doQuery(query monitor.MetricInputQuery) (*mq.Metrics, error) { +func doQuery(userCred mcclient.TokenCredential, query monitor.MetricInputQuery) (*mq.Metrics, error) { conditions := make([]*monitor.AlertCondition, 0) for _, q := range query.MetricQuery { condition := monitor.AlertCondition{ @@ -297,7 +294,7 @@ func doQuery(query monitor.MetricInputQuery) (*mq.Metrics, error) { if err != nil { return nil, err } - metrics, err := metricQ.ExecuteQuery() + metrics, err := metricQ.ExecuteQuery(userCred, query.ForceCheckSeries) if err != nil { return nil, err } @@ -441,7 +438,7 @@ func checkQueryGroupBy(query *monitor.AlertQuery, inputQuery *monitor.MetricInpu if metricMeasurement != nil { tagId = monitor.MEASUREMENT_TAG_ID[metricMeasurement.ResType] } - if len(tagId) == 0 { + if len(tagId) == 0 || (len(inputQuery.Slimit) != 0 && len(inputQuery.Soffset) != 0) { tagId = "*" } query.Model.GroupBy = append(query.Model.GroupBy, @@ -452,13 +449,12 @@ func checkQueryGroupBy(query *monitor.AlertQuery, inputQuery *monitor.MetricInpu } func setSerieRowName(series *tsdb.TimeSeriesSlice, groupTag []string) { - //Add rowname,The front end displays the curve according to rowname + // Add rowname,The front end displays the curve according to rowname var index, unknownIndex = 1, 1 for i, serie := range *series { - //setRowName by groupTag + // setRowName by groupTag if len(groupTag) != 0 { for key, val := range serie.Tags { - if strings.Contains(strings.Join(groupTag, ","), key) { serie.RawName = fmt.Sprintf("%s", val) (*series)[i] = serie @@ -468,7 +464,7 @@ func setSerieRowName(series *tsdb.TimeSeriesSlice, groupTag []string) { continue } measurement := strings.Split(serie.Name, ".")[0] - //sep measurement set RowName by spe param + // sep measurement set RowName by spe param measurements, _ := MetricMeasurementManager.getMeasurementByName(measurement) if len(measurements) != 0 { if key, ok := monitor.MEASUREMENT_TAG_KEYWORD[measurements[0].ResType]; ok { @@ -478,7 +474,7 @@ func setSerieRowName(series *tsdb.TimeSeriesSlice, groupTag []string) { continue } } - //other condition set RowName + // other condition set RowName for key, val := range serie.Tags { if strings.Contains(key, "id") { serie.RawName = fmt.Sprintf("%d: %s", index, val)