Graphite: Backend events endpoint (#110598)

* Add lint rules

* Backend decoupling

- Add standalone files
- Add graphite query type
- Add logger to Service
- Create logger in the ProvideService method
- Use a pointer for the HTTP client provider
- Update logger usage everywhere
- Update tracer type
- Replace simplejson with json
- Add dummy CallResource and CheckHealth methods
- Update tests

* Update ConfigEditor imports

* Update types imports

* Update datasource

- Switch to using semver package
- Update imports

* Update store imports

* Update helper imports and notification creation

* Update context import

* Update version numbers and logic

* Copy array_move from core

* Test updates

* Add required files and update plugin.json

* Update core references and packages

* Remove commented code

* Update wire

* Lint

* Fix import

* Copy null type

* More lint

* Update snapshot

* Refactor backend

- Split query logic into separate file
- Move utils to separate file

* Add health-check logic

- Support backend healthcheck if the FF is enabled

* Remove query import support as unneeded

* Add test

* Add util function for decoding responses

* Add events types

* Add resource handler

* Add events handler and generic resource req handler

* Tests

* Update frontend

- Add types
- Update events function to support backend requests

* Lint and typing

* Lint

* Add tests

* Review

* Review

* Fix packages

* Fix merge issues
This commit is contained in:
Andreas Christou
2025-09-11 17:08:19 +01:00
committed by GitHub
parent aecc2c9fe7
commit 85e92ce04b
7 changed files with 532 additions and 23 deletions
+12 -5
View File
@@ -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)
}
+155
View File
@@ -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)
}
}
+267
View File
@@ -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"])
}
+13
View File
@@ -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"`
}
+46
View File
@@ -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
}
@@ -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<GraphiteEvents>> {
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<GraphiteEvents[]>({
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<T>(
options: BackendSrvRequest & {
inspect?: any;
}
@@ -1002,7 +1010,7 @@ export class GraphiteDatasource
options.inspect = { type: 'graphite' };
return getBackendSrv()
.fetch(options)
.fetch<T>(options)
.pipe(
catchError((err) => {
return throwError(() => {
@@ -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;
}