mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
fix(monitor): query resources according by listing region resources (#16074)
This commit is contained in:
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user