diff --git a/pkg/tsdb/graphite/graphite.go b/pkg/tsdb/graphite/graphite.go index ad37275406c..406ba10f106 100644 --- a/pkg/tsdb/graphite/graphite.go +++ b/pkg/tsdb/graphite/graphite.go @@ -10,13 +10,16 @@ import ( "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/backend/resource/httpadapter" "go.opentelemetry.io/otel/trace" ) type Service struct { - im instancemgmt.InstanceManager - tracer trace.Tracer - logger log.Logger + im instancemgmt.InstanceManager + tracer trace.Tracer + logger log.Logger + resourceHandler backend.CallResourceHandler + HTTPClient *http.Client } const ( @@ -26,11 +29,15 @@ const ( func ProvideService(httpClientProvider *httpclient.Provider, tracer trace.Tracer) *Service { logger := backend.NewLoggerWith("logger", "graphite") - return &Service{ + s := &Service{ im: datasource.NewInstanceManager(newInstanceSettings(httpClientProvider)), tracer: tracer, logger: logger, } + + s.resourceHandler = httpadapter.New(s.newResourceMux()) + + return s } type datasourceInfo struct { @@ -85,5 +92,5 @@ func (s *Service) QueryData(ctx context.Context, req *backend.QueryDataRequest) } func (s *Service) CallResource(ctx context.Context, req *backend.CallResourceRequest, sender backend.CallResourceResponseSender) error { - return nil + return s.resourceHandler.CallResource(ctx, req, sender) } diff --git a/pkg/tsdb/graphite/resource_handler.go b/pkg/tsdb/graphite/resource_handler.go new file mode 100644 index 00000000000..51eea3775d5 --- /dev/null +++ b/pkg/tsdb/graphite/resource_handler.go @@ -0,0 +1,155 @@ +package graphite + +import ( + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "net/url" + + "github.com/grafana/grafana-plugin-sdk-go/backend" + "github.com/grafana/grafana-plugin-sdk-go/backend/tracing" + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/codes" +) + +type resourceHandler func(context.Context, *datasourceInfo, []byte) ([]byte, int, error) + +func (s *Service) newResourceMux() *http.ServeMux { + mux := http.NewServeMux() + mux.HandleFunc("/events", s.handleResourceReq(s.handleEvents)) + return mux +} + +func (s *Service) handleResourceReq(handlerFn resourceHandler) func(rw http.ResponseWriter, req *http.Request) { + return func(rw http.ResponseWriter, req *http.Request) { + s.logger.Debug("Received resource call", "url", req.URL.String(), "method", req.Method) + + pluginCtx := backend.PluginConfigFromContext(req.Context()) + ctx := req.Context() + dsInfo, err := s.getDSInfo(ctx, pluginCtx) + if err != nil { + writeErrorResponse(rw, http.StatusInternalServerError, fmt.Sprintf("unexpected error %v", err)) + return + } + + defer func() { + if err := req.Body.Close(); err != nil { + s.logger.Warn("Failed to close response body", "err", err) + writeErrorResponse(rw, http.StatusInternalServerError, fmt.Sprintf("unexpected error %v", err)) + return + } + }() + requestBody, err := io.ReadAll(req.Body) + if err != nil { + s.logger.Error("Failed to read events request body", "error", err) + writeErrorResponse(rw, http.StatusInternalServerError, fmt.Sprintf("unexpected error %v", err)) + return + } + + if handlerFn == nil { + writeErrorResponse(rw, http.StatusInternalServerError, "responseFn should not be nil") + return + } + + response, statusCode, err := handlerFn(ctx, dsInfo, requestBody) + if err != nil { + writeErrorResponse(rw, statusCode, fmt.Sprintf("failed to handle resource request: %v", err)) + return + } + + rw.WriteHeader(statusCode) + _, err = rw.Write(response) + if err != nil { + writeErrorResponse(rw, http.StatusInternalServerError, fmt.Sprintf("failed to write events response: %v", err)) + return + } + } +} + +func (s *Service) handleEvents(ctx context.Context, dsInfo *datasourceInfo, requestBody []byte) ([]byte, int, error) { + eventsRequestJson := GraphiteEventsRequest{} + err := json.Unmarshal(requestBody, &eventsRequestJson) + if err != nil { + s.logger.Error("Failed to unmarshal events request body to JSON", "error", err) + return nil, http.StatusInternalServerError, fmt.Errorf("unexpected error %v", err) + } + + eventsUrl, err := url.Parse(fmt.Sprintf("%s/events/get_data", dsInfo.URL)) + if err != nil { + return nil, http.StatusInternalServerError, fmt.Errorf("unexpected error %v", err) + } + + queryValues := eventsUrl.Query() + queryValues.Set("from", eventsRequestJson.From) + queryValues.Set("until", eventsRequestJson.Until) + if eventsRequestJson.Tags != "" { + queryValues.Set("tags", eventsRequestJson.Tags) + } + + eventsUrl.RawQuery = queryValues.Encode() + + p := eventsUrl.String() + graphiteReq, err := http.NewRequestWithContext(ctx, http.MethodGet, p, nil) + if err != nil { + s.logger.Info("Failed to create request", "error", err) + return nil, http.StatusInternalServerError, fmt.Errorf("failed to create request: %v", err) + } + + _, span := tracing.DefaultTracer().Start(ctx, "graphite events") + defer span.End() + span.SetAttributes( + attribute.Int64("datasource_id", dsInfo.Id), + ) + res, err := dsInfo.HTTPClient.Do(graphiteReq) + if res != nil { + span.SetAttributes(attribute.Int("graphite.response.code", res.StatusCode)) + } + if err != nil { + span.RecordError(err) + span.SetStatus(codes.Error, err.Error()) + return nil, http.StatusInternalServerError, fmt.Errorf("failed to complete events request: %v", err) + } + + defer func() { + err := res.Body.Close() + if err != nil { + s.logger.Warn("Failed to close response body", "error", err) + } + }() + + encoding := res.Header.Get("Content-Encoding") + body, err := decode(encoding, res.Body) + if err != nil { + return nil, res.StatusCode, fmt.Errorf("failed to read events response: %v", err) + } + + events := []GraphiteEventsResponse{} + err = json.Unmarshal(body, &events) + if err != nil { + return nil, http.StatusInternalServerError, fmt.Errorf("failed to unmarshal events response: %v", err) + } + + // We construct this struct to avoid frontend changes. + graphiteEventsResponse, err := json.Marshal(map[string][]GraphiteEventsResponse{ + "data": events, + }) + if err != nil { + return nil, http.StatusInternalServerError, fmt.Errorf("failed to marshal events response: %s", err) + } + + return graphiteEventsResponse, res.StatusCode, nil +} + +func writeErrorResponse(rw http.ResponseWriter, code int, msg string) { + rw.WriteHeader(code) + errorBody := map[string]string{ + "error": msg, + } + jsonRes, _ := json.Marshal(errorBody) + _, err := rw.Write(jsonRes) + if err != nil { + backend.Logger.Error("Unable to write HTTP response", "error", err) + } +} diff --git a/pkg/tsdb/graphite/resource_handler_test.go b/pkg/tsdb/graphite/resource_handler_test.go new file mode 100644 index 00000000000..a1aeee3697d --- /dev/null +++ b/pkg/tsdb/graphite/resource_handler_test.go @@ -0,0 +1,267 @@ +package graphite + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "io" + "net/http" + "net/http/httptest" + "testing" + + "github.com/grafana/grafana-plugin-sdk-go/backend" + "github.com/grafana/grafana-plugin-sdk-go/backend/instancemgmt" + "github.com/grafana/grafana-plugin-sdk-go/backend/log" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +type mockRoundTripper struct { + respBody []byte + status int + err error +} + +func (m *mockRoundTripper) RoundTrip(req *http.Request) (*http.Response, error) { + if m.err != nil { + return nil, m.err + } + resp := &http.Response{ + StatusCode: m.status, + Body: io.NopCloser(bytes.NewBuffer(m.respBody)), + Header: make(http.Header), + } + return resp, nil +} + +type mockInstanceManager struct { + instance instancemgmt.Instance + err error +} + +func (m *mockInstanceManager) Get(ctx context.Context, pluginCtx backend.PluginContext) (instancemgmt.Instance, error) { + return m.instance, m.err +} +func (m *mockInstanceManager) Dispose(_ string) {} +func (m *mockInstanceManager) Do(ctx context.Context, pluginCtx backend.PluginContext, fn instancemgmt.InstanceCallbackFunc) error { + return nil +} + +func TestHandleEvents(t *testing.T) { + mockEvents := []GraphiteEventsResponse{ + {When: 1234567890, What: "event1", Tags: []string{"tag1"}, Data: "data1"}, + {When: 1234567891, What: "event2", Tags: []string{"tag2"}, Data: "data2"}, + } + mockResp, _ := json.Marshal(mockEvents) + + tests := []struct { + name string + dsInfo *datasourceInfo + requestBody []byte + expectedStatus int + expectError bool + errorContains string + expectedEvents []GraphiteEventsResponse + }{ + { + name: "Success with tags", + dsInfo: &datasourceInfo{ + Id: 1, + URL: "http://example.com", + HTTPClient: &http.Client{Transport: &mockRoundTripper{respBody: mockResp, status: 200}}, + }, + requestBody: func() []byte { + request := GraphiteEventsRequest{From: "now-1h", Until: "now", Tags: "foo"} + body, _ := json.Marshal(request) + return body + }(), + expectedStatus: 200, + expectError: false, + expectedEvents: mockEvents, + }, + { + name: "Success without tags", + dsInfo: &datasourceInfo{ + Id: 1, + URL: "http://example.com", + HTTPClient: &http.Client{Transport: &mockRoundTripper{respBody: mockResp, status: 200}}, + }, + requestBody: func() []byte { + request := GraphiteEventsRequest{From: "now-1h", Until: "now"} + body, _ := json.Marshal(request) + return body + }(), + expectedStatus: 200, + expectError: false, + expectedEvents: mockEvents, + }, + { + name: "Invalid request body", + dsInfo: &datasourceInfo{Id: 1, URL: "http://example.com"}, + requestBody: []byte(`{"invalid": json}`), + expectedStatus: http.StatusInternalServerError, + expectError: true, + errorContains: "unexpected error", + }, + { + name: "Invalid URL", + dsInfo: &datasourceInfo{ + Id: 1, + URL: "ht tp://invalid url", // Invalid URL + }, + requestBody: func() []byte { + request := GraphiteEventsRequest{From: "now-1h", Until: "now"} + body, _ := json.Marshal(request) + return body + }(), + expectedStatus: http.StatusInternalServerError, + expectError: true, + errorContains: "unexpected error", + }, + { + name: "HTTP client error", + dsInfo: &datasourceInfo{ + Id: 1, + URL: "http://example.com", + HTTPClient: &http.Client{Transport: &mockRoundTripper{err: errors.New("network error")}}, + }, + requestBody: func() []byte { + request := GraphiteEventsRequest{From: "now-1h", Until: "now"} + body, _ := json.Marshal(request) + return body + }(), + expectedStatus: http.StatusInternalServerError, + expectError: true, + errorContains: "failed to complete events request", + }, + { + name: "Invalid response JSON", + dsInfo: &datasourceInfo{ + Id: 1, + URL: "http://example.com", + HTTPClient: &http.Client{Transport: &mockRoundTripper{respBody: []byte("invalid json"), status: 200}}, + }, + requestBody: func() []byte { + request := GraphiteEventsRequest{From: "now-1h", Until: "now"} + body, _ := json.Marshal(request) + return body + }(), + expectedStatus: http.StatusInternalServerError, + expectError: true, + errorContains: "failed to unmarshal events response", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + svc := &Service{logger: log.NewNullLogger()} + + respBody, status, err := svc.handleEvents(context.Background(), tt.dsInfo, tt.requestBody) + + assert.Equal(t, tt.expectedStatus, status) + + if tt.expectError { + assert.Error(t, err) + assert.Nil(t, respBody) + if tt.errorContains != "" { + assert.Contains(t, err.Error(), tt.errorContains) + } + } else { + require.NoError(t, err) + assert.NotNil(t, respBody) + + if tt.expectedEvents != nil { + var result map[string][]GraphiteEventsResponse + require.NoError(t, json.Unmarshal(respBody, &result)) + assert.Equal(t, tt.expectedEvents, result["data"]) + } + } + }) + } +} + +func TestHandleResourceReq_Success(t *testing.T) { + mockEvents := []GraphiteEventsResponse{{When: 1234567890, What: "event1"}} + mockResp, _ := json.Marshal(mockEvents) + + dsInfo := datasourceInfo{ + Id: 1, + URL: "http://example.com", + HTTPClient: &http.Client{Transport: &mockRoundTripper{respBody: mockResp, status: 200}}, + } + + svc := &Service{ + logger: log.NewNullLogger(), + im: &mockInstanceManager{instance: dsInfo}, + } + + request := GraphiteEventsRequest{From: "now-1h", Until: "now"} + requestBody, _ := json.Marshal(request) + + req := httptest.NewRequest("POST", "/events", bytes.NewBuffer(requestBody)) + req = req.WithContext(backend.WithPluginContext(context.Background(), backend.PluginContext{})) + rr := httptest.NewRecorder() + + handler := svc.handleResourceReq(svc.handleEvents) + handler(rr, req) + + assert.Equal(t, http.StatusOK, rr.Code) + + var result map[string][]GraphiteEventsResponse + require.NoError(t, json.Unmarshal(rr.Body.Bytes(), &result)) + assert.Equal(t, mockEvents, result["data"]) +} + +func TestHandleResourceReq_GetDSInfoError(t *testing.T) { + svc := &Service{ + logger: log.NewNullLogger(), + im: &mockInstanceManager{err: errors.New("datasource not found")}, + } + + req := httptest.NewRequest("POST", "/events", bytes.NewBufferString("{}")) + req = req.WithContext(backend.WithPluginContext(context.Background(), backend.PluginContext{})) + rr := httptest.NewRecorder() + + handler := svc.handleResourceReq(svc.handleEvents) + handler(rr, req) + + assert.Equal(t, http.StatusInternalServerError, rr.Code) + + var errorResp map[string]string + require.NoError(t, json.Unmarshal(rr.Body.Bytes(), &errorResp)) + assert.Contains(t, errorResp["error"], "unexpected error") +} + +func TestHandleResourceReq_NilHandler(t *testing.T) { + dsInfo := datasourceInfo{Id: 1, URL: "http://example.com"} + + svc := &Service{ + logger: log.NewNullLogger(), + im: &mockInstanceManager{instance: dsInfo}, + } + + req := httptest.NewRequest("POST", "/events", bytes.NewBufferString("{}")) + req = req.WithContext(backend.WithPluginContext(context.Background(), backend.PluginContext{})) + rr := httptest.NewRecorder() + + handler := svc.handleResourceReq(nil) + handler(rr, req) + + assert.Equal(t, http.StatusInternalServerError, rr.Code) + + var errorResp map[string]string + require.NoError(t, json.Unmarshal(rr.Body.Bytes(), &errorResp)) + assert.Equal(t, "responseFn should not be nil", errorResp["error"]) +} + +func TestWriteErrorResponse(t *testing.T) { + rr := httptest.NewRecorder() + writeErrorResponse(rr, http.StatusBadRequest, "test error message") + + assert.Equal(t, http.StatusBadRequest, rr.Code) + + var errorResp map[string]string + require.NoError(t, json.Unmarshal(rr.Body.Bytes(), &errorResp)) + assert.Equal(t, "test error message", errorResp["error"]) +} diff --git a/pkg/tsdb/graphite/types.go b/pkg/tsdb/graphite/types.go index 40e47d5d1b0..9265b6c4855 100644 --- a/pkg/tsdb/graphite/types.go +++ b/pkg/tsdb/graphite/types.go @@ -18,3 +18,16 @@ type GraphiteQuery struct { Tags []string `json:"tags,omitempty"` FromAnnotations *bool `json:"fromAnnotations,omitempty"` } + +type GraphiteEventsRequest struct { + Tags string `json:"tags,omitempty"` + From string `json:"from"` + Until string `json:"until"` +} + +type GraphiteEventsResponse struct { + When int64 `json:"when"` + What string `json:"what"` + Tags []string `json:"tags"` + Data string `json:"data"` +} diff --git a/pkg/tsdb/graphite/utils.go b/pkg/tsdb/graphite/utils.go index 7a3bba9035b..4fb22cc6803 100644 --- a/pkg/tsdb/graphite/utils.go +++ b/pkg/tsdb/graphite/utils.go @@ -1 +1,47 @@ package graphite + +import ( + "compress/flate" + "compress/gzip" + "fmt" + "io" + + "github.com/andybalholm/brotli" + "github.com/grafana/grafana-plugin-sdk-go/backend" +) + +func decode(encoding string, original io.ReadCloser) ([]byte, error) { + var reader io.Reader + var err error + switch encoding { + case "gzip": + reader, err = gzip.NewReader(original) + if err != nil { + return nil, err + } + defer func() { + if err := reader.(io.ReadCloser).Close(); err != nil { + backend.Logger.Warn("Failed to close reader body", "err", err) + } + }() + case "deflate": + reader = flate.NewReader(original) + defer func() { + if err := reader.(io.ReadCloser).Close(); err != nil { + backend.Logger.Warn("Failed to close reader body", "err", err) + } + }() + case "br": + reader = brotli.NewReader(original) + case "": + reader = original + default: + return nil, fmt.Errorf("unexpected encoding type %v", err) + } + + body, err := io.ReadAll(reader) + if err != nil { + return nil, err + } + return body, nil +} diff --git a/public/app/plugins/datasource/graphite/datasource.ts b/public/app/plugins/datasource/graphite/datasource.ts index 75bc5c166ad..8d1df35715d 100644 --- a/public/app/plugins/datasource/graphite/datasource.ts +++ b/public/app/plugins/datasource/graphite/datasource.ts @@ -41,6 +41,7 @@ import { getRollupNotice, getRuntimeConsolidationNotice } from './meta'; import { prepareAnnotation } from './migrations'; // Types import { + GraphiteEvents, GraphiteLokiMapping, GraphiteMetricLokiMatcher, GraphiteOptions, @@ -457,7 +458,7 @@ export class GraphiteDatasource return this.events({ range: range, tags: tags }).then((results) => { const list = []; if (!isArray(results.data)) { - console.error(`Unable to get annotations from ${results.url}.`); + console.error(`Unable to get annotations.`); return []; } for (let i = 0; i < results.data.length; i++) { @@ -482,23 +483,30 @@ export class GraphiteDatasource } } - events(options: { range: TimeRange; tags: string; timezone?: TimeZone }) { + async events(options: { + range: TimeRange; + tags: string; + timezone?: TimeZone; + }): Promise<{ data: GraphiteEvents[] } | FetchResponse> { try { - let tags = ''; - if (options.tags) { - tags = '&tags=' + options.tags; + const tags = options.tags || ''; + const from = this.translateTime(options.range.raw.from, false, options.timezone); + const until = this.translateTime(options.range.raw.to, true, options.timezone); + if (config.featureToggles.graphiteBackendMode) { + return await this.postResource<{ data: GraphiteEvents[] }>('events', { + from: typeof from === 'string' ? from : `${from}`, + until: typeof until === 'string' ? until : `${until}`, + tags, + }); + } else { + const tagsQueryParam = tags === '' ? '' : `&tags=${tags}`; + return lastValueFrom( + this.doGraphiteRequest({ + method: 'GET', + url: `/events/get_data?from=${from}&until=${until}${tagsQueryParam}`, + }) + ); } - return lastValueFrom( - this.doGraphiteRequest({ - method: 'GET', - url: - '/events/get_data?from=' + - this.translateTime(options.range.raw.from, false, options.timezone) + - '&until=' + - this.translateTime(options.range.raw.to, true, options.timezone) + - tags, - }) - ); } catch (err) { return Promise.reject(err); } @@ -985,7 +993,7 @@ export class GraphiteDatasource return lastValueFrom(this.query(query)).then(() => ({ status: 'success', message: 'Data source is working' })); } - doGraphiteRequest( + doGraphiteRequest( options: BackendSrvRequest & { inspect?: any; } @@ -1002,7 +1010,7 @@ export class GraphiteDatasource options.inspect = { type: 'graphite' }; return getBackendSrv() - .fetch(options) + .fetch(options) .pipe( catchError((err) => { return throwError(() => { diff --git a/public/app/plugins/datasource/graphite/types.ts b/public/app/plugins/datasource/graphite/types.ts index 75471aaf0ef..338f49f5922 100644 --- a/public/app/plugins/datasource/graphite/types.ts +++ b/public/app/plugins/datasource/graphite/types.ts @@ -102,3 +102,16 @@ export type GraphiteQueryEditorDependencies = { export interface GraphiteQueryRequest extends DataQueryRequest { format: string; } + +export interface GraphiteEventsRequest { + from: number; + until: number; + tags: string; +} + +export interface GraphiteEvents { + when: number; + what: string; + tags: string[]; + data: string; +}