Plugins: Refactor Plugin Management (#40477)
* add core plugin flow * add instrumentation * move func * remove cruft * support external backend plugins * refactor + clean up * remove comments * refactor loader * simplify core plugin path arg * cleanup loggers * move signature validator to plugins package * fix sig packaging * cleanup plugin model * remove unnecessary plugin field * add start+stop for pm * fix failures * add decommissioned state * export fields just to get things flowing * fix comments * set static routes * make image loading idempotent * merge with backend plugin manager * re-use funcs * reorder imports + remove unnecessary interface * add some TODOs + remove unused func * remove unused instrumentation func * simplify client usage * remove import alias * re-use backendplugin.Plugin interface * re order funcs * improve var name * fix log statements * refactor data model * add logic for dupe check during loading * cleanup state setting * refactor loader * cleanup manager interface * add rendering flow * refactor loading + init * add renderer support * fix renderer plugin * reformat imports * track errors * fix plugin signature inheritance * name param in interface * update func comment * fix func arg name * introduce class concept * remove func * fix external plugin check * apply changes from pm-experiment * fix core plugins * fix imports * rename interface * comment API interface * add support for testdata plugin * enable alerting + use correct core plugin contracts * slim manager API * fix param name * fix filter * support static routes * fix rendering * tidy rendering * get tests compiling * fix install+uninstall * start finder test * add finder test coverage * start loader tests * add test for core plugins * load core + bundled test * add test for nested plugin loading * add test files * clean interface + fix registering some core plugins * refactoring * reformat and create sub packages * simplify core plugin init * fix ctx cancel scenario * migrate initializer * remove Init() funcs * add test starter * new logger * flesh out initializer tests * refactoring * remove unused svc * refactor rendering flow * fixup loader tests * add enabled helper func * fix logger name * fix data fetchers * fix case where plugin dir doesn't exist * improve coverage + move dupe checking to loader * remove noisy debug logs * register core plugins automagically * add support for renderer in catalog * make private func + fix req validation * use interface * re-add check for renderer in catalog * tidy up from moving to auto reg core plugins * core plugin registrar * guards * copy over core plugins for test infra * all tests green * renames * propagate new interfaces * kill old manager * get compiling * tidy up * update naming * refactor manager test + cleanup * add more cases to finder test * migrate validator to field * more coverage * refactor dupe checking * add test for plugin class * add coverage for initializer * split out rendering * move * fixup tests * fix uss test * fix frontend settings * fix grafanads test * add check when checking sig errors * fix enabled map * fixup * allow manual setup of CM * rename to cloud-monitoring * remove TODO * add installer interface for testing * loader interface returns * tests passing * refactor + add more coverage * support 'stackdriver' * fix frontend settings loading * improve naming based on package name * small tidy * refactor test * fix renderer start * make cloud-monitoring plugin ID clearer * add plugin update test * add integration tests * don't break all if sig can't be calculated * add root URL check test * add more signature verification tests * update DTO name * update enabled plugins comment * update comments * fix linter * revert fe naming change * fix errors endpoint * reset error code field name * re-order test to help verify * assert -> require * pm check * add missing entry + re-order * re-check * dump icon log * verify manager contents first * reformat * apply PR feedback * apply style changes * fix one vs all loading err * improve log output * only start when no signature error * move log * rework plugin update check * fix test * fix multi loading from cfg.PluginSettings * improve log output #2 * add error abstraction to capture errors without registering a plugin * add debug log * add unsigned warning * e2e test attempt * fix logger * set home path * prevent panic * alternate * ugh.. fix home path * return renderer even if not started * make renderer plugin managed * add fallback renderer icon, update renderer badge + prevent changes when renderer is installed * fix icon loading * rollback renderer changes * use correct field * remove unneccessary block * remove newline * remove unused func * fix bundled plugins base + module fields * remove unused field since refactor * add authorizer abstraction * loader only returns plugins expected to run * fix multi log output
This commit is contained in:
@@ -18,7 +18,7 @@ import (
|
||||
"github.com/grafana/grafana/pkg/components/simplejson"
|
||||
)
|
||||
|
||||
func (p *TestDataPlugin) handleCsvContentScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
func (s *Service) handleCsvContentScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
resp := backend.NewQueryDataResponse()
|
||||
|
||||
for _, q := range req.Queries {
|
||||
@@ -43,7 +43,7 @@ func (p *TestDataPlugin) handleCsvContentScenario(ctx context.Context, req *back
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) handleCsvFileScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
func (s *Service) handleCsvFileScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
resp := backend.NewQueryDataResponse()
|
||||
|
||||
for _, q := range req.Queries {
|
||||
@@ -58,7 +58,7 @@ func (p *TestDataPlugin) handleCsvFileScenario(ctx context.Context, req *backend
|
||||
continue
|
||||
}
|
||||
|
||||
frame, err := p.loadCsvFile(fileName)
|
||||
frame, err := s.loadCsvFile(fileName)
|
||||
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -72,14 +72,14 @@ func (p *TestDataPlugin) handleCsvFileScenario(ctx context.Context, req *backend
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) loadCsvFile(fileName string) (*data.Frame, error) {
|
||||
func (s *Service) loadCsvFile(fileName string) (*data.Frame, error) {
|
||||
validFileName := regexp.MustCompile(`([\w_]+)\.csv`)
|
||||
|
||||
if !validFileName.MatchString(fileName) {
|
||||
return nil, fmt.Errorf("invalid csv file name: %q", fileName)
|
||||
}
|
||||
|
||||
filePath := filepath.Join(p.cfg.StaticRootPath, "testdata", fileName)
|
||||
filePath := filepath.Join(s.cfg.StaticRootPath, "testdata", fileName)
|
||||
|
||||
// Can ignore gosec G304 here, because we check the file pattern above
|
||||
// nolint:gosec
|
||||
@@ -90,7 +90,7 @@ func (p *TestDataPlugin) loadCsvFile(fileName string) (*data.Frame, error) {
|
||||
|
||||
defer func() {
|
||||
if err := fileReader.Close(); err != nil {
|
||||
p.logger.Warn("Failed to close file", "err", err, "path", fileName)
|
||||
s.logger.Warn("Failed to close file", "err", err, "path", fileName)
|
||||
}
|
||||
}()
|
||||
|
||||
|
||||
@@ -17,7 +17,7 @@ func TestCSVFileScenario(t *testing.T) {
|
||||
cfg.DataPath = t.TempDir()
|
||||
cfg.StaticRootPath = "../../../public"
|
||||
|
||||
p := &TestDataPlugin{
|
||||
s := &Service{
|
||||
cfg: cfg,
|
||||
}
|
||||
|
||||
@@ -50,7 +50,7 @@ func TestCSVFileScenario(t *testing.T) {
|
||||
}
|
||||
|
||||
t.Run("Should not allow non file name chars", func(t *testing.T) {
|
||||
_, err := p.loadCsvFile("../population_by_state.csv")
|
||||
_, err := s.loadCsvFile("../population_by_state.csv")
|
||||
require.Error(t, err)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -11,7 +11,7 @@ import (
|
||||
"github.com/grafana/grafana/pkg/components/simplejson"
|
||||
)
|
||||
|
||||
func (p *TestDataPlugin) handleFlightPathScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
func (s *Service) handleFlightPathScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
resp := backend.NewQueryDataResponse()
|
||||
|
||||
for _, q := range req.Queries {
|
||||
|
||||
@@ -16,7 +16,7 @@ import (
|
||||
|
||||
func TestFlightPathScenario(t *testing.T) {
|
||||
cfg := setting.NewCfg()
|
||||
p := &TestDataPlugin{
|
||||
s := &Service{
|
||||
cfg: cfg,
|
||||
}
|
||||
|
||||
@@ -37,7 +37,7 @@ func TestFlightPathScenario(t *testing.T) {
|
||||
},
|
||||
}
|
||||
|
||||
rsp, err := p.handleFlightPathScenario(context.Background(), qr)
|
||||
rsp, err := s.handleFlightPathScenario(context.Background(), qr)
|
||||
require.NoError(t, err)
|
||||
require.NotNil(t, rsp)
|
||||
for k, v := range rsp.Responses {
|
||||
|
||||
@@ -15,40 +15,40 @@ import (
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend/resource/httpadapter"
|
||||
)
|
||||
|
||||
func (p *TestDataPlugin) registerRoutes(mux *http.ServeMux) {
|
||||
mux.HandleFunc("/", p.testGetHandler)
|
||||
mux.HandleFunc("/scenarios", p.getScenariosHandler)
|
||||
mux.HandleFunc("/stream", p.testStreamHandler)
|
||||
mux.Handle("/test", createJSONHandler(p.logger))
|
||||
mux.Handle("/test/json", createJSONHandler(p.logger))
|
||||
mux.HandleFunc("/boom", p.testPanicHandler)
|
||||
func (s *Service) RegisterRoutes(mux *http.ServeMux) {
|
||||
mux.HandleFunc("/", s.testGetHandler)
|
||||
mux.HandleFunc("/scenarios", s.getScenariosHandler)
|
||||
mux.HandleFunc("/stream", s.testStreamHandler)
|
||||
mux.Handle("/test", createJSONHandler(s.logger))
|
||||
mux.Handle("/test/json", createJSONHandler(s.logger))
|
||||
mux.HandleFunc("/boom", s.testPanicHandler)
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) testGetHandler(rw http.ResponseWriter, req *http.Request) {
|
||||
p.logger.Debug("Received resource call", "url", req.URL.String(), "method", req.Method)
|
||||
func (s *Service) testGetHandler(rw http.ResponseWriter, req *http.Request) {
|
||||
s.logger.Debug("Received resource call", "url", req.URL.String(), "method", req.Method)
|
||||
|
||||
if req.Method != http.MethodGet {
|
||||
return
|
||||
}
|
||||
|
||||
if _, err := rw.Write([]byte("Hello world from test datasource!")); err != nil {
|
||||
p.logger.Error("Failed to write response", "error", err)
|
||||
s.logger.Error("Failed to write response", "error", err)
|
||||
return
|
||||
}
|
||||
rw.WriteHeader(http.StatusOK)
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) getScenariosHandler(rw http.ResponseWriter, req *http.Request) {
|
||||
func (s *Service) getScenariosHandler(rw http.ResponseWriter, req *http.Request) {
|
||||
result := make([]interface{}, 0)
|
||||
|
||||
scenarioIds := make([]string, 0)
|
||||
for id := range p.scenarios {
|
||||
for id := range s.scenarios {
|
||||
scenarioIds = append(scenarioIds, id)
|
||||
}
|
||||
sort.Strings(scenarioIds)
|
||||
|
||||
for _, scenarioID := range scenarioIds {
|
||||
scenario := p.scenarios[scenarioID]
|
||||
scenario := s.scenarios[scenarioID]
|
||||
result = append(result, map[string]interface{}{
|
||||
"id": scenario.ID,
|
||||
"name": scenario.Name,
|
||||
@@ -59,18 +59,18 @@ func (p *TestDataPlugin) getScenariosHandler(rw http.ResponseWriter, req *http.R
|
||||
|
||||
bytes, err := json.Marshal(&result)
|
||||
if err != nil {
|
||||
p.logger.Error("Failed to marshal response body to JSON", "error", err)
|
||||
s.logger.Error("Failed to marshal response body to JSON", "error", err)
|
||||
}
|
||||
|
||||
rw.Header().Set("Content-Type", "application/json")
|
||||
rw.WriteHeader(http.StatusOK)
|
||||
if _, err := rw.Write(bytes); err != nil {
|
||||
p.logger.Error("Failed to write response", "error", err)
|
||||
s.logger.Error("Failed to write response", "error", err)
|
||||
}
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) testStreamHandler(rw http.ResponseWriter, req *http.Request) {
|
||||
p.logger.Debug("Received resource call", "url", req.URL.String(), "method", req.Method)
|
||||
func (s *Service) testStreamHandler(rw http.ResponseWriter, req *http.Request) {
|
||||
s.logger.Debug("Received resource call", "url", req.URL.String(), "method", req.Method)
|
||||
|
||||
if req.Method != http.MethodGet {
|
||||
return
|
||||
@@ -95,7 +95,7 @@ func (p *TestDataPlugin) testStreamHandler(rw http.ResponseWriter, req *http.Req
|
||||
|
||||
for i := 1; i <= count; i++ {
|
||||
if _, err := io.WriteString(rw, fmt.Sprintf("Message #%d", i)); err != nil {
|
||||
p.logger.Error("Failed to write response", "error", err)
|
||||
s.logger.Error("Failed to write response", "error", err)
|
||||
return
|
||||
}
|
||||
rw.(http.Flusher).Flush()
|
||||
@@ -152,6 +152,6 @@ func createJSONHandler(logger log.Logger) http.Handler {
|
||||
})
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) testPanicHandler(rw http.ResponseWriter, req *http.Request) {
|
||||
func (s *Service) testPanicHandler(rw http.ResponseWriter, req *http.Request) {
|
||||
panic("BOOM")
|
||||
}
|
||||
|
||||
@@ -54,197 +54,197 @@ type Scenario struct {
|
||||
handler backend.QueryDataHandlerFunc
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) registerScenario(scenario *Scenario) {
|
||||
p.scenarios[scenario.ID] = scenario
|
||||
p.queryMux.HandleFunc(scenario.ID, scenario.handler)
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) registerScenarios() {
|
||||
p.registerScenario(&Scenario{
|
||||
func (s *Service) registerScenarios() {
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(exponentialHeatmapBucketDataQuery),
|
||||
Name: "Exponential heatmap bucket data",
|
||||
handler: p.handleExponentialHeatmapBucketDataScenario,
|
||||
handler: s.handleExponentialHeatmapBucketDataScenario,
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(linearHeatmapBucketDataQuery),
|
||||
Name: "Linear heatmap bucket data",
|
||||
handler: p.handleLinearHeatmapBucketDataScenario,
|
||||
handler: s.handleLinearHeatmapBucketDataScenario,
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(randomWalkQuery),
|
||||
Name: "Random Walk",
|
||||
handler: p.handleRandomWalkScenario,
|
||||
handler: s.handleRandomWalkScenario,
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(predictablePulseQuery),
|
||||
Name: "Predictable Pulse",
|
||||
handler: p.handlePredictablePulseScenario,
|
||||
handler: s.handlePredictablePulseScenario,
|
||||
Description: `Predictable Pulse returns a pulse wave where there is a datapoint every timeStepSeconds.
|
||||
The wave cycles at timeStepSeconds*(onCount+offCount).
|
||||
The cycle of the wave is based off of absolute time (from the epoch) which makes it predictable.
|
||||
Timestamps will line up evenly on timeStepSeconds (For example, 60 seconds means times will all end in :00 seconds).`,
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(predictableCSVWaveQuery),
|
||||
Name: "Predictable CSV Wave",
|
||||
handler: p.handlePredictableCSVWaveScenario,
|
||||
handler: s.handlePredictableCSVWaveScenario,
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(randomWalkTableQuery),
|
||||
Name: "Random Walk Table",
|
||||
handler: p.handleRandomWalkTableScenario,
|
||||
handler: s.handleRandomWalkTableScenario,
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(randomWalkSlowQuery),
|
||||
Name: "Slow Query",
|
||||
StringInput: "5s",
|
||||
handler: p.handleRandomWalkSlowScenario,
|
||||
handler: s.handleRandomWalkSlowScenario,
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(noDataPointsQuery),
|
||||
Name: "No Data Points",
|
||||
handler: p.handleClientSideScenario,
|
||||
handler: s.handleClientSideScenario,
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(datapointsOutsideRangeQuery),
|
||||
Name: "Datapoints Outside Range",
|
||||
handler: p.handleDatapointsOutsideRangeScenario,
|
||||
handler: s.handleDatapointsOutsideRangeScenario,
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(csvMetricValuesQuery),
|
||||
Name: "CSV Metric Values",
|
||||
StringInput: "1,20,90,30,5,0",
|
||||
handler: p.handleCSVMetricValuesScenario,
|
||||
handler: s.handleCSVMetricValuesScenario,
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(streamingClientQuery),
|
||||
Name: "Streaming Client",
|
||||
handler: p.handleClientSideScenario,
|
||||
handler: s.handleClientSideScenario,
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(liveQuery),
|
||||
Name: "Grafana Live",
|
||||
handler: p.handleClientSideScenario,
|
||||
handler: s.handleClientSideScenario,
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(flightPath),
|
||||
Name: "Flight path",
|
||||
handler: p.handleFlightPathScenario,
|
||||
handler: s.handleFlightPathScenario,
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(usaQueryKey),
|
||||
Name: "USA generated data",
|
||||
handler: p.handleUSAScenario,
|
||||
handler: s.handleUSAScenario,
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(grafanaAPIQuery),
|
||||
Name: "Grafana API",
|
||||
handler: p.handleClientSideScenario,
|
||||
handler: s.handleClientSideScenario,
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(arrowQuery),
|
||||
Name: "Load Apache Arrow Data",
|
||||
handler: p.handleArrowScenario,
|
||||
handler: s.handleArrowScenario,
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(annotationsQuery),
|
||||
Name: "Annotations",
|
||||
handler: p.handleClientSideScenario,
|
||||
handler: s.handleClientSideScenario,
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(tableStaticQuery),
|
||||
Name: "Table Static",
|
||||
handler: p.handleTableStaticScenario,
|
||||
handler: s.handleTableStaticScenario,
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(randomWalkWithErrorQuery),
|
||||
Name: "Random Walk (with error)",
|
||||
handler: p.handleRandomWalkWithErrorScenario,
|
||||
handler: s.handleRandomWalkWithErrorScenario,
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(serverError500Query),
|
||||
Name: "Server Error (500)",
|
||||
handler: p.handleServerError500Scenario,
|
||||
handler: s.handleServerError500Scenario,
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(logsQuery),
|
||||
Name: "Logs",
|
||||
handler: p.handleLogsScenario,
|
||||
handler: s.handleLogsScenario,
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(nodeGraphQuery),
|
||||
Name: "Node Graph",
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(csvFileQueryType),
|
||||
Name: "CSV File",
|
||||
handler: p.handleCsvFileScenario,
|
||||
handler: s.handleCsvFileScenario,
|
||||
})
|
||||
|
||||
p.registerScenario(&Scenario{
|
||||
s.registerScenario(&Scenario{
|
||||
ID: string(csvContentQueryType),
|
||||
Name: "CSV Content",
|
||||
handler: p.handleCsvContentScenario,
|
||||
handler: s.handleCsvContentScenario,
|
||||
})
|
||||
|
||||
p.queryMux.HandleFunc("", p.handleFallbackScenario)
|
||||
s.queryMux.HandleFunc("", s.handleFallbackScenario)
|
||||
}
|
||||
|
||||
func (s *Service) registerScenario(scenario *Scenario) {
|
||||
s.scenarios[scenario.ID] = scenario
|
||||
s.queryMux.HandleFunc(scenario.ID, scenario.handler)
|
||||
}
|
||||
|
||||
// handleFallbackScenario handles the scenario where queryType is not set and fallbacks to scenarioId.
|
||||
func (p *TestDataPlugin) handleFallbackScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
func (s *Service) handleFallbackScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
scenarioQueries := map[string][]backend.DataQuery{}
|
||||
|
||||
for _, q := range req.Queries {
|
||||
model, err := simplejson.NewJson(q.JSON)
|
||||
if err != nil {
|
||||
p.logger.Error("Failed to unmarshal query model to JSON", "error", err)
|
||||
s.logger.Error("Failed to unmarshal query model to JSON", "error", err)
|
||||
continue
|
||||
}
|
||||
|
||||
scenarioID := model.Get("scenarioId").MustString(string(randomWalkQuery))
|
||||
if _, exist := p.scenarios[scenarioID]; exist {
|
||||
if _, exist := s.scenarios[scenarioID]; exist {
|
||||
if _, ok := scenarioQueries[scenarioID]; !ok {
|
||||
scenarioQueries[scenarioID] = []backend.DataQuery{}
|
||||
}
|
||||
|
||||
scenarioQueries[scenarioID] = append(scenarioQueries[scenarioID], q)
|
||||
} else {
|
||||
p.logger.Error("Scenario not found", "scenarioId", scenarioID)
|
||||
s.logger.Error("Scenario not found", "scenarioId", scenarioID)
|
||||
}
|
||||
}
|
||||
|
||||
resp := backend.NewQueryDataResponse()
|
||||
for scenarioID, queries := range scenarioQueries {
|
||||
if scenario, exist := p.scenarios[scenarioID]; exist {
|
||||
if scenario, exist := s.scenarios[scenarioID]; exist {
|
||||
sReq := &backend.QueryDataRequest{
|
||||
PluginContext: req.PluginContext,
|
||||
Headers: req.Headers,
|
||||
Queries: queries,
|
||||
}
|
||||
if sResp, err := scenario.handler(ctx, sReq); err != nil {
|
||||
p.logger.Error("Failed to handle scenario", "scenarioId", scenarioID, "error", err)
|
||||
s.logger.Error("Failed to handle scenario", "scenarioId", scenarioID, "error", err)
|
||||
} else {
|
||||
for refID, dr := range sResp.Responses {
|
||||
resp.Responses[refID] = dr
|
||||
@@ -256,7 +256,7 @@ func (p *TestDataPlugin) handleFallbackScenario(ctx context.Context, req *backen
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) handleRandomWalkScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
func (s *Service) handleRandomWalkScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
resp := backend.NewQueryDataResponse()
|
||||
|
||||
for _, q := range req.Queries {
|
||||
@@ -276,7 +276,7 @@ func (p *TestDataPlugin) handleRandomWalkScenario(ctx context.Context, req *back
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) handleDatapointsOutsideRangeScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
func (s *Service) handleDatapointsOutsideRangeScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
resp := backend.NewQueryDataResponse()
|
||||
|
||||
for _, q := range req.Queries {
|
||||
@@ -300,7 +300,7 @@ func (p *TestDataPlugin) handleDatapointsOutsideRangeScenario(ctx context.Contex
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) handleCSVMetricValuesScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
func (s *Service) handleCSVMetricValuesScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
resp := backend.NewQueryDataResponse()
|
||||
|
||||
for _, q := range req.Queries {
|
||||
@@ -344,7 +344,7 @@ func (p *TestDataPlugin) handleCSVMetricValuesScenario(ctx context.Context, req
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) handleRandomWalkWithErrorScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
func (s *Service) handleRandomWalkWithErrorScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
resp := backend.NewQueryDataResponse()
|
||||
|
||||
for _, q := range req.Queries {
|
||||
@@ -362,7 +362,7 @@ func (p *TestDataPlugin) handleRandomWalkWithErrorScenario(ctx context.Context,
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) handleRandomWalkSlowScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
func (s *Service) handleRandomWalkSlowScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
resp := backend.NewQueryDataResponse()
|
||||
|
||||
for _, q := range req.Queries {
|
||||
@@ -383,7 +383,7 @@ func (p *TestDataPlugin) handleRandomWalkSlowScenario(ctx context.Context, req *
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) handleRandomWalkTableScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
func (s *Service) handleRandomWalkTableScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
resp := backend.NewQueryDataResponse()
|
||||
|
||||
for _, q := range req.Queries {
|
||||
@@ -400,7 +400,7 @@ func (p *TestDataPlugin) handleRandomWalkTableScenario(ctx context.Context, req
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) handlePredictableCSVWaveScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
func (s *Service) handlePredictableCSVWaveScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
resp := backend.NewQueryDataResponse()
|
||||
|
||||
for _, q := range req.Queries {
|
||||
@@ -421,7 +421,7 @@ func (p *TestDataPlugin) handlePredictableCSVWaveScenario(ctx context.Context, r
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) handlePredictablePulseScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
func (s *Service) handlePredictablePulseScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
resp := backend.NewQueryDataResponse()
|
||||
|
||||
for _, q := range req.Queries {
|
||||
@@ -442,15 +442,15 @@ func (p *TestDataPlugin) handlePredictablePulseScenario(ctx context.Context, req
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) handleServerError500Scenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
func (s *Service) handleServerError500Scenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
panic("Test Data Panic!")
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) handleClientSideScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
func (s *Service) handleClientSideScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
return backend.NewQueryDataResponse(), nil
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) handleArrowScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
func (s *Service) handleArrowScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
resp := backend.NewQueryDataResponse()
|
||||
|
||||
for _, q := range req.Queries {
|
||||
@@ -474,7 +474,7 @@ func (p *TestDataPlugin) handleArrowScenario(ctx context.Context, req *backend.Q
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) handleExponentialHeatmapBucketDataScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
func (s *Service) handleExponentialHeatmapBucketDataScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
resp := backend.NewQueryDataResponse()
|
||||
|
||||
for _, q := range req.Queries {
|
||||
@@ -489,7 +489,7 @@ func (p *TestDataPlugin) handleExponentialHeatmapBucketDataScenario(ctx context.
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) handleLinearHeatmapBucketDataScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
func (s *Service) handleLinearHeatmapBucketDataScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
resp := backend.NewQueryDataResponse()
|
||||
|
||||
for _, q := range req.Queries {
|
||||
@@ -504,7 +504,7 @@ func (p *TestDataPlugin) handleLinearHeatmapBucketDataScenario(ctx context.Conte
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) handleTableStaticScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
func (s *Service) handleTableStaticScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
resp := backend.NewQueryDataResponse()
|
||||
|
||||
for _, q := range req.Queries {
|
||||
@@ -533,7 +533,7 @@ func (p *TestDataPlugin) handleTableStaticScenario(ctx context.Context, req *bac
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) handleLogsScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
func (s *Service) handleLogsScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
resp := backend.NewQueryDataResponse()
|
||||
|
||||
for _, q := range req.Queries {
|
||||
|
||||
@@ -15,7 +15,7 @@ import (
|
||||
)
|
||||
|
||||
func TestTestdataScenarios(t *testing.T) {
|
||||
p := &TestDataPlugin{}
|
||||
s := &Service{}
|
||||
|
||||
t.Run("random walk ", func(t *testing.T) {
|
||||
t.Run("Should start at the requested value", func(t *testing.T) {
|
||||
@@ -42,7 +42,7 @@ func TestTestdataScenarios(t *testing.T) {
|
||||
Queries: []backend.DataQuery{query},
|
||||
}
|
||||
|
||||
resp, err := p.handleRandomWalkScenario(context.Background(), req)
|
||||
resp, err := s.handleRandomWalkScenario(context.Background(), req)
|
||||
require.NoError(t, err)
|
||||
require.NotNil(t, resp)
|
||||
|
||||
@@ -85,7 +85,7 @@ func TestTestdataScenarios(t *testing.T) {
|
||||
Queries: []backend.DataQuery{query},
|
||||
}
|
||||
|
||||
resp, err := p.handleRandomWalkTableScenario(context.Background(), req)
|
||||
resp, err := s.handleRandomWalkTableScenario(context.Background(), req)
|
||||
require.NoError(t, err)
|
||||
require.NotNil(t, resp)
|
||||
|
||||
@@ -141,7 +141,7 @@ func TestTestdataScenarios(t *testing.T) {
|
||||
Queries: []backend.DataQuery{query},
|
||||
}
|
||||
|
||||
resp, err := p.handleRandomWalkTableScenario(context.Background(), req)
|
||||
resp, err := s.handleRandomWalkTableScenario(context.Background(), req)
|
||||
require.NoError(t, err)
|
||||
require.NotNil(t, resp)
|
||||
|
||||
|
||||
@@ -9,34 +9,11 @@ import (
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
"github.com/grafana/grafana-plugin-sdk-go/data"
|
||||
|
||||
"github.com/grafana/grafana/pkg/infra/log"
|
||||
)
|
||||
|
||||
type testStreamHandler struct {
|
||||
logger log.Logger
|
||||
frame *data.Frame
|
||||
// If Live Pipeline enabled we are sending the whole frame to have a chance to process stream with rules.
|
||||
livePipelineEnabled bool
|
||||
}
|
||||
|
||||
func newTestStreamHandler(logger log.Logger, livePipelineEnabled bool) *testStreamHandler {
|
||||
frame := data.NewFrame("testdata",
|
||||
data.NewField("Time", nil, make([]time.Time, 1)),
|
||||
data.NewField("Value", nil, make([]float64, 1)),
|
||||
data.NewField("Min", nil, make([]float64, 1)),
|
||||
data.NewField("Max", nil, make([]float64, 1)),
|
||||
)
|
||||
return &testStreamHandler{
|
||||
frame: frame,
|
||||
logger: logger,
|
||||
livePipelineEnabled: livePipelineEnabled,
|
||||
}
|
||||
}
|
||||
|
||||
func (p *testStreamHandler) SubscribeStream(_ context.Context, req *backend.SubscribeStreamRequest) (*backend.SubscribeStreamResponse, error) {
|
||||
p.logger.Debug("Allowing access to stream", "path", req.Path, "user", req.PluginContext.User)
|
||||
initialData, err := backend.NewInitialFrame(p.frame, data.IncludeSchemaOnly)
|
||||
func (s *Service) SubscribeStream(_ context.Context, req *backend.SubscribeStreamRequest) (*backend.SubscribeStreamResponse, error) {
|
||||
s.logger.Debug("Allowing access to stream", "path", req.Path, "user", req.PluginContext.User)
|
||||
initialData, err := backend.NewInitialFrame(s.frame, data.IncludeSchemaOnly)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -50,7 +27,7 @@ func (p *testStreamHandler) SubscribeStream(_ context.Context, req *backend.Subs
|
||||
}
|
||||
}
|
||||
|
||||
if p.livePipelineEnabled {
|
||||
if s.cfg.FeatureToggles["live-pipeline"] {
|
||||
// While developing Live pipeline avoid sending initial data.
|
||||
initialData = nil
|
||||
}
|
||||
@@ -61,15 +38,15 @@ func (p *testStreamHandler) SubscribeStream(_ context.Context, req *backend.Subs
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (p *testStreamHandler) PublishStream(_ context.Context, req *backend.PublishStreamRequest) (*backend.PublishStreamResponse, error) {
|
||||
p.logger.Debug("Attempt to publish into stream", "path", req.Path, "user", req.PluginContext.User)
|
||||
func (s *Service) PublishStream(_ context.Context, req *backend.PublishStreamRequest) (*backend.PublishStreamResponse, error) {
|
||||
s.logger.Debug("Attempt to publish into stream", "path", req.Path, "user", req.PluginContext.User)
|
||||
return &backend.PublishStreamResponse{
|
||||
Status: backend.PublishStreamStatusPermissionDenied,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (p *testStreamHandler) RunStream(ctx context.Context, request *backend.RunStreamRequest, sender *backend.StreamSender) error {
|
||||
p.logger.Debug("New stream call", "path", request.Path)
|
||||
func (s *Service) RunStream(ctx context.Context, request *backend.RunStreamRequest, sender *backend.StreamSender) error {
|
||||
s.logger.Debug("New stream call", "path", request.Path)
|
||||
var conf testStreamConfig
|
||||
switch request.Path {
|
||||
case "random-2s-stream":
|
||||
@@ -93,7 +70,7 @@ func (p *testStreamHandler) RunStream(ctx context.Context, request *backend.RunS
|
||||
default:
|
||||
return fmt.Errorf("testdata plugin does not support path: %s", request.Path)
|
||||
}
|
||||
return p.runTestStream(ctx, request.Path, conf, sender)
|
||||
return s.runTestStream(ctx, request.Path, conf, sender)
|
||||
}
|
||||
|
||||
type testStreamConfig struct {
|
||||
@@ -102,7 +79,7 @@ type testStreamConfig struct {
|
||||
Flight *flightConfig
|
||||
}
|
||||
|
||||
func (p *testStreamHandler) runTestStream(ctx context.Context, path string, conf testStreamConfig, sender *backend.StreamSender) error {
|
||||
func (s *Service) runTestStream(ctx context.Context, path string, conf testStreamConfig, sender *backend.StreamSender) error {
|
||||
spread := 50.0
|
||||
walker := rand.Float64() * 100
|
||||
|
||||
@@ -118,7 +95,7 @@ func (p *testStreamHandler) runTestStream(ctx context.Context, path string, conf
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
p.logger.Debug("Stop streaming data for path", "path", path)
|
||||
s.logger.Debug("Stop streaming data for path", "path", path)
|
||||
return ctx.Err()
|
||||
case t := <-ticker.C:
|
||||
if rand.Float64() < conf.Drop {
|
||||
@@ -126,7 +103,7 @@ func (p *testStreamHandler) runTestStream(ctx context.Context, path string, conf
|
||||
}
|
||||
|
||||
mode := data.IncludeDataOnly
|
||||
if p.livePipelineEnabled {
|
||||
if s.cfg.FeatureToggles["live-pipeline"] {
|
||||
mode = data.IncludeAll
|
||||
}
|
||||
|
||||
@@ -139,11 +116,11 @@ func (p *testStreamHandler) runTestStream(ctx context.Context, path string, conf
|
||||
delta := rand.Float64() - 0.5
|
||||
walker += delta
|
||||
|
||||
p.frame.Fields[0].Set(0, t)
|
||||
p.frame.Fields[1].Set(0, walker) // Value
|
||||
p.frame.Fields[2].Set(0, walker-((rand.Float64()*spread)+0.01)) // Min
|
||||
p.frame.Fields[3].Set(0, walker+((rand.Float64()*spread)+0.01)) // Max
|
||||
if err := sender.SendFrame(p.frame, mode); err != nil {
|
||||
s.frame.Fields[0].Set(0, t)
|
||||
s.frame.Fields[1].Set(0, walker) // Value
|
||||
s.frame.Fields[2].Set(0, walker-((rand.Float64()*spread)+0.01)) // Min
|
||||
s.frame.Fields[3].Set(0, walker+((rand.Float64()*spread)+0.01)) // Max
|
||||
if err := sender.SendFrame(s.frame, mode); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2,49 +2,54 @@ package testdatasource
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"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/resource/httpadapter"
|
||||
"github.com/grafana/grafana-plugin-sdk-go/data"
|
||||
"github.com/grafana/grafana/pkg/infra/log"
|
||||
"github.com/grafana/grafana/pkg/plugins/backendplugin"
|
||||
"github.com/grafana/grafana/pkg/plugins"
|
||||
"github.com/grafana/grafana/pkg/plugins/backendplugin/coreplugin"
|
||||
"github.com/grafana/grafana/pkg/setting"
|
||||
)
|
||||
|
||||
func ProvideService(cfg *setting.Cfg, manager backendplugin.Manager) (*TestDataPlugin, error) {
|
||||
resourceMux := http.NewServeMux()
|
||||
p := new(cfg, resourceMux)
|
||||
func ProvideService(cfg *setting.Cfg, registrar plugins.CoreBackendRegistrar) (*Service, error) {
|
||||
s := &Service{
|
||||
queryMux: datasource.NewQueryTypeMux(),
|
||||
scenarios: map[string]*Scenario{},
|
||||
frame: data.NewFrame("testdata",
|
||||
data.NewField("Time", nil, make([]time.Time, 1)),
|
||||
data.NewField("Value", nil, make([]float64, 1)),
|
||||
data.NewField("Min", nil, make([]float64, 1)),
|
||||
data.NewField("Max", nil, make([]float64, 1)),
|
||||
),
|
||||
logger: log.New("tsdb.testdata"),
|
||||
cfg: cfg,
|
||||
}
|
||||
|
||||
s.registerScenarios()
|
||||
|
||||
rMux := http.NewServeMux()
|
||||
s.RegisterRoutes(rMux)
|
||||
|
||||
factory := coreplugin.New(backend.ServeOpts{
|
||||
QueryDataHandler: p.queryMux,
|
||||
CallResourceHandler: httpadapter.New(resourceMux),
|
||||
StreamHandler: newTestStreamHandler(p.logger, cfg.FeatureToggles["live-pipeline"]),
|
||||
QueryDataHandler: s.queryMux,
|
||||
CallResourceHandler: httpadapter.New(rMux),
|
||||
StreamHandler: s,
|
||||
})
|
||||
err := manager.Register("testdata", factory)
|
||||
err := registrar.LoadAndRegister("testdata", factory)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return p, nil
|
||||
return s, nil
|
||||
}
|
||||
|
||||
func new(cfg *setting.Cfg, resourceMux *http.ServeMux) *TestDataPlugin {
|
||||
p := &TestDataPlugin{
|
||||
logger: log.New("tsdb.testdata"),
|
||||
cfg: cfg,
|
||||
scenarios: map[string]*Scenario{},
|
||||
queryMux: datasource.NewQueryTypeMux(),
|
||||
}
|
||||
|
||||
p.registerScenarios()
|
||||
p.registerRoutes(resourceMux)
|
||||
|
||||
return p
|
||||
}
|
||||
|
||||
type TestDataPlugin struct {
|
||||
type Service struct {
|
||||
cfg *setting.Cfg
|
||||
logger log.Logger
|
||||
scenarios map[string]*Scenario
|
||||
frame *data.Frame
|
||||
queryMux *datasource.QueryTypeMux
|
||||
}
|
||||
|
||||
@@ -36,7 +36,7 @@ type usaQuery struct {
|
||||
period time.Duration
|
||||
}
|
||||
|
||||
func (p *TestDataPlugin) handleUSAScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
func (s *Service) handleUSAScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
resp := backend.NewQueryDataResponse()
|
||||
|
||||
for _, q := range req.Queries {
|
||||
|
||||
@@ -17,7 +17,7 @@ import (
|
||||
|
||||
func TestUSAScenario(t *testing.T) {
|
||||
cfg := setting.NewCfg()
|
||||
p := &TestDataPlugin{
|
||||
p := &Service{
|
||||
cfg: cfg,
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user