From 169ffc15c6a1813d93bac5f8e355d469c296d4d1 Mon Sep 17 00:00:00 2001 From: Gareth Date: Fri, 12 Dec 2025 18:36:52 +0900 Subject: [PATCH] OpenTSDB: Run suggest queries through the data source backend (#114990) * OpenTSDB: Run suggest queries through the data source backend * use mux --- pkg/tsdb/opentsdb/callresource.go | 67 +++++ pkg/tsdb/opentsdb/opentsdb.go | 257 +---------------- pkg/tsdb/opentsdb/opentsdb_test.go | 46 ++- pkg/tsdb/opentsdb/standalone/datasource.go | 17 +- pkg/tsdb/opentsdb/utils.go | 273 ++++++++++++++++++ .../plugins/datasource/opentsdb/datasource.ts | 6 +- 6 files changed, 389 insertions(+), 277 deletions(-) create mode 100644 pkg/tsdb/opentsdb/callresource.go diff --git a/pkg/tsdb/opentsdb/callresource.go b/pkg/tsdb/opentsdb/callresource.go new file mode 100644 index 00000000000..be0f81b9c80 --- /dev/null +++ b/pkg/tsdb/opentsdb/callresource.go @@ -0,0 +1,67 @@ +package opentsdb + +import ( + "fmt" + "net/http" + "net/url" + "path" + + "github.com/grafana/grafana-plugin-sdk-go/backend" +) + +func (s *Service) HandleSuggestQuery(rw http.ResponseWriter, req *http.Request) { + logger := logger.FromContext(req.Context()) + + dsInfo, err := s.getDSInfo(req.Context(), backend.PluginConfigFromContext(req.Context())) + if err != nil { + http.Error(rw, fmt.Sprintf("failed to get datasource info: %v", err), http.StatusInternalServerError) + return + } + + u, err := url.Parse(dsInfo.URL) + if err != nil { + http.Error(rw, fmt.Sprintf("failed to parse datasource URL: %v", err), http.StatusInternalServerError) + return + } + + u.Path = path.Join(u.Path, "api/suggest") + u.RawQuery = req.URL.RawQuery + httpReq, err := http.NewRequestWithContext(req.Context(), http.MethodGet, u.String(), nil) + if err != nil { + http.Error(rw, fmt.Sprintf("failed to create request: %v", err), http.StatusInternalServerError) + return + } + + res, err := dsInfo.HTTPClient.Do(httpReq) + if err != nil { + http.Error(rw, fmt.Sprintf("failed to execute request: %v", err), http.StatusInternalServerError) + return + } + + defer func() { + if err := res.Body.Close(); err != nil { + logger.Error("Failed to close response body", "error", err) + } + }() + + responseBody, err := DecodeResponseBody(res, logger) + if err != nil { + http.Error(rw, fmt.Sprintf("failed to decode response: %v", err), http.StatusInternalServerError) + return + } + + for name, values := range res.Header { + if name == "Content-Encoding" || name == "Content-Length" { + continue + } + for _, value := range values { + rw.Header().Add(name, value) + } + } + + rw.WriteHeader(res.StatusCode) + if _, err := rw.Write(responseBody); err != nil { + logger.Error("Failed to write response", "error", err) + return + } +} diff --git a/pkg/tsdb/opentsdb/opentsdb.go b/pkg/tsdb/opentsdb/opentsdb.go index d00242f6432..a694445e1cd 100644 --- a/pkg/tsdb/opentsdb/opentsdb.go +++ b/pkg/tsdb/opentsdb/opentsdb.go @@ -4,21 +4,15 @@ import ( "context" "encoding/json" "fmt" - "io" "net/http" "net/url" "path" - "sort" - "strconv" - "strings" - "time" "github.com/grafana/grafana-plugin-sdk-go/backend" "github.com/grafana/grafana-plugin-sdk-go/backend/datasource" "github.com/grafana/grafana-plugin-sdk-go/backend/httpclient" "github.com/grafana/grafana-plugin-sdk-go/backend/instancemgmt" - "github.com/grafana/grafana-plugin-sdk-go/backend/log" - "github.com/grafana/grafana-plugin-sdk-go/data" + "github.com/grafana/grafana-plugin-sdk-go/backend/resource/httpadapter" ) var logger = backend.NewLoggerWith("tsdb.opentsdb") @@ -155,6 +149,14 @@ func (s *Service) CheckHealth(ctx context.Context, req *backend.CheckHealthReque }, nil } +func (s *Service) CallResource(ctx context.Context, req *backend.CallResourceRequest, sender backend.CallResourceResponseSender) error { + mux := http.NewServeMux() + mux.HandleFunc("/api/suggest", s.HandleSuggestQuery) + + handler := httpadapter.New(mux) + return handler.CallResource(ctx, req, sender) +} + func (s *Service) QueryData(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) { logger := logger.FromContext(ctx) @@ -166,16 +168,15 @@ func (s *Service) QueryData(ctx context.Context, req *backend.QueryDataRequest) result := backend.NewQueryDataResponse() for _, query := range req.Queries { - // Build OpenTsdbQuery with per-query time range tsdbQuery := OpenTsdbQuery{ Start: query.TimeRange.From.Unix(), End: query.TimeRange.To.Unix(), Queries: []map[string]any{ - s.buildMetric(query), + BuildMetric(query), }, } - httpReq, err := s.createRequest(ctx, logger, dsInfo, tsdbQuery) + httpReq, err := CreateRequest(ctx, logger, dsInfo, tsdbQuery) if err != nil { return nil, err } @@ -191,251 +192,17 @@ func (s *Service) QueryData(ctx context.Context, req *backend.QueryDataRequest) } }() - queryRes, err := s.parseResponse(logger, httpRes, query.RefID, dsInfo.TSDBVersion) + queryRes, err := ParseResponse(logger, httpRes, query.RefID, dsInfo.TSDBVersion) if err != nil { return nil, err } - // Attach parsed result for this query's RefID result.Responses[query.RefID] = queryRes.Responses[query.RefID] } return result, nil } -func (s *Service) createRequest(ctx context.Context, logger log.Logger, dsInfo *datasourceInfo, data OpenTsdbQuery) (*http.Request, error) { - u, err := url.Parse(dsInfo.URL) - if err != nil { - return nil, err - } - u.Path = path.Join(u.Path, "api/query") - if dsInfo.TSDBVersion == 4 { - queryParams := u.Query() - queryParams.Set("arrays", "true") - u.RawQuery = queryParams.Encode() - } - - postData, err := json.Marshal(data) - if err != nil { - logger.Info("Failed marshaling data", "error", err) - return nil, fmt.Errorf("failed to create request: %w", err) - } - - req, err := http.NewRequestWithContext(ctx, http.MethodPost, u.String(), strings.NewReader(string(postData))) - if err != nil { - logger.Info("Failed to create request", "error", err) - return nil, fmt.Errorf("failed to create request: %w", err) - } - - req.Header.Set("Content-Type", "application/json") - return req, nil -} - -func createInitialFrame(val OpenTsdbCommon, length int, refID string) *data.Frame { - labels := data.Labels{} - for label, value := range val.Tags { - labels[label] = value - } - - tagKeys := make([]string, 0, len(val.Tags)+len(val.AggregateTags)) - for tagKey := range val.Tags { - tagKeys = append(tagKeys, tagKey) - } - sort.Strings(tagKeys) - tagKeys = append(tagKeys, val.AggregateTags...) - - frame := data.NewFrameOfFieldTypes(val.Metric, length, data.FieldTypeTime, data.FieldTypeFloat64) - frame.Meta = &data.FrameMeta{ - Type: data.FrameTypeTimeSeriesMulti, - TypeVersion: data.FrameTypeVersion{0, 1}, - Custom: map[string]any{"tagKeys": tagKeys}, - } - frame.RefID = refID - timeField := frame.Fields[0] - timeField.Name = data.TimeSeriesTimeFieldName - dataField := frame.Fields[1] - dataField.Name = val.Metric - dataField.Labels = labels - - return frame -} - -// Parse response function for OpenTSDB version 2.4 -func parseResponse24(responseData []OpenTsdbResponse24, refID string, frames data.Frames) data.Frames { - for _, val := range responseData { - frame := createInitialFrame(val.OpenTsdbCommon, len(val.DataPoints), refID) - - for i, point := range val.DataPoints { - frame.SetRow(i, time.Unix(int64(point[0]), 0).UTC(), point[1]) - } - - frames = append(frames, frame) - } - - return frames -} - -// Parse response function for OpenTSDB versions < 2.4 -func parseResponseLT24(responseData []OpenTsdbResponse, refID string, frames data.Frames) (data.Frames, error) { - for _, val := range responseData { - frame := createInitialFrame(val.OpenTsdbCommon, len(val.DataPoints), refID) - - // Order the timestamps in ascending order to avoid issues like https://github.com/grafana/grafana/issues/38729 - timestamps := make([]string, 0, len(val.DataPoints)) - for timestamp := range val.DataPoints { - timestamps = append(timestamps, timestamp) - } - sort.Strings(timestamps) - - for i, timeString := range timestamps { - timestamp, err := strconv.ParseInt(timeString, 10, 64) - if err != nil { - logger.Info("Failed to unmarshal opentsdb timestamp", "timestamp", timeString) - return frames, err - } - frame.SetRow(i, time.Unix(timestamp, 0).UTC(), val.DataPoints[timeString]) - } - - frames = append(frames, frame) - } - - return frames, nil -} - -func (s *Service) parseResponse(logger log.Logger, res *http.Response, refID string, tsdbVersion float32) (*backend.QueryDataResponse, error) { - resp := backend.NewQueryDataResponse() - - body, err := io.ReadAll(res.Body) - if err != nil { - return nil, err - } - defer func() { - if err := res.Body.Close(); err != nil { - logger.Warn("Failed to close response body", "err", err) - } - }() - - if res.StatusCode/100 != 2 { - logger.Info("Request failed", "status", res.Status, "body", string(body)) - return nil, fmt.Errorf("request failed, status: %s", res.Status) - } - - frames := data.Frames{} - - var responseData []OpenTsdbResponse - var responseData24 []OpenTsdbResponse24 - if tsdbVersion == 4 { - err = json.Unmarshal(body, &responseData24) - if err != nil { - logger.Info("Failed to unmarshal opentsdb response", "error", err, "status", res.Status, "body", string(body)) - return nil, err - } - - frames = parseResponse24(responseData24, refID, frames) - } else { - err = json.Unmarshal(body, &responseData) - if err != nil { - logger.Info("Failed to unmarshal opentsdb response", "error", err, "status", res.Status, "body", string(body)) - return nil, err - } - - frames, err = parseResponseLT24(responseData, refID, frames) - if err != nil { - return nil, err - } - } - - result := resp.Responses[refID] - result.Frames = frames - resp.Responses[refID] = result - return resp, nil -} - -func (s *Service) buildMetric(query backend.DataQuery) map[string]any { - metric := make(map[string]any) - - var model QueryModel - if err := json.Unmarshal(query.JSON, &model); err != nil { - return nil - } - - // Setting metric and aggregator - metric["metric"] = model.Metric - metric["aggregator"] = model.Aggregator - - // Setting downsampling options - if !model.DisableDownsampling { - downsampleInterval := model.DownsampleInterval - if downsampleInterval == "" { - if ms := query.Interval.Milliseconds(); ms > 0 { - downsampleInterval = FormatDownsampleInterval(ms) - } else { - downsampleInterval = "1m" - } - } else if strings.Contains(downsampleInterval, ".") && strings.HasSuffix(downsampleInterval, "s") { - if val, err := strconv.ParseFloat(strings.TrimSuffix(downsampleInterval, "s"), 64); err == nil { - downsampleInterval = strconv.FormatInt(int64(val*1000), 10) + "ms" - } - } - - downsample := downsampleInterval + "-" + model.DownsampleAggregator - if model.DownsampleFillPolicy != "" && model.DownsampleFillPolicy != "none" { - metric["downsample"] = downsample + "-" + model.DownsampleFillPolicy - } else { - metric["downsample"] = downsample - } - } - - // Setting rate options - if model.ShouldComputeRate { - metric["rate"] = true - rateOptions := make(map[string]any) - rateOptions["counter"] = model.IsCounter - - var counterMax *float64 - if model.CounterMax != "" { - if val, err := strconv.ParseFloat(model.CounterMax, 64); err == nil { - counterMax = &val - } - } - if counterMax != nil { - rateOptions["counterMax"] = *counterMax - } - - var counterResetValue *float64 - if model.CounterResetValue != "" { - if val, err := strconv.ParseFloat(model.CounterResetValue, 64); err == nil { - counterResetValue = &val - } - } - if counterResetValue != nil { - rateOptions["resetValue"] = *counterResetValue - } - - if counterMax == nil && (counterResetValue == nil || *counterResetValue == 0) { - rateOptions["dropResets"] = true - } - - metric["rateOptions"] = rateOptions - } - - // Setting tags - if len(model.Tags) > 0 { - metric["tags"] = model.Tags - } - - // Setting filters - if len(model.Filters) > 0 { - metric["filters"] = model.Filters - } - - if model.ExplicitTags { - metric["explicitTags"] = true - } - - return metric -} - func (s *Service) getDSInfo(ctx context.Context, pluginCtx backend.PluginContext) (*datasourceInfo, error) { i, err := s.im.Get(ctx, pluginCtx) if err != nil { diff --git a/pkg/tsdb/opentsdb/opentsdb_test.go b/pkg/tsdb/opentsdb/opentsdb_test.go index a0150b9e75f..72bbe6bf8b9 100644 --- a/pkg/tsdb/opentsdb/opentsdb_test.go +++ b/pkg/tsdb/opentsdb/opentsdb_test.go @@ -71,8 +71,6 @@ func TestCheckHealth(t *testing.T) { } func TestBuildMetric(t *testing.T) { - service := &Service{} - t.Run("Metric with no downsampleInterval should use query interval", func(t *testing.T) { query := backend.DataQuery{ JSON: []byte(` @@ -88,7 +86,7 @@ func TestBuildMetric(t *testing.T) { Interval: 30 * time.Second, } - metric := service.buildMetric(query) + metric := BuildMetric(query) require.Equal(t, "30s-avg", metric["downsample"], "should use query interval formatted as seconds") }) @@ -106,7 +104,7 @@ func TestBuildMetric(t *testing.T) { ), } - metric := service.buildMetric(query) + metric := BuildMetric(query) require.Equal(t, "500ms-avg", metric["downsample"], "should convert 0.5s to 500ms") }) @@ -125,7 +123,7 @@ func TestBuildMetric(t *testing.T) { Interval: 500 * time.Millisecond, } - metric := service.buildMetric(query) + metric := BuildMetric(query) require.Equal(t, "500ms-avg", metric["downsample"], "should use query interval formatted as milliseconds") }) @@ -144,7 +142,7 @@ func TestBuildMetric(t *testing.T) { Interval: 5 * time.Minute, } - metric := service.buildMetric(query) + metric := BuildMetric(query) require.Equal(t, "5m-sum", metric["downsample"], "should use query interval formatted as minutes") }) @@ -163,7 +161,7 @@ func TestBuildMetric(t *testing.T) { Interval: 2 * time.Hour, } - metric := service.buildMetric(query) + metric := BuildMetric(query) require.Equal(t, "2h-max", metric["downsample"], "should use query interval formatted as hours") }) @@ -182,7 +180,7 @@ func TestBuildMetric(t *testing.T) { Interval: 48 * time.Hour, } - metric := service.buildMetric(query) + metric := BuildMetric(query) require.Equal(t, "2d-min", metric["downsample"], "should use query interval formatted as days") }) @@ -201,7 +199,7 @@ func TestBuildMetric(t *testing.T) { ), } - metric := service.buildMetric(query) + metric := BuildMetric(query) require.True(t, metric["explicitTags"].(bool), "explicitTags should be true") metricTags := metric["tags"].(map[string]any) @@ -223,16 +221,14 @@ func TestBuildMetric(t *testing.T) { ), } - metric := service.buildMetric(query) + metric := BuildMetric(query) require.Nil(t, metric["explicitTags"], "explicitTags should not be present when false") }) } func TestOpenTsdbExecutor(t *testing.T) { - service := &Service{} - t.Run("create request", func(t *testing.T) { - req, err := service.createRequest(context.Background(), logger, &datasourceInfo{}, OpenTsdbQuery{}) + req, err := CreateRequest(context.Background(), logger, &datasourceInfo{}, OpenTsdbQuery{}) require.NoError(t, err) assert.Equal(t, "POST", req.Method) @@ -247,7 +243,7 @@ func TestOpenTsdbExecutor(t *testing.T) { response := `{ invalid }` tsdbVersion := float32(4) - result, err := service.parseResponse(logger, &http.Response{Body: io.NopCloser(strings.NewReader(response))}, "A", tsdbVersion) + result, err := ParseResponse(logger, &http.Response{Body: io.NopCloser(strings.NewReader(response))}, "A", tsdbVersion) require.Nil(t, result) require.Error(t, err) }) @@ -284,7 +280,7 @@ func TestOpenTsdbExecutor(t *testing.T) { resp := http.Response{Body: io.NopCloser(strings.NewReader(response))} resp.StatusCode = 200 - result, err := service.parseResponse(logger, &resp, "A", tsdbVersion) + result, err := ParseResponse(logger, &resp, "A", tsdbVersion) require.NoError(t, err) frame := result.Responses["A"] @@ -326,7 +322,7 @@ func TestOpenTsdbExecutor(t *testing.T) { resp := http.Response{Body: io.NopCloser(strings.NewReader(response))} resp.StatusCode = 200 - result, err := service.parseResponse(logger, &resp, "A", tsdbVersion) + result, err := ParseResponse(logger, &resp, "A", tsdbVersion) require.NoError(t, err) frame := result.Responses["A"] @@ -399,7 +395,7 @@ func TestOpenTsdbExecutor(t *testing.T) { resp := http.Response{Body: io.NopCloser(strings.NewReader(response))} resp.StatusCode = 200 - result, err := service.parseResponse(logger, &resp, "A", tsdbVersion) + result, err := ParseResponse(logger, &resp, "A", tsdbVersion) require.NoError(t, err) frame := result.Responses["A"] @@ -444,7 +440,7 @@ func TestOpenTsdbExecutor(t *testing.T) { resp := http.Response{Body: io.NopCloser(strings.NewReader(response))} resp.StatusCode = 200 - result, err := service.parseResponse(logger, &resp, myRefid, tsdbVersion) + result, err := ParseResponse(logger, &resp, myRefid, tsdbVersion) require.NoError(t, err) if diff := cmp.Diff(testFrame, result.Responses[myRefid].Frames[0], data.FrameTestCompareOptions()...); diff != "" { @@ -473,7 +469,7 @@ func TestOpenTsdbExecutor(t *testing.T) { resp := http.Response{Body: io.NopCloser(strings.NewReader(response))} resp.StatusCode = 200 - result, err := service.parseResponse(logger, &resp, "A", tsdbVersion) + result, err := ParseResponse(logger, &resp, "A", tsdbVersion) require.NoError(t, err) frame := result.Responses["A"].Frames[0] @@ -505,7 +501,7 @@ func TestOpenTsdbExecutor(t *testing.T) { ), } - metric := service.buildMetric(query) + metric := BuildMetric(query) require.Len(t, metric, 3) require.Equal(t, "cpu.average.percent", metric["metric"]) @@ -527,7 +523,7 @@ func TestOpenTsdbExecutor(t *testing.T) { ), } - metric := service.buildMetric(query) + metric := BuildMetric(query) require.Len(t, metric, 2) require.Equal(t, "cpu.average.percent", metric["metric"]) @@ -548,7 +544,7 @@ func TestOpenTsdbExecutor(t *testing.T) { ), } - metric := service.buildMetric(query) + metric := BuildMetric(query) require.Len(t, metric, 3) require.Equal(t, "cpu.average.percent", metric["metric"]) @@ -574,7 +570,7 @@ func TestOpenTsdbExecutor(t *testing.T) { ), } - metric := service.buildMetric(query) + metric := BuildMetric(query) require.Len(t, metric, 3) require.Equal(t, "cpu.average.percent", metric["metric"]) @@ -605,7 +601,7 @@ func TestOpenTsdbExecutor(t *testing.T) { ), } - metric := service.buildMetric(query) + metric := BuildMetric(query) require.Len(t, metric, 5) require.Equal(t, "cpu.average.percent", metric["metric"]) @@ -640,7 +636,7 @@ func TestOpenTsdbExecutor(t *testing.T) { ), } - metric := service.buildMetric(query) + metric := BuildMetric(query) t.Log(metric) require.Len(t, metric, 5) diff --git a/pkg/tsdb/opentsdb/standalone/datasource.go b/pkg/tsdb/opentsdb/standalone/datasource.go index c2eacaf1d53..97241d651d6 100644 --- a/pkg/tsdb/opentsdb/standalone/datasource.go +++ b/pkg/tsdb/opentsdb/standalone/datasource.go @@ -10,8 +10,9 @@ import ( ) var ( - _ backend.QueryDataHandler = (*Datasource)(nil) - _ backend.CheckHealthHandler = (*Datasource)(nil) + _ backend.CheckHealthHandler = (*Datasource)(nil) + _ backend.CallResourceHandler = (*Datasource)(nil) + _ backend.QueryDataHandler = (*Datasource)(nil) ) type Datasource struct { @@ -24,10 +25,14 @@ func NewDatasource(context.Context, backend.DataSourceInstanceSettings) (instanc }, nil } -func (d *Datasource) QueryData(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) { - return d.Service.QueryData(ctx, req) -} - func (d *Datasource) CheckHealth(ctx context.Context, req *backend.CheckHealthRequest) (*backend.CheckHealthResult, error) { return d.Service.CheckHealth(ctx, req) } + +func (d *Datasource) CallResource(ctx context.Context, req *backend.CallResourceRequest, sender backend.CallResourceResponseSender) error { + return d.Service.CallResource(ctx, req, sender) +} + +func (d *Datasource) QueryData(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) { + return d.Service.QueryData(ctx, req) +} diff --git a/pkg/tsdb/opentsdb/utils.go b/pkg/tsdb/opentsdb/utils.go index ae57b3e787a..ddfa8122fce 100644 --- a/pkg/tsdb/opentsdb/utils.go +++ b/pkg/tsdb/opentsdb/utils.go @@ -1,8 +1,22 @@ package opentsdb import ( + "compress/gzip" + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "net/url" + "path" + "sort" "strconv" + "strings" "time" + + "github.com/grafana/grafana-plugin-sdk-go/backend" + "github.com/grafana/grafana-plugin-sdk-go/backend/log" + "github.com/grafana/grafana-plugin-sdk-go/data" ) func FormatDownsampleInterval(ms int64) string { @@ -29,3 +43,262 @@ func FormatDownsampleInterval(ms int64) string { days := int64(duration / (24 * time.Hour)) return strconv.FormatInt(days, 10) + "d" } + +func BuildMetric(query backend.DataQuery) map[string]any { + metric := make(map[string]any) + + var model QueryModel + if err := json.Unmarshal(query.JSON, &model); err != nil { + return nil + } + + // Setting metric and aggregator + metric["metric"] = model.Metric + metric["aggregator"] = model.Aggregator + + // Setting downsampling options + if !model.DisableDownsampling { + downsampleInterval := model.DownsampleInterval + if downsampleInterval == "" { + if ms := query.Interval.Milliseconds(); ms > 0 { + downsampleInterval = FormatDownsampleInterval(ms) + } else { + downsampleInterval = "1m" + } + } else if strings.Contains(downsampleInterval, ".") && strings.HasSuffix(downsampleInterval, "s") { + if val, err := strconv.ParseFloat(strings.TrimSuffix(downsampleInterval, "s"), 64); err == nil { + downsampleInterval = strconv.FormatInt(int64(val*1000), 10) + "ms" + } + } + + downsample := downsampleInterval + "-" + model.DownsampleAggregator + if model.DownsampleFillPolicy != "" && model.DownsampleFillPolicy != "none" { + metric["downsample"] = downsample + "-" + model.DownsampleFillPolicy + } else { + metric["downsample"] = downsample + } + } + + // Setting rate options + if model.ShouldComputeRate { + metric["rate"] = true + rateOptions := make(map[string]any) + rateOptions["counter"] = model.IsCounter + + var counterMax *float64 + if model.CounterMax != "" { + if val, err := strconv.ParseFloat(model.CounterMax, 64); err == nil { + counterMax = &val + } + } + if counterMax != nil { + rateOptions["counterMax"] = *counterMax + } + + var counterResetValue *float64 + if model.CounterResetValue != "" { + if val, err := strconv.ParseFloat(model.CounterResetValue, 64); err == nil { + counterResetValue = &val + } + } + if counterResetValue != nil { + rateOptions["resetValue"] = *counterResetValue + } + + if counterMax == nil && (counterResetValue == nil || *counterResetValue == 0) { + rateOptions["dropResets"] = true + } + + metric["rateOptions"] = rateOptions + } + + // Setting tags + if len(model.Tags) > 0 { + metric["tags"] = model.Tags + } + + // Setting filters + if len(model.Filters) > 0 { + metric["filters"] = model.Filters + } + + if model.ExplicitTags { + metric["explicitTags"] = true + } + + return metric +} + +func CreateRequest(ctx context.Context, logger log.Logger, dsInfo *datasourceInfo, data OpenTsdbQuery) (*http.Request, error) { + u, err := url.Parse(dsInfo.URL) + if err != nil { + return nil, err + } + u.Path = path.Join(u.Path, "api/query") + if dsInfo.TSDBVersion == 4 { + queryParams := u.Query() + queryParams.Set("arrays", "true") + u.RawQuery = queryParams.Encode() + } + + postData, err := json.Marshal(data) + if err != nil { + logger.Info("Failed marshaling data", "error", err) + return nil, fmt.Errorf("failed to create request: %w", err) + } + + req, err := http.NewRequestWithContext(ctx, http.MethodPost, u.String(), strings.NewReader(string(postData))) + if err != nil { + logger.Info("Failed to create request", "error", err) + return nil, fmt.Errorf("failed to create request: %w", err) + } + + req.Header.Set("Content-Type", "application/json") + return req, nil +} + +func DecodeResponseBody(res *http.Response, logger log.Logger) ([]byte, error) { + encoding := res.Header.Get("Content-Encoding") + var reader io.Reader + + switch encoding { + case "gzip": + gzipReader, err := gzip.NewReader(res.Body) + if err != nil { + return nil, fmt.Errorf("failed to create gzip reader: %w", err) + } + defer func() { + if err := gzipReader.Close(); err != nil { + logger.Warn("Failed to close gzip reader", "error", err) + } + }() + reader = gzipReader + default: + reader = res.Body + } + + body, err := io.ReadAll(reader) + if err != nil { + return nil, fmt.Errorf("failed to read response body: %w", err) + } + + return body, nil +} + +func CreateDataFrame(val OpenTsdbCommon, length int, refID string) *data.Frame { + labels := data.Labels{} + for label, value := range val.Tags { + labels[label] = value + } + + tagKeys := make([]string, 0, len(val.Tags)+len(val.AggregateTags)) + for tagKey := range val.Tags { + tagKeys = append(tagKeys, tagKey) + } + sort.Strings(tagKeys) + tagKeys = append(tagKeys, val.AggregateTags...) + + frame := data.NewFrameOfFieldTypes(val.Metric, length, data.FieldTypeTime, data.FieldTypeFloat64) + frame.Meta = &data.FrameMeta{ + Type: data.FrameTypeTimeSeriesMulti, + TypeVersion: data.FrameTypeVersion{0, 1}, + Custom: map[string]any{"tagKeys": tagKeys}, + } + frame.RefID = refID + timeField := frame.Fields[0] + timeField.Name = data.TimeSeriesTimeFieldName + dataField := frame.Fields[1] + dataField.Name = val.Metric + dataField.Labels = labels + + return frame +} + +func ParseResponse(logger log.Logger, res *http.Response, refID string, tsdbVersion float32) (*backend.QueryDataResponse, error) { + resp := backend.NewQueryDataResponse() + + body, err := io.ReadAll(res.Body) + if err != nil { + return nil, err + } + defer func() { + if err := res.Body.Close(); err != nil { + logger.Warn("Failed to close response body", "err", err) + } + }() + + if res.StatusCode/100 != 2 { + logger.Info("Request failed", "status", res.Status, "body", string(body)) + return nil, fmt.Errorf("request failed, status: %s", res.Status) + } + + frames := data.Frames{} + + var responseData []OpenTsdbResponse + var responseData24 []OpenTsdbResponse24 + if tsdbVersion == 4 { + err = json.Unmarshal(body, &responseData24) + if err != nil { + logger.Info("Failed to unmarshal opentsdb response", "error", err, "status", res.Status, "body", string(body)) + return nil, err + } + + frames = ParseResponse24(responseData24, refID, frames) + } else { + err = json.Unmarshal(body, &responseData) + if err != nil { + logger.Info("Failed to unmarshal opentsdb response", "error", err, "status", res.Status, "body", string(body)) + return nil, err + } + + frames, err = ParseResponseLT24(responseData, refID, frames) + if err != nil { + return nil, err + } + } + + result := resp.Responses[refID] + result.Frames = frames + resp.Responses[refID] = result + return resp, nil +} + +func ParseResponse24(responseData []OpenTsdbResponse24, refID string, frames data.Frames) data.Frames { + for _, val := range responseData { + frame := CreateDataFrame(val.OpenTsdbCommon, len(val.DataPoints), refID) + + for i, point := range val.DataPoints { + frame.SetRow(i, time.Unix(int64(point[0]), 0).UTC(), point[1]) + } + + frames = append(frames, frame) + } + + return frames +} + +func ParseResponseLT24(responseData []OpenTsdbResponse, refID string, frames data.Frames) (data.Frames, error) { + for _, val := range responseData { + frame := CreateDataFrame(val.OpenTsdbCommon, len(val.DataPoints), refID) + + // Order the timestamps in ascending order to avoid issues like https://github.com/grafana/grafana/issues/38729 + timestamps := make([]string, 0, len(val.DataPoints)) + for timestamp := range val.DataPoints { + timestamps = append(timestamps, timestamp) + } + sort.Strings(timestamps) + + for i, timeString := range timestamps { + timestamp, err := strconv.ParseInt(timeString, 10, 64) + if err != nil { + logger.Info("Failed to unmarshal opentsdb timestamp", "timestamp", timeString) + return frames, err + } + frame.SetRow(i, time.Unix(timestamp, 0).UTC(), val.DataPoints[timeString]) + } + + frames = append(frames, frame) + } + + return frames, nil +} diff --git a/public/app/plugins/datasource/opentsdb/datasource.ts b/public/app/plugins/datasource/opentsdb/datasource.ts index 1413daacfc2..24356eefbac 100644 --- a/public/app/plugins/datasource/opentsdb/datasource.ts +++ b/public/app/plugins/datasource/opentsdb/datasource.ts @@ -13,7 +13,7 @@ import { map as _map, toPairs, } from 'lodash'; -import { lastValueFrom, merge, Observable, of } from 'rxjs'; +import { from, lastValueFrom, merge, Observable, of } from 'rxjs'; import { catchError, map } from 'rxjs/operators'; import { @@ -290,6 +290,10 @@ export default class OpenTsDatasource extends DataSourceWithBackend { return result.data;