From 043d6cd58441f676470208a94d2c588391c66853 Mon Sep 17 00:00:00 2001 From: Marcus Efraimsson Date: Fri, 29 Jan 2021 18:33:23 +0100 Subject: [PATCH] Backend Plugins: Convert test data source to use SDK contracts (#29916) Converts the core testdata data source to use the SDK contracts and by that implementing a backend plugin in core Grafana in similar manner as an external one. Co-authored-by: Will Browne Co-authored-by: Marcus Efraimsson Co-authored-by: Ryan McKinley --- e2e/suite1/specs/panelEdit_queries.spec.ts | 2 +- .../specs/variables/load-options-from-url.ts | 6 +- .../specs/variables/set-options-from-ui.ts | 6 +- .../src/utils/DataSourceWithBackend.ts | 2 +- .../src/utils/queryResponse.test.ts | 216 ++- .../src/utils/queryResponse.ts | 115 +- pkg/api/api.go | 4 - pkg/api/metrics.go | 37 +- .../backendplugin/coreplugin/core_plugin.go | 3 +- .../coreplugin/core_plugin_test.go | 4 +- pkg/tsdb/testdatasource/resource_handler.go | 157 ++ pkg/tsdb/testdatasource/scenarios.go | 1444 ++++++++++------- pkg/tsdb/testdatasource/scenarios_test.go | 220 ++- pkg/tsdb/testdatasource/testdata.go | 58 +- .../plugins/datasource/testdata/datasource.ts | 77 +- 15 files changed, 1456 insertions(+), 895 deletions(-) create mode 100644 pkg/tsdb/testdatasource/resource_handler.go diff --git a/e2e/suite1/specs/panelEdit_queries.spec.ts b/e2e/suite1/specs/panelEdit_queries.spec.ts index cf9ae4d4482..861feeaba2b 100644 --- a/e2e/suite1/specs/panelEdit_queries.spec.ts +++ b/e2e/suite1/specs/panelEdit_queries.spec.ts @@ -29,7 +29,7 @@ e2e.scenario({ e2e() .route({ method: 'POST', - url: '/api/tsdb/query', + url: '/api/ds/query', }) .as('apiPostQuery'); diff --git a/e2e/suite1/specs/variables/load-options-from-url.ts b/e2e/suite1/specs/variables/load-options-from-url.ts index 35c2da718b0..172ccc23e9b 100644 --- a/e2e/suite1/specs/variables/load-options-from-url.ts +++ b/e2e/suite1/specs/variables/load-options-from-url.ts @@ -10,7 +10,7 @@ describe('Variables - Load options from Url', () => { e2e() .route({ method: 'POST', - url: '/api/tsdb/query', + url: '/api/ds/query', }) .as('query'); @@ -63,7 +63,7 @@ describe('Variables - Load options from Url', () => { e2e() .route({ method: 'POST', - url: '/api/tsdb/query', + url: '/api/ds/query', }) .as('query'); @@ -127,7 +127,7 @@ describe('Variables - Load options from Url', () => { e2e() .route({ method: 'POST', - url: '/api/tsdb/query', + url: '/api/ds/query', }) .as('query'); diff --git a/e2e/suite1/specs/variables/set-options-from-ui.ts b/e2e/suite1/specs/variables/set-options-from-ui.ts index 3f0e6b0aaef..4a22c959244 100644 --- a/e2e/suite1/specs/variables/set-options-from-ui.ts +++ b/e2e/suite1/specs/variables/set-options-from-ui.ts @@ -10,7 +10,7 @@ describe('Variables - Set options from ui', () => { e2e() .route({ method: 'POST', - url: '/api/tsdb/query', + url: '/api/ds/query', }) .as('query'); @@ -68,7 +68,7 @@ describe('Variables - Set options from ui', () => { e2e() .route({ method: 'POST', - url: '/api/tsdb/query', + url: '/api/ds/query', }) .as('query'); @@ -123,7 +123,7 @@ describe('Variables - Set options from ui', () => { e2e() .route({ method: 'POST', - url: '/api/tsdb/query', + url: '/api/ds/query', }) .as('query'); diff --git a/packages/grafana-runtime/src/utils/DataSourceWithBackend.ts b/packages/grafana-runtime/src/utils/DataSourceWithBackend.ts index a0fdce73e31..6c942f70c22 100644 --- a/packages/grafana-runtime/src/utils/DataSourceWithBackend.ts +++ b/packages/grafana-runtime/src/utils/DataSourceWithBackend.ts @@ -112,7 +112,7 @@ export class DataSourceWithBackend< }) .pipe( map((rsp: any) => { - return toDataQueryResponse(rsp); + return toDataQueryResponse(rsp, queries as DataQuery[]); }), catchError((err) => { return of(toDataQueryResponse(err)); diff --git a/packages/grafana-runtime/src/utils/queryResponse.test.ts b/packages/grafana-runtime/src/utils/queryResponse.test.ts index b47b7fe9a35..0adb915d2c5 100644 --- a/packages/grafana-runtime/src/utils/queryResponse.test.ts +++ b/packages/grafana-runtime/src/utils/queryResponse.test.ts @@ -1,16 +1,25 @@ -import { toDataFrameDTO } from '@grafana/data'; - +import { DataQuery, toDataFrameDTO, DataFrame } from '@grafana/data'; import { toDataQueryResponse } from './queryResponse'; /* eslint-disable */ const resp = { data: { results: { - GC: { + A: { + refId: 'A', + series: null, + tables: null, dataframes: [ - 'QVJST1cxAACsAQAAEAAAAAAACgAOAAwACwAEAAoAAAAUAAAAAAAAAQMACgAMAAAACAAEAAoAAAAIAAAAUAAAAAIAAAAoAAAABAAAAOD+//8IAAAADAAAAAIAAABHQwAABQAAAHJlZklkAAAAAP///wgAAAAMAAAAAAAAAAAAAAAEAAAAbmFtZQAAAAACAAAAlAAAAAQAAACG////FAAAAGAAAABgAAAAAAADAWAAAAACAAAALAAAAAQAAABQ////CAAAABAAAAAGAAAAbnVtYmVyAAAEAAAAdHlwZQAAAAB0////CAAAAAwAAAAAAAAAAAAAAAQAAABuYW1lAAAAAAAAAABm////AAACAAAAAAAAABIAGAAUABMAEgAMAAAACAAEABIAAAAUAAAAbAAAAHQAAAAAAAoBdAAAAAIAAAA0AAAABAAAANz///8IAAAAEAAAAAQAAAB0aW1lAAAAAAQAAAB0eXBlAAAAAAgADAAIAAQACAAAAAgAAAAQAAAABAAAAFRpbWUAAAAABAAAAG5hbWUAAAAAAAAAAAAABgAIAAYABgAAAAAAAwAEAAAAVGltZQAAAAC8AAAAFAAAAAAAAAAMABYAFAATAAwABAAMAAAA0AAAAAAAAAAUAAAAAAAAAwMACgAYAAwACAAEAAoAAAAUAAAAWAAAAA0AAAAAAAAAAAAAAAQAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAABoAAAAAAAAAGgAAAAAAAAAAAAAAAAAAABoAAAAAAAAAGgAAAAAAAAAAAAAAAIAAAANAAAAAAAAAAAAAAAAAAAADQAAAAAAAAAAAAAAAAAAAAAAAAAAFp00e2XHFQAIo158ZccVAPqoiH1lxxUA7K6yfmXHFQDetNx/ZccVANC6BoFlxxUAwsAwgmXHFQC0xlqDZccVAKbMhIRlxxUAmNKuhWXHFQCK2NiGZccVAHzeAohlxxUAbuQsiWXHFQAAAAAAAAhAAAAAAAAACEAAAAAAAAAIQAAAAAAAABRAAAAAAAAAFEAAAAAAAAAUQAAAAAAAAAhAAAAAAAAACEAAAAAAAAAIQAAAAAAAABRAAAAAAAAAFEAAAAAAAAAUQAAAAAAAAAhAEAAAAAwAFAASAAwACAAEAAwAAAAQAAAALAAAADgAAAAAAAMAAQAAALgBAAAAAAAAwAAAAAAAAADQAAAAAAAAAAAAAAAAAAAAAAAKAAwAAAAIAAQACgAAAAgAAABQAAAAAgAAACgAAAAEAAAA4P7//wgAAAAMAAAAAgAAAEdDAAAFAAAAcmVmSWQAAAAA////CAAAAAwAAAAAAAAAAAAAAAQAAABuYW1lAAAAAAIAAACUAAAABAAAAIb///8UAAAAYAAAAGAAAAAAAAMBYAAAAAIAAAAsAAAABAAAAFD///8IAAAAEAAAAAYAAABudW1iZXIAAAQAAAB0eXBlAAAAAHT///8IAAAADAAAAAAAAAAAAAAABAAAAG5hbWUAAAAAAAAAAGb///8AAAIAAAAAAAAAEgAYABQAEwASAAwAAAAIAAQAEgAAABQAAABsAAAAdAAAAAAACgF0AAAAAgAAADQAAAAEAAAA3P///wgAAAAQAAAABAAAAHRpbWUAAAAABAAAAHR5cGUAAAAACAAMAAgABAAIAAAACAAAABAAAAAEAAAAVGltZQAAAAAEAAAAbmFtZQAAAAAAAAAAAAAGAAgABgAGAAAAAAADAAQAAABUaW1lAAAAANgBAABBUlJPVzE=', + 'QVJST1cxAAD/////cAEAABAAAAAAAAoADgAMAAsABAAKAAAAFAAAAAAAAAEDAAoADAAAAAgABAAKAAAACAAAAFAAAAACAAAAKAAAAAQAAAAg////CAAAAAwAAAABAAAAQQAAAAUAAAByZWZJZAAAAED///8IAAAADAAAAAAAAAAAAAAABAAAAG5hbWUAAAAAAgAAAHwAAAAEAAAAnv///xQAAABAAAAAQAAAAAAAAwFAAAAAAQAAAAQAAACM////CAAAABQAAAAIAAAAQS1zZXJpZXMAAAAABAAAAG5hbWUAAAAAAAAAAIb///8AAAIACAAAAEEtc2VyaWVzAAASABgAFAATABIADAAAAAgABAASAAAAFAAAAEQAAABMAAAAAAAKAUwAAAABAAAADAAAAAgADAAIAAQACAAAAAgAAAAQAAAABAAAAHRpbWUAAAAABAAAAG5hbWUAAAAAAAAAAAAABgAIAAYABgAAAAAAAwAEAAAAdGltZQAAAAAAAAAA/////7gAAAAUAAAAAAAAAAwAFgAUABMADAAEAAwAAABgAAAAAAAAABQAAAAAAAADAwAKABgADAAIAAQACgAAABQAAABYAAAABgAAAAAAAAAAAAAABAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAADAAAAAAAAAAMAAAAAAAAAAAAAAAAAAAADAAAAAAAAAAMAAAAAAAAAAAAAAAAgAAAAYAAAAAAAAAAAAAAAAAAAAGAAAAAAAAAAAAAAAAAAAAQMC/OcElXhZAOAEFxCVeFkCwQtDGJV4WQCiEm8klXhZAoMVmzCVeFkAYBzLPJV4WAAAAAAAA8D8AAAAAAAA0QAAAAAAAgFZAAAAAAAAAPkAAAAAAAAAUQAAAAAAAAAAAEAAAAAwAFAASAAwACAAEAAwAAAAQAAAALAAAADgAAAAAAAMAAQAAAIABAAAAAAAAwAAAAAAAAABgAAAAAAAAAAAAAAAAAAAAAAAKAAwAAAAIAAQACgAAAAgAAABQAAAAAgAAACgAAAAEAAAAIP///wgAAAAMAAAAAQAAAEEAAAAFAAAAcmVmSWQAAABA////CAAAAAwAAAAAAAAAAAAAAAQAAABuYW1lAAAAAAIAAAB8AAAABAAAAJ7///8UAAAAQAAAAEAAAAAAAAMBQAAAAAEAAAAEAAAAjP///wgAAAAUAAAACAAAAEEtc2VyaWVzAAAAAAQAAABuYW1lAAAAAAAAAACG////AAACAAgAAABBLXNlcmllcwAAEgAYABQAEwASAAwAAAAIAAQAEgAAABQAAABEAAAATAAAAAAACgFMAAAAAQAAAAwAAAAIAAwACAAEAAgAAAAIAAAAEAAAAAQAAAB0aW1lAAAAAAQAAABuYW1lAAAAAAAAAAAAAAYACAAGAAYAAAAAAAMABAAAAHRpbWUAAAAAmAEAAEFSUk9XMQ==', + ], + }, + B: { + refId: 'B', + series: null, + tables: null, + dataframes: [ + 'QVJST1cxAAD/////cAEAABAAAAAAAAoADgAMAAsABAAKAAAAFAAAAAAAAAEDAAoADAAAAAgABAAKAAAACAAAAFAAAAACAAAAKAAAAAQAAAAg////CAAAAAwAAAABAAAAQgAAAAUAAAByZWZJZAAAAED///8IAAAADAAAAAAAAAAAAAAABAAAAG5hbWUAAAAAAgAAAHwAAAAEAAAAnv///xQAAABAAAAAQAAAAAAAAwFAAAAAAQAAAAQAAACM////CAAAABQAAAAIAAAAQi1zZXJpZXMAAAAABAAAAG5hbWUAAAAAAAAAAIb///8AAAIACAAAAEItc2VyaWVzAAASABgAFAATABIADAAAAAgABAASAAAAFAAAAEQAAABMAAAAAAAKAUwAAAABAAAADAAAAAgADAAIAAQACAAAAAgAAAAQAAAABAAAAHRpbWUAAAAABAAAAG5hbWUAAAAAAAAAAAAABgAIAAYABgAAAAAAAwAEAAAAdGltZQAAAAAAAAAA/////7gAAAAUAAAAAAAAAAwAFgAUABMADAAEAAwAAABgAAAAAAAAABQAAAAAAAADAwAKABgADAAIAAQACgAAABQAAABYAAAABgAAAAAAAAAAAAAABAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAADAAAAAAAAAAMAAAAAAAAAAAAAAAAAAAADAAAAAAAAAAMAAAAAAAAAAAAAAAAgAAAAYAAAAAAAAAAAAAAAAAAAAGAAAAAAAAAAAAAAAAAAAAQMC/OcElXhZAOAEFxCVeFkCwQtDGJV4WQCiEm8klXhZAoMVmzCVeFkAYBzLPJV4WAAAAAAAA8D8AAAAAAAA0QAAAAAAAgFZAAAAAAAAAPkAAAAAAAAAUQAAAAAAAAAAAEAAAAAwAFAASAAwACAAEAAwAAAAQAAAALAAAADgAAAAAAAMAAQAAAIABAAAAAAAAwAAAAAAAAABgAAAAAAAAAAAAAAAAAAAAAAAKAAwAAAAIAAQACgAAAAgAAABQAAAAAgAAACgAAAAEAAAAIP///wgAAAAMAAAAAQAAAEIAAAAFAAAAcmVmSWQAAABA////CAAAAAwAAAAAAAAAAAAAAAQAAABuYW1lAAAAAAIAAAB8AAAABAAAAJ7///8UAAAAQAAAAEAAAAAAAAMBQAAAAAEAAAAEAAAAjP///wgAAAAUAAAACAAAAEItc2VyaWVzAAAAAAQAAABuYW1lAAAAAAAAAACG////AAACAAgAAABCLXNlcmllcwAAEgAYABQAEwASAAwAAAAIAAQAEgAAABQAAABEAAAATAAAAAAACgFMAAAAAQAAAAwAAAAIAAwACAAEAAgAAAAIAAAAEAAAAAQAAAB0aW1lAAAAAAQAAABuYW1lAAAAAAAAAAAAAAYACAAGAAYAAAAAAAMABAAAAHRpbWUAAAAAmAEAAEFSUk9XMQ==', ], - frames: null as any, }, }, }, @@ -41,9 +50,9 @@ describe('Query Response parser', () => { test('should parse output with dataframe', () => { const res = toDataQueryResponse(resp); const frames = res.data; - for (const frame of frames) { - expect(frame.refId).toEqual('GC'); - } + expect(frames).toHaveLength(2); + expect(frames[0].refId).toEqual('A'); + expect(frames[1].refId).toEqual('B'); const norm = frames.map((f) => toDataFrameDTO(f)); expect(norm).toMatchInlineSnapshot(` @@ -53,49 +62,155 @@ describe('Query Response parser', () => { Object { "config": Object {}, "labels": undefined, - "name": "Time", + "name": "time", "type": "time", "values": Array [ - 1569334575000, - 1569334580000, - 1569334585000, - 1569334590000, - 1569334595000, - 1569334600000, - 1569334605000, - 1569334610000, - 1569334615000, - 1569334620000, - 1569334625000, - 1569334630000, - 1569334635000, + 1611767228473, + 1611767240473, + 1611767252473, + 1611767264473, + 1611767276473, + 1611767288473, ], }, Object { "config": Object {}, "labels": undefined, - "name": "", + "name": "A-series", "type": "number", "values": Array [ - 3, - 3, - 3, + 1, + 20, + 90, + 30, 5, - 5, - 5, - 3, - 3, - 3, - 5, - 5, - 5, - 3, + 0, ], }, ], "meta": undefined, "name": undefined, - "refId": "GC", + "refId": "A", + }, + Object { + "fields": Array [ + Object { + "config": Object {}, + "labels": undefined, + "name": "time", + "type": "time", + "values": Array [ + 1611767228473, + 1611767240473, + 1611767252473, + 1611767264473, + 1611767276473, + 1611767288473, + ], + }, + Object { + "config": Object {}, + "labels": undefined, + "name": "B-series", + "type": "number", + "values": Array [ + 1, + 20, + 90, + 30, + 5, + 0, + ], + }, + ], + "meta": undefined, + "name": undefined, + "refId": "B", + }, + ] + `); + }); + + test('should parse output with dataframe in order of queries', () => { + const queries: DataQuery[] = [{ refId: 'B' }, { refId: 'A' }]; + const res = toDataQueryResponse(resp, queries); + const frames = res.data; + expect(frames).toHaveLength(2); + expect(frames[0].refId).toEqual('B'); + expect(frames[1].refId).toEqual('A'); + + const norm = frames.map((f) => toDataFrameDTO(f)); + expect(norm).toMatchInlineSnapshot(` + Array [ + Object { + "fields": Array [ + Object { + "config": Object {}, + "labels": undefined, + "name": "time", + "type": "time", + "values": Array [ + 1611767228473, + 1611767240473, + 1611767252473, + 1611767264473, + 1611767276473, + 1611767288473, + ], + }, + Object { + "config": Object {}, + "labels": undefined, + "name": "B-series", + "type": "number", + "values": Array [ + 1, + 20, + 90, + 30, + 5, + 0, + ], + }, + ], + "meta": undefined, + "name": undefined, + "refId": "B", + }, + Object { + "fields": Array [ + Object { + "config": Object {}, + "labels": undefined, + "name": "time", + "type": "time", + "values": Array [ + 1611767228473, + 1611767240473, + 1611767252473, + 1611767264473, + 1611767276473, + 1611767288473, + ], + }, + Object { + "config": Object {}, + "labels": undefined, + "name": "A-series", + "type": "number", + "values": Array [ + 1, + 20, + 90, + 30, + 5, + 0, + ], + }, + ], + "meta": undefined, + "name": undefined, + "refId": "A", }, ] `); @@ -106,6 +221,35 @@ describe('Query Response parser', () => { expect(frames.length).toEqual(0); }); + test('keeps query order', () => { + const resp = { + data: { + results: { + X: { + series: [ + { name: 'Requests/s', points: [[13.594958983547151, 1611839862951]], tables: null, dataframes: null }, + ], + }, + B: { + series: [ + { name: 'Requests/s', points: [[13.594958983547151, 1611839862951]], tables: null, dataframes: null }, + ], + }, + A: { + series: [ + { name: 'Requests/s', points: [[13.594958983547151, 1611839862951]], tables: null, dataframes: null }, + ], + }, + }, + }, + }; + + const queries: DataQuery[] = [{ refId: 'A' }, { refId: 'B' }]; + + const ids = (toDataQueryResponse(resp, queries).data as DataFrame[]).map((f) => f.refId); + expect(ids).toEqual(['A', 'B', 'X']); + }); + test('resultWithError', () => { // Generated from: // qdr.Responses[q.GetRefID()] = backend.DataResponse{ diff --git a/packages/grafana-runtime/src/utils/queryResponse.ts b/packages/grafana-runtime/src/utils/queryResponse.ts index 0ce0424a482..83f21d872cb 100644 --- a/packages/grafana-runtime/src/utils/queryResponse.ts +++ b/packages/grafana-runtime/src/utils/queryResponse.ts @@ -11,6 +11,7 @@ import { DataFrame, MetricFindValue, FieldType, + DataQuery, } from '@grafana/data'; interface DataResponse { @@ -24,56 +25,84 @@ interface DataResponse { /** * Parse the results from /api/ds/query into a DataQueryResponse * + * @param res - the HTTP response data. + * @param queries - optional DataQuery array that will order the response based on the order of query refId's. + * * @public */ -export function toDataQueryResponse(res: any): DataQueryResponse { +export function toDataQueryResponse(res: any, queries?: DataQuery[]): DataQueryResponse { const rsp: DataQueryResponse = { data: [], state: LoadingState.Done }; if (res.data?.results) { const results: KeyValue = res.data.results; - for (const refId of Object.keys(results)) { + const resultIDs = Object.keys(results); + const refIDs = queries ? queries.map((q) => q.refId) : resultIDs; + const usedResultIDs = new Set(resultIDs); + const data: DataResponse[] = []; + + for (const refId of refIDs) { const dr = results[refId] as DataResponse; - if (dr) { - if (dr.error) { - if (!rsp.error) { - rsp.error = { - refId, - message: dr.error, - }; + if (!dr) { + continue; + } + dr.refId = refId; + usedResultIDs.delete(refId); + data.push(dr); + } + + // Add any refIds that do not match the query targets + if (usedResultIDs.size) { + for (const refId of usedResultIDs) { + const dr = results[refId] as DataResponse; + if (!dr) { + continue; + } + dr.refId = refId; + usedResultIDs.delete(refId); + data.push(dr); + } + } + + for (const dr of data) { + if (dr.error) { + if (!rsp.error) { + rsp.error = { + refId: dr.refId, + message: dr.error, + }; + rsp.state = LoadingState.Error; + } + } + + if (dr.series?.length) { + for (const s of dr.series) { + if (!s.refId) { + s.refId = dr.refId; + } + rsp.data.push(toDataFrame(s)); + } + } + + if (dr.tables?.length) { + for (const s of dr.tables) { + if (!s.refId) { + s.refId = dr.refId; + } + rsp.data.push(toDataFrame(s)); + } + } + + if (dr.dataframes) { + for (const b64 of dr.dataframes) { + try { + const t = base64StringToArrowTable(b64); + const f = arrowTableToDataFrame(t); + if (!f.refId) { + f.refId = dr.refId; + } + rsp.data.push(f); + } catch (err) { rsp.state = LoadingState.Error; - } - } - - if (dr.series && dr.series.length) { - for (const s of dr.series) { - if (!s.refId) { - s.refId = refId; - } - rsp.data.push(toDataFrame(s)); - } - } - - if (dr.tables && dr.tables.length) { - for (const s of dr.tables) { - if (!s.refId) { - s.refId = refId; - } - rsp.data.push(toDataFrame(s)); - } - } - - if (dr.dataframes) { - for (const b64 of dr.dataframes) { - try { - const t = base64StringToArrowTable(b64); - const f = arrowTableToDataFrame(t); - if (!f.refId) { - f.refId = refId; - } - rsp.data.push(f); - } catch (err) { - rsp.state = LoadingState.Error; - rsp.error = toDataQueryError(err); - } + rsp.error = toDataQueryError(err); } } } diff --git a/pkg/api/api.go b/pkg/api/api.go index 5fecd0f8ccf..3ef80ebb014 100644 --- a/pkg/api/api.go +++ b/pkg/api/api.go @@ -350,7 +350,6 @@ func (hs *HTTPServer) registerRoutes() { // metrics apiRoute.Post("/tsdb/query", bind(dtos.MetricRequest{}), routing.Wrap(hs.QueryMetrics)) - apiRoute.Get("/tsdb/testdata/scenarios", routing.Wrap(GetTestDataScenarios)) apiRoute.Get("/tsdb/testdata/gensql", reqGrafanaAdmin, routing.Wrap(GenerateSQLTestData)) apiRoute.Get("/tsdb/testdata/random-walk", routing.Wrap(GetTestDataRandomWalk)) @@ -395,9 +394,6 @@ func (hs *HTTPServer) registerRoutes() { annotationsRoute.Post("/graphite", reqEditorRole, bind(dtos.PostGraphiteAnnotationsCmd{}), routing.Wrap(PostGraphiteAnnotation)) }) - // error test - r.Get("/metrics/error", routing.Wrap(GenerateError)) - // short urls apiRoute.Post("/short-urls", bind(dtos.CreateShortURLCmd{}), routing.Wrap(hs.createShortURL)) }, reqSignedIn) diff --git a/pkg/api/metrics.go b/pkg/api/metrics.go index 8674709ec45..8ad60076d28 100644 --- a/pkg/api/metrics.go +++ b/pkg/api/metrics.go @@ -3,7 +3,6 @@ package api import ( "context" "errors" - "sort" "github.com/grafana/grafana/pkg/expr" "github.com/grafana/grafana/pkg/models" @@ -13,7 +12,6 @@ import ( "github.com/grafana/grafana/pkg/bus" "github.com/grafana/grafana/pkg/components/simplejson" "github.com/grafana/grafana/pkg/tsdb" - "github.com/grafana/grafana/pkg/tsdb/testdatasource" "github.com/grafana/grafana/pkg/util" ) @@ -202,36 +200,6 @@ func (hs *HTTPServer) QueryMetrics(c *models.ReqContext, reqDto dtos.MetricReque return response.JSON(statusCode, &resp) } -// GET /api/tsdb/testdata/scenarios -func GetTestDataScenarios(c *models.ReqContext) response.Response { - result := make([]interface{}, 0) - - scenarioIds := make([]string, 0) - for id := range testdatasource.ScenarioRegistry { - scenarioIds = append(scenarioIds, id) - } - sort.Strings(scenarioIds) - - for _, scenarioId := range scenarioIds { - scenario := testdatasource.ScenarioRegistry[scenarioId] - result = append(result, map[string]interface{}{ - "id": scenario.Id, - "name": scenario.Name, - "description": scenario.Description, - "stringInput": scenario.StringInput, - }) - } - - return response.JSON(200, &result) -} - -// GenerateError generates a index out of range error -func GenerateError(c *models.ReqContext) response.Response { - var array []string - // nolint: govet - return response.JSON(200, array[20]) -} - // GET /api/tsdb/testdata/gensql func GenerateSQLTestData(c *models.ReqContext) response.Response { if err := bus.Dispatch(&models.InsertSQLTestDataCommand{}); err != nil { @@ -250,7 +218,10 @@ func GetTestDataRandomWalk(c *models.ReqContext) response.Response { timeRange := tsdb.NewTimeRange(from, to) request := &tsdb.TsdbQuery{TimeRange: timeRange} - dsInfo := &models.DataSource{Type: "testdata"} + dsInfo := &models.DataSource{ + Type: "testdata", + JsonData: simplejson.New(), + } request.Queries = append(request.Queries, &tsdb.Query{ RefId: "A", IntervalMs: intervalMs, diff --git a/pkg/plugins/backendplugin/coreplugin/core_plugin.go b/pkg/plugins/backendplugin/coreplugin/core_plugin.go index d746e2b46bd..7732aadf71e 100644 --- a/pkg/plugins/backendplugin/coreplugin/core_plugin.go +++ b/pkg/plugins/backendplugin/coreplugin/core_plugin.go @@ -11,7 +11,6 @@ import ( ) // corePlugin represents a plugin that's part of Grafana core. -// nolint:unused type corePlugin struct { pluginID string logger log.Logger @@ -55,7 +54,7 @@ func (cp *corePlugin) Stop(ctx context.Context) error { } func (cp *corePlugin) IsManaged() bool { - return false + return true } func (cp *corePlugin) Exited() bool { diff --git a/pkg/plugins/backendplugin/coreplugin/core_plugin_test.go b/pkg/plugins/backendplugin/coreplugin/core_plugin_test.go index d9fa8be7f2a..18732afc5b0 100644 --- a/pkg/plugins/backendplugin/coreplugin/core_plugin_test.go +++ b/pkg/plugins/backendplugin/coreplugin/core_plugin_test.go @@ -19,7 +19,7 @@ func TestCorePlugin(t *testing.T) { require.NotNil(t, p) require.NoError(t, p.Start(context.Background())) require.NoError(t, p.Stop(context.Background())) - require.False(t, p.IsManaged()) + require.True(t, p.IsManaged()) require.False(t, p.Exited()) _, err = p.CollectMetrics(context.Background()) @@ -50,7 +50,7 @@ func TestCorePlugin(t *testing.T) { require.NotNil(t, p) require.NoError(t, p.Start(context.Background())) require.NoError(t, p.Stop(context.Background())) - require.False(t, p.IsManaged()) + require.True(t, p.IsManaged()) require.False(t, p.Exited()) _, err = p.CollectMetrics(context.Background()) diff --git a/pkg/tsdb/testdatasource/resource_handler.go b/pkg/tsdb/testdatasource/resource_handler.go new file mode 100644 index 00000000000..7929b6fc285 --- /dev/null +++ b/pkg/tsdb/testdatasource/resource_handler.go @@ -0,0 +1,157 @@ +package testdatasource + +import ( + "encoding/json" + "fmt" + "io" + "io/ioutil" + "net/http" + "sort" + "strconv" + "time" + + "github.com/grafana/grafana/pkg/infra/log" + + "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 (p *testDataPlugin) testGetHandler(rw http.ResponseWriter, req *http.Request) { + p.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) + return + } + rw.WriteHeader(http.StatusOK) +} + +func (p *testDataPlugin) getScenariosHandler(rw http.ResponseWriter, req *http.Request) { + result := make([]interface{}, 0) + + scenarioIds := make([]string, 0) + for id := range p.scenarios { + scenarioIds = append(scenarioIds, id) + } + sort.Strings(scenarioIds) + + for _, scenarioID := range scenarioIds { + scenario := p.scenarios[scenarioID] + result = append(result, map[string]interface{}{ + "id": scenario.ID, + "name": scenario.Name, + "description": scenario.Description, + "stringInput": scenario.StringInput, + }) + } + + bytes, err := json.Marshal(&result) + if err != nil { + p.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) + } +} + +func (p *testDataPlugin) testStreamHandler(rw http.ResponseWriter, req *http.Request) { + p.logger.Debug("Received resource call", "url", req.URL.String(), "method", req.Method) + + if req.Method != http.MethodGet { + return + } + + count := 10 + countstr := req.URL.Query().Get("count") + if countstr != "" { + if i, err := strconv.Atoi(countstr); err == nil { + count = i + } + } + + sleep := req.URL.Query().Get("sleep") + sleepDuration, err := time.ParseDuration(sleep) + if err != nil { + sleepDuration = time.Millisecond + } + + rw.Header().Set("Content-Type", "text/plain") + rw.WriteHeader(http.StatusOK) + + 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) + return + } + rw.(http.Flusher).Flush() + time.Sleep(sleepDuration) + } +} + +func createJSONHandler(logger log.Logger) http.Handler { + return http.HandlerFunc(func(rw http.ResponseWriter, req *http.Request) { + logger.Debug("Received resource call", "url", req.URL.String(), "method", req.Method) + + var reqData map[string]interface{} + if req.Body != nil { + defer func() { + if err := req.Body.Close(); err != nil { + logger.Warn("Failed to close response body", "err", err) + } + }() + b, err := ioutil.ReadAll(req.Body) + if err != nil { + logger.Error("Failed to read request body to bytes", "error", err) + } else { + err := json.Unmarshal(b, &reqData) + if err != nil { + logger.Error("Failed to unmarshal request body to JSON", "error", err) + } + + logger.Debug("Received resource call body", "body", reqData) + } + } + + config := httpadapter.PluginConfigFromContext(req.Context()) + + data := map[string]interface{}{ + "message": "Hello world from test datasource!", + "request": map[string]interface{}{ + "method": req.Method, + "url": req.URL, + "headers": req.Header, + "body": reqData, + "config": config, + }, + } + bytes, err := json.Marshal(&data) + if err != nil { + 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 { + logger.Error("Failed to write response", "error", err) + } + }) +} + +func (p *testDataPlugin) testPanicHandler(rw http.ResponseWriter, req *http.Request) { + panic("BOOM") +} diff --git a/pkg/tsdb/testdatasource/scenarios.go b/pkg/tsdb/testdatasource/scenarios.go index d250f0d76d5..a6b2b5ce5b5 100644 --- a/pkg/tsdb/testdatasource/scenarios.go +++ b/pkg/tsdb/testdatasource/scenarios.go @@ -1,6 +1,7 @@ package testdatasource import ( + "context" "encoding/json" "fmt" "math" @@ -9,583 +10,641 @@ import ( "strings" "time" + "github.com/grafana/grafana-plugin-sdk-go/backend" + "github.com/grafana/grafana-plugin-sdk-go/data" + "github.com/grafana/grafana/pkg/components/null" "github.com/grafana/grafana/pkg/components/simplejson" "github.com/grafana/grafana/pkg/util/errutil" - - "github.com/grafana/grafana/pkg/components/null" - "github.com/grafana/grafana/pkg/infra/log" - "github.com/grafana/grafana/pkg/tsdb" ) -type ScenarioHandler func(query *tsdb.Query, context *tsdb.TsdbQuery) *tsdb.QueryResult +const ( + randomWalkQuery queryType = "random_walk" + randomWalkSlowQuery queryType = "slow_query" + randomWalkWithErrorQuery queryType = "random_walk_with_error" + randomWalkTableQuery queryType = "random_walk_table" + exponentialHeatmapBucketDataQuery queryType = "exponential_heatmap_bucket_data" + linearHeatmapBucketDataQuery queryType = "linear_heatmap_bucket_data" + noDataPointsQuery queryType = "no_data_points" + datapointsOutsideRangeQuery queryType = "datapoints_outside_range" + csvMetricValuesQuery queryType = "csv_metric_values" + manualEntryQuery queryType = "manual_entry" + predictablePulseQuery queryType = "predictable_pulse" + predictableCSVWaveQuery queryType = "predictable_csv_wave" + streamingClientQuery queryType = "streaming_client" + liveQuery queryType = "live" + grafanaAPIQuery queryType = "grafana_api" + arrowQuery queryType = "arrow" + annotationsQuery queryType = "annotations" + tableStaticQuery queryType = "table_static" + serverError500Query queryType = "server_error_500" + logsQuery queryType = "logs" + nodeGraphQuery queryType = "node_graph" +) + +type queryType string type Scenario struct { - Id string `json:"id"` - Name string `json:"name"` - StringInput string `json:"stringOption"` - Description string `json:"description"` - Handler ScenarioHandler `json:"-"` + ID string `json:"id"` + Name string `json:"name"` + StringInput string `json:"stringInput"` + Description string `json:"description"` + handler backend.QueryDataHandlerFunc } -var ScenarioRegistry map[string]*Scenario - -func init() { - ScenarioRegistry = make(map[string]*Scenario) - logger := log.New("tsdb.testdata") - - logger.Debug("Initializing TestData Scenario") - - registerScenario(&Scenario{ - Id: "exponential_heatmap_bucket_data", - Name: "Exponential heatmap bucket data", - - Handler: func(query *tsdb.Query, context *tsdb.TsdbQuery) *tsdb.QueryResult { - to := context.TimeRange.GetToAsMsEpoch() - - var series []*tsdb.TimeSeries - start := 1 - factor := 2 - for i := 0; i < 10; i++ { - timeWalkerMs := context.TimeRange.GetFromAsMsEpoch() - ts := &tsdb.TimeSeries{Name: strconv.Itoa(start)} - start *= factor - - points := make(tsdb.TimeSeriesPoints, 0) - for j := int64(0); j < 100 && timeWalkerMs < to; j++ { - v := float64(rand.Int63n(100)) - points = append(points, tsdb.NewTimePoint(null.FloatFrom(v), float64(timeWalkerMs))) - timeWalkerMs += query.IntervalMs * 50 - } - - ts.Points = points - series = append(series, ts) - } - - queryRes := tsdb.NewQueryResult() - queryRes.Series = append(queryRes.Series, series...) - return queryRes - }, - }) - - registerScenario(&Scenario{ - Id: "linear_heatmap_bucket_data", - Name: "Linear heatmap bucket data", - - Handler: func(query *tsdb.Query, context *tsdb.TsdbQuery) *tsdb.QueryResult { - to := context.TimeRange.GetToAsMsEpoch() - - var series []*tsdb.TimeSeries - for i := 0; i < 10; i++ { - timeWalkerMs := context.TimeRange.GetFromAsMsEpoch() - ts := &tsdb.TimeSeries{Name: strconv.Itoa(i * 10)} - - points := make(tsdb.TimeSeriesPoints, 0) - for j := int64(0); j < 100 && timeWalkerMs < to; j++ { - v := float64(rand.Int63n(100)) - points = append(points, tsdb.NewTimePoint(null.FloatFrom(v), float64(timeWalkerMs))) - timeWalkerMs += query.IntervalMs * 50 - } - - ts.Points = points - series = append(series, ts) - } - - queryRes := tsdb.NewQueryResult() - queryRes.Series = append(queryRes.Series, series...) - return queryRes - }, - }) - - registerScenario(&Scenario{ - Id: "random_walk", - Name: "Random Walk", - - Handler: func(query *tsdb.Query, context *tsdb.TsdbQuery) *tsdb.QueryResult { - queryRes := tsdb.NewQueryResult() - - seriesCount := query.Model.Get("seriesCount").MustInt(1) - - for i := 0; i < seriesCount; i++ { - queryRes.Series = append(queryRes.Series, getRandomWalk(query, context, i)) - } - - return queryRes - }, - }) - - registerScenario(&Scenario{ - Id: "predictable_pulse", - Name: "Predictable Pulse", - Handler: getPredictablePulse, - Description: PredictablePulseDesc, - }) - - registerScenario(&Scenario{ - Id: "predictable_csv_wave", - Name: "Predictable CSV Wave", - Handler: getPredictableCSVWave, - }) - - registerScenario(&Scenario{ - Id: "random_walk_table", - Name: "Random Walk Table", - Handler: getRandomWalkTable, - }) - - registerScenario(&Scenario{ - Id: "slow_query", - Name: "Slow Query", - StringInput: "5s", - Handler: func(query *tsdb.Query, context *tsdb.TsdbQuery) *tsdb.QueryResult { - stringInput := query.Model.Get("stringInput").MustString() - parsedInterval, _ := time.ParseDuration(stringInput) - time.Sleep(parsedInterval) - - queryRes := tsdb.NewQueryResult() - queryRes.Series = append(queryRes.Series, getRandomWalk(query, context, 0)) - return queryRes - }, - }) - - registerScenario(&Scenario{ - Id: "no_data_points", - Name: "No Data Points", - Handler: func(query *tsdb.Query, context *tsdb.TsdbQuery) *tsdb.QueryResult { - return tsdb.NewQueryResult() - }, - }) - - registerScenario(&Scenario{ - Id: "datapoints_outside_range", - Name: "Datapoints Outside Range", - Handler: func(query *tsdb.Query, context *tsdb.TsdbQuery) *tsdb.QueryResult { - queryRes := tsdb.NewQueryResult() - - series := newSeriesForQuery(query, 0) - outsideTime := context.TimeRange.MustGetFrom().Add(-1*time.Hour).Unix() * 1000 - - series.Points = append(series.Points, tsdb.NewTimePoint(null.FloatFrom(10), float64(outsideTime))) - queryRes.Series = append(queryRes.Series, series) - - return queryRes - }, - }) - - registerScenario(&Scenario{ - Id: "manual_entry", - Name: "Manual Entry", - Handler: func(query *tsdb.Query, context *tsdb.TsdbQuery) *tsdb.QueryResult { - queryRes := tsdb.NewQueryResult() - - points := query.Model.Get("points").MustArray() - - series := newSeriesForQuery(query, 0) - startTime := context.TimeRange.GetFromAsMsEpoch() - endTime := context.TimeRange.GetToAsMsEpoch() - - for _, val := range points { - pointValues := val.([]interface{}) - - var value null.Float - var time int64 - - if valueFloat, err := strconv.ParseFloat(string(pointValues[0].(json.Number)), 64); err == nil { - value = null.FloatFrom(valueFloat) - } - - timeInt, err := strconv.ParseInt(string(pointValues[1].(json.Number)), 10, 64) - if err != nil { - continue - } - time = timeInt - - if time >= startTime && time <= endTime { - series.Points = append(series.Points, tsdb.NewTimePoint(value, float64(time))) - } - } - - queryRes.Series = append(queryRes.Series, series) - - return queryRes - }, - }) - - registerScenario(&Scenario{ - Id: "csv_metric_values", - Name: "CSV Metric Values", - StringInput: "1,20,90,30,5,0", - Handler: func(query *tsdb.Query, context *tsdb.TsdbQuery) *tsdb.QueryResult { - queryRes := tsdb.NewQueryResult() - - stringInput := query.Model.Get("stringInput").MustString() - stringInput = strings.ReplaceAll(stringInput, " ", "") - - values := []null.Float{} - for _, strVal := range strings.Split(stringInput, ",") { - if strVal == "null" { - values = append(values, null.FloatFromPtr(nil)) - } - if val, err := strconv.ParseFloat(strVal, 64); err == nil { - values = append(values, null.FloatFrom(val)) - } - } - - if len(values) == 0 { - return queryRes - } - - series := newSeriesForQuery(query, 0) - startTime := context.TimeRange.GetFromAsMsEpoch() - endTime := context.TimeRange.GetToAsMsEpoch() - var step int64 = 0 - if len(values) > 1 { - step = (endTime - startTime) / int64(len(values)-1) - } - - for _, val := range values { - series.Points = append(series.Points, tsdb.TimePoint{val, null.FloatFrom(float64(startTime))}) - startTime += step - } - - queryRes.Series = append(queryRes.Series, series) - - return queryRes - }, - }) - - registerScenario(&Scenario{ - Id: "streaming_client", - Name: "Streaming Client", - Handler: func(query *tsdb.Query, context *tsdb.TsdbQuery) *tsdb.QueryResult { - // Real work is in javascript client - return tsdb.NewQueryResult() - }, - }) - - registerScenario(&Scenario{ - Id: "live", - Name: "Grafana Live", - Handler: func(query *tsdb.Query, context *tsdb.TsdbQuery) *tsdb.QueryResult { - // Real work is in javascript client - return tsdb.NewQueryResult() - }, - }) - - registerScenario(&Scenario{ - Id: "grafana_api", - Name: "Grafana API", - Handler: func(query *tsdb.Query, context *tsdb.TsdbQuery) *tsdb.QueryResult { - // Real work is in javascript client - return tsdb.NewQueryResult() - }, - }) - - registerScenario(&Scenario{ - Id: "arrow", - Name: "Load Apache Arrow Data", - Handler: func(query *tsdb.Query, context *tsdb.TsdbQuery) *tsdb.QueryResult { - // Real work is in javascript client - return tsdb.NewQueryResult() - }, - }) - - registerScenario(&Scenario{ - Id: "annotations", - Name: "Annotations", - Handler: func(query *tsdb.Query, context *tsdb.TsdbQuery) *tsdb.QueryResult { - return tsdb.NewQueryResult() - }, - }) - - registerScenario(&Scenario{ - Id: "table_static", - Name: "Table Static", - - Handler: func(query *tsdb.Query, context *tsdb.TsdbQuery) *tsdb.QueryResult { - timeWalkerMs := context.TimeRange.GetFromAsMsEpoch() - to := context.TimeRange.GetToAsMsEpoch() - - table := tsdb.Table{ - Columns: []tsdb.TableColumn{ - {Text: "Time"}, - {Text: "Message"}, - {Text: "Description"}, - {Text: "Value"}, - }, - Rows: []tsdb.RowValues{}, - } - for i := int64(0); i < 10 && timeWalkerMs < to; i++ { - table.Rows = append(table.Rows, tsdb.RowValues{float64(timeWalkerMs), "This is a message", "Description", 23.1}) - timeWalkerMs += query.IntervalMs - } - - queryRes := tsdb.NewQueryResult() - queryRes.Tables = append(queryRes.Tables, &table) - return queryRes - }, - }) - - registerScenario(&Scenario{ - Id: "random_walk_with_error", - Name: "Random Walk (with error)", - - Handler: func(query *tsdb.Query, context *tsdb.TsdbQuery) *tsdb.QueryResult { - queryRes := tsdb.NewQueryResult() - queryRes.Series = append(queryRes.Series, getRandomWalk(query, context, 0)) - queryRes.ErrorString = "This is an error. It can include URLs http://grafana.com/" - return queryRes - }, - }) - - registerScenario(&Scenario{ - Id: "server_error_500", - Name: "Server Error (500)", - - Handler: func(query *tsdb.Query, context *tsdb.TsdbQuery) *tsdb.QueryResult { - panic("Test Data Panic!") - }, - }) - - registerScenario(&Scenario{ - Id: "logs", - Name: "Logs", - - Handler: func(query *tsdb.Query, context *tsdb.TsdbQuery) *tsdb.QueryResult { - from := context.TimeRange.GetFromAsMsEpoch() - to := context.TimeRange.GetToAsMsEpoch() - lines := query.Model.Get("lines").MustInt64(10) - includeLevelColumn := query.Model.Get("levelColumn").MustBool(false) - - logLevelGenerator := newRandomStringProvider([]string{ - "emerg", - "alert", - "crit", - "critical", - "warn", - "warning", - "err", - "eror", - "error", - "info", - "notice", - "dbug", - "debug", - "trace", - "", - }) - containerIDGenerator := newRandomStringProvider([]string{ - "f36a9eaa6d34310686f2b851655212023a216de955cbcc764210cefa71179b1a", - "5a354a630364f3742c602f315132e16def594fe68b1e4a195b2fce628e24c97a", - }) - hostnameGenerator := newRandomStringProvider([]string{ - "srv-001", - "srv-002", - }) - - table := tsdb.Table{ - Columns: []tsdb.TableColumn{ - {Text: "time"}, - {Text: "message"}, - {Text: "container_id"}, - {Text: "hostname"}, - }, - Rows: []tsdb.RowValues{}, - } - - if includeLevelColumn { - table.Columns = append(table.Columns, tsdb.TableColumn{Text: "level"}) - } - - for i := int64(0); i < lines && to > from; i++ { - row := tsdb.RowValues{float64(to)} - - logLevel := logLevelGenerator.Next() - timeFormatted := time.Unix(to/1000, 0).Format(time.RFC3339) - lvlString := "" - if !includeLevelColumn { - lvlString = fmt.Sprintf("lvl=%s ", logLevel) - } - - row = append(row, fmt.Sprintf("t=%s %smsg=\"Request Completed\" logger=context userId=1 orgId=1 uname=admin method=GET path=/api/datasources/proxy/152/api/prom/label status=502 remote_addr=[::1] time_ms=1 size=0 referer=\"http://localhost:3000/explore?left=%%5B%%22now-6h%%22,%%22now%%22,%%22Prometheus%%202.x%%22,%%7B%%7D,%%7B%%22ui%%22:%%5Btrue,true,true,%%22none%%22%%5D%%7D%%5D\"", timeFormatted, lvlString)) - row = append(row, containerIDGenerator.Next()) - row = append(row, hostnameGenerator.Next()) - - if includeLevelColumn { - row = append(row, logLevel) - } - - table.Rows = append(table.Rows, row) - to -= query.IntervalMs - } - - queryRes := tsdb.NewQueryResult() - queryRes.Tables = append(queryRes.Tables, &table) - return queryRes - }, - }) - - registerScenario(&Scenario{ - Id: "node_graph", - Name: "Node Graph", - // Data generated in JS - }) +func (p *testDataPlugin) registerScenario(scenario *Scenario) { + p.scenarios[scenario.ID] = scenario + p.queryMux.HandleFunc(scenario.ID, scenario.handler) } -// PredictablePulseDesc is the description for the Predictable Pulse scenerio. -const PredictablePulseDesc = `Predictable Pulse returns a pulse wave where there is a datapoint every timeStepSeconds. +func (p *testDataPlugin) registerScenarios() { + p.registerScenario(&Scenario{ + ID: string(exponentialHeatmapBucketDataQuery), + Name: "Exponential heatmap bucket data", + handler: p.handleExponentialHeatmapBucketDataScenario, + }) + + p.registerScenario(&Scenario{ + ID: string(linearHeatmapBucketDataQuery), + Name: "Linear heatmap bucket data", + handler: p.handleLinearHeatmapBucketDataScenario, + }) + + p.registerScenario(&Scenario{ + ID: string(randomWalkQuery), + Name: "Random Walk", + handler: p.handleRandomWalkScenario, + }) + + p.registerScenario(&Scenario{ + ID: string(predictablePulseQuery), + Name: "Predictable Pulse", + handler: p.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).` +Timestamps will line up evenly on timeStepSeconds (For example, 60 seconds means times will all end in :00 seconds).`, + }) -func getPredictablePulse(query *tsdb.Query, context *tsdb.TsdbQuery) *tsdb.QueryResult { - queryRes := tsdb.NewQueryResult() + p.registerScenario(&Scenario{ + ID: string(predictableCSVWaveQuery), + Name: "Predictable CSV Wave", + handler: p.handlePredictableCSVWaveScenario, + }) - // Process Input - var timeStep int64 - var onCount int64 - var offCount int64 - var onValue null.Float - var offValue null.Float + p.registerScenario(&Scenario{ + ID: string(randomWalkTableQuery), + Name: "Random Walk Table", + handler: p.handleRandomWalkTableScenario, + }) - options := query.Model.Get("pulseWave") + p.registerScenario(&Scenario{ + ID: string(randomWalkSlowQuery), + Name: "Slow Query", + StringInput: "5s", + handler: p.handleRandomWalkSlowScenario, + }) - var err error - if timeStep, err = options.Get("timeStep").Int64(); err != nil { - queryRes.Error = fmt.Errorf("failed to parse timeStep value '%v' into integer: %v", options.Get("timeStep"), err) - return queryRes - } - if onCount, err = options.Get("onCount").Int64(); err != nil { - queryRes.Error = fmt.Errorf("failed to parse onCount value '%v' into integer: %v", options.Get("onCount"), err) - return queryRes - } - if offCount, err = options.Get("offCount").Int64(); err != nil { - queryRes.Error = fmt.Errorf("failed to parse offCount value '%v' into integer: %v", options.Get("offCount"), err) - return queryRes + p.registerScenario(&Scenario{ + ID: string(noDataPointsQuery), + Name: "No Data Points", + handler: p.handleClientSideScenario, + }) + + p.registerScenario(&Scenario{ + ID: string(datapointsOutsideRangeQuery), + Name: "Datapoints Outside Range", + handler: p.handleDatapointsOutsideRangeScenario, + }) + + p.registerScenario(&Scenario{ + ID: string(manualEntryQuery), + Name: "Manual Entry", + handler: p.handleManualEntryScenario, + }) + + p.registerScenario(&Scenario{ + ID: string(csvMetricValuesQuery), + Name: "CSV Metric Values", + StringInput: "1,20,90,30,5,0", + handler: p.handleCSVMetricValuesScenario, + }) + + p.registerScenario(&Scenario{ + ID: string(streamingClientQuery), + Name: "Streaming Client", + handler: p.handleClientSideScenario, + }) + + p.registerScenario(&Scenario{ + ID: string(liveQuery), + Name: "Grafana Live", + handler: p.handleClientSideScenario, + }) + + p.registerScenario(&Scenario{ + ID: string(grafanaAPIQuery), + Name: "Grafana API", + handler: p.handleClientSideScenario, + }) + + p.registerScenario(&Scenario{ + ID: string(arrowQuery), + Name: "Load Apache Arrow Data", + handler: p.handleClientSideScenario, + }) + + p.registerScenario(&Scenario{ + ID: string(annotationsQuery), + Name: "Annotations", + handler: p.handleClientSideScenario, + }) + + p.registerScenario(&Scenario{ + ID: string(tableStaticQuery), + Name: "Table Static", + handler: p.handleTableStaticScenario, + }) + + p.registerScenario(&Scenario{ + ID: string(randomWalkWithErrorQuery), + Name: "Random Walk (with error)", + handler: p.handleRandomWalkWithErrorScenario, + }) + + p.registerScenario(&Scenario{ + ID: string(serverError500Query), + Name: "Server Error (500)", + handler: p.handleServerError500Scenario, + }) + + p.registerScenario(&Scenario{ + ID: string(logsQuery), + Name: "Logs", + handler: p.handleLogsScenario, + }) + + p.registerScenario(&Scenario{ + ID: string(nodeGraphQuery), + Name: "Node Graph", + }) + + p.queryMux.HandleFunc("", p.handleFallbackScenario) +} + +// 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) { + 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) + continue + } + + scenarioID := model.Get("scenarioId").MustString(string(randomWalkQuery)) + if _, exist := p.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) + } } - onValue, err = fromStringOrNumber(options.Get("onValue")) - if err != nil { - queryRes.Error = fmt.Errorf("failed to parse onValue value '%v' into float: %v", options.Get("onValue"), err) - return queryRes - } - offValue, err = fromStringOrNumber(options.Get("offValue")) - if err != nil { - queryRes.Error = fmt.Errorf("failed to parse offValue value '%v' into float: %v", options.Get("offValue"), err) - return queryRes - } - - timeStep *= 1000 // Seconds to Milliseconds - onFor := func(mod int64) (null.Float, error) { // How many items in the cycle should get the on value - var i int64 - for i = 0; i < onCount; i++ { - if mod == i*timeStep { - return onValue, nil + resp := backend.NewQueryDataResponse() + for scenarioID, queries := range scenarioQueries { + if scenario, exist := p.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) + } else { + for refID, dr := range sResp.Responses { + resp.Responses[refID] = dr + } } } - return offValue, nil - } - points, err := predictableSeries(context.TimeRange, timeStep, onCount+offCount, onFor) - if err != nil { - queryRes.Error = err - return queryRes } - series := newSeriesForQuery(query, 0) - series.Points = *points - series.Tags = parseLabels(query) - - queryRes.Series = append(queryRes.Series, series) - return queryRes + return resp, nil } -func getPredictableCSVWave(query *tsdb.Query, context *tsdb.TsdbQuery) *tsdb.QueryResult { - queryRes := tsdb.NewQueryResult() +func (p *testDataPlugin) handleRandomWalkScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) { + resp := backend.NewQueryDataResponse() - // Process Input - var timeStep int64 - - options := query.Model.Get("csvWave") - - var err error - if timeStep, err = options.Get("timeStep").Int64(); err != nil { - queryRes.Error = fmt.Errorf("failed to parse timeStep value '%v' into integer: %v", options.Get("timeStep"), err) - return queryRes - } - rawValues := options.Get("valuesCSV").MustString() - rawValues = strings.TrimRight(strings.TrimSpace(rawValues), ",") // Strip Trailing Comma - rawValesCSV := strings.Split(rawValues, ",") - values := make([]null.Float, len(rawValesCSV)) - for i, rawValue := range rawValesCSV { - val, err := null.FloatFromString(strings.TrimSpace(rawValue), "null") + for _, q := range req.Queries { + model, err := simplejson.NewJson(q.JSON) if err != nil { - queryRes.Error = errutil.Wrapf(err, "failed to parse value '%v' into nullable float", rawValue) - return queryRes + continue + } + seriesCount := model.Get("seriesCount").MustInt(1) + + for i := 0; i < seriesCount; i++ { + respD := resp.Responses[q.RefID] + respD.Frames = append(respD.Frames, randomWalk(q, model, i)) + resp.Responses[q.RefID] = respD } - values[i] = val } - timeStep *= 1000 // Seconds to Milliseconds - valuesLen := int64(len(values)) - getValue := func(mod int64) (null.Float, error) { - var i int64 - for i = 0; i < valuesLen; i++ { - if mod == i*timeStep { - return values[i], nil + return resp, nil +} + +func (p *testDataPlugin) handleDatapointsOutsideRangeScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) { + resp := backend.NewQueryDataResponse() + + for _, q := range req.Queries { + model, err := simplejson.NewJson(q.JSON) + if err != nil { + continue + } + + frame := newSeriesForQuery(q, model, 0) + outsideTime := q.TimeRange.From.Add(-1 * time.Hour) + frame.Fields = data.Fields{ + data.NewField("time", nil, []time.Time{outsideTime}), + data.NewField("value", nil, []float64{10}), + } + + respD := resp.Responses[q.RefID] + respD.Frames = append(respD.Frames, frame) + resp.Responses[q.RefID] = respD + } + + return resp, nil +} + +func (p *testDataPlugin) handleManualEntryScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) { + resp := backend.NewQueryDataResponse() + + for _, q := range req.Queries { + model, err := simplejson.NewJson(q.JSON) + if err != nil { + continue + } + points := model.Get("points").MustArray() + + frame := newSeriesForQuery(q, model, 0) + startTime := q.TimeRange.From.UnixNano() / int64(time.Millisecond) + endTime := q.TimeRange.To.UnixNano() / int64(time.Millisecond) + + timeVec := make([]*time.Time, 0) + floatVec := make([]*float64, 0) + + for _, val := range points { + pointValues := val.([]interface{}) + + var value float64 + + if valueFloat, err := strconv.ParseFloat(string(pointValues[0].(json.Number)), 64); err == nil { + value = valueFloat + } + + timeInt, err := strconv.ParseInt(string(pointValues[1].(json.Number)), 10, 64) + if err != nil { + continue + } + t := time.Unix(timeInt/int64(1e+3), (timeInt%int64(1e+3))*int64(1e+6)) + + if timeInt >= startTime && timeInt <= endTime { + timeVec = append(timeVec, &t) + floatVec = append(floatVec, &value) } } - return null.Float{}, fmt.Errorf("did not get value at point in waveform - should not be here") - } - points, err := predictableSeries(context.TimeRange, timeStep, valuesLen, getValue) - if err != nil { - queryRes.Error = err - return queryRes - } - series := newSeriesForQuery(query, 0) - series.Points = *points - series.Tags = parseLabels(query) - - queryRes.Series = append(queryRes.Series, series) - return queryRes -} - -func predictableSeries(timeRange *tsdb.TimeRange, timeStep, length int64, getValue func(mod int64) (null.Float, error)) (*tsdb.TimeSeriesPoints, error) { - points := make(tsdb.TimeSeriesPoints, 0) - - from := timeRange.GetFromAsMsEpoch() - to := timeRange.GetToAsMsEpoch() - - timeCursor := from - (from % timeStep) // Truncate Start - wavePeriod := timeStep * length - maxPoints := 10000 // Don't return too many points - - for i := 0; i < maxPoints && timeCursor < to; i++ { - val, err := getValue(timeCursor % wavePeriod) - if err != nil { - return &points, err + frame.Fields = data.Fields{ + data.NewField("time", nil, timeVec), + data.NewField("value", nil, floatVec), } - point := tsdb.NewTimePoint(val, float64(timeCursor)) - points = append(points, point) - timeCursor += timeStep + + respD := resp.Responses[q.RefID] + respD.Frames = append(respD.Frames, frame) + resp.Responses[q.RefID] = respD } - return &points, nil + + return resp, nil } -func getRandomWalk(query *tsdb.Query, tsdbQuery *tsdb.TsdbQuery, index int) *tsdb.TimeSeries { - timeWalkerMs := tsdbQuery.TimeRange.GetFromAsMsEpoch() - to := tsdbQuery.TimeRange.GetToAsMsEpoch() - series := newSeriesForQuery(query, index) +func (p *testDataPlugin) handleCSVMetricValuesScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) { + resp := backend.NewQueryDataResponse() - startValue := query.Model.Get("startValue").MustFloat64(rand.Float64() * 100) - spread := query.Model.Get("spread").MustFloat64(1) - noise := query.Model.Get("noise").MustFloat64(0) + for _, q := range req.Queries { + model, err := simplejson.NewJson(q.JSON) + if err != nil { + continue + } - min, err := query.Model.Get("min").Float64() + stringInput := model.Get("stringInput").MustString() + stringInput = strings.ReplaceAll(stringInput, " ", "") + + var values []*float64 + for _, strVal := range strings.Split(stringInput, ",") { + if strVal == "null" { + values = append(values, nil) + } + if val, err := strconv.ParseFloat(strVal, 64); err == nil { + values = append(values, &val) + } + } + + if len(values) == 0 { + return resp, nil + } + + frame := data.NewFrame("", + data.NewField("time", nil, []*time.Time{}), + data.NewField(frameNameForQuery(q, model, 0), nil, []*float64{})) + startTime := q.TimeRange.From.UnixNano() / int64(time.Millisecond) + endTime := q.TimeRange.To.UnixNano() / int64(time.Millisecond) + var step int64 = 0 + if len(values) > 1 { + step = (endTime - startTime) / int64(len(values)-1) + } + + for _, val := range values { + t := time.Unix(startTime/int64(1e+3), (startTime%int64(1e+3))*int64(1e+6)) + frame.AppendRow(&t, val) + startTime += step + } + + respD := resp.Responses[q.RefID] + respD.Frames = append(respD.Frames, frame) + resp.Responses[q.RefID] = respD + } + + return resp, nil +} + +func (p *testDataPlugin) handleRandomWalkWithErrorScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) { + resp := backend.NewQueryDataResponse() + + for _, q := range req.Queries { + model, err := simplejson.NewJson(q.JSON) + if err != nil { + continue + } + + respD := resp.Responses[q.RefID] + respD.Frames = append(respD.Frames, randomWalk(q, model, 0)) + respD.Error = fmt.Errorf("this is an error and it can include URLs http://grafana.com/") + resp.Responses[q.RefID] = respD + } + + return resp, nil +} + +func (p *testDataPlugin) handleRandomWalkSlowScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) { + resp := backend.NewQueryDataResponse() + + for _, q := range req.Queries { + model, err := simplejson.NewJson(q.JSON) + if err != nil { + continue + } + + stringInput := model.Get("stringInput").MustString() + parsedInterval, _ := time.ParseDuration(stringInput) + time.Sleep(parsedInterval) + + respD := resp.Responses[q.RefID] + respD.Frames = append(respD.Frames, randomWalk(q, model, 0)) + resp.Responses[q.RefID] = respD + } + + return resp, nil +} + +func (p *testDataPlugin) handleRandomWalkTableScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) { + resp := backend.NewQueryDataResponse() + + for _, q := range req.Queries { + model, err := simplejson.NewJson(q.JSON) + if err != nil { + continue + } + + respD := resp.Responses[q.RefID] + respD.Frames = append(respD.Frames, randomWalkTable(q, model)) + resp.Responses[q.RefID] = respD + } + + return resp, nil +} + +func (p *testDataPlugin) handlePredictableCSVWaveScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) { + resp := backend.NewQueryDataResponse() + + for _, q := range req.Queries { + model, err := simplejson.NewJson(q.JSON) + if err != nil { + continue + } + + respD := resp.Responses[q.RefID] + frame, err := predictableCSVWave(q, model) + if err != nil { + continue + } + respD.Frames = append(respD.Frames, frame) + resp.Responses[q.RefID] = respD + } + + return resp, nil +} + +func (p *testDataPlugin) handlePredictablePulseScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) { + resp := backend.NewQueryDataResponse() + + for _, q := range req.Queries { + model, err := simplejson.NewJson(q.JSON) + if err != nil { + continue + } + + respD := resp.Responses[q.RefID] + frame, err := predictablePulse(q, model) + if err != nil { + continue + } + respD.Frames = append(respD.Frames, frame) + resp.Responses[q.RefID] = respD + } + + return resp, nil +} + +func (p *testDataPlugin) 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) { + return backend.NewQueryDataResponse(), nil +} + +func (p *testDataPlugin) handleExponentialHeatmapBucketDataScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) { + resp := backend.NewQueryDataResponse() + + for _, q := range req.Queries { + respD := resp.Responses[q.RefID] + frame := randomHeatmapData(q, func(index int) float64 { + return math.Exp2(float64(index)) + }) + respD.Frames = append(respD.Frames, frame) + resp.Responses[q.RefID] = respD + } + + return resp, nil +} + +func (p *testDataPlugin) handleLinearHeatmapBucketDataScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) { + resp := backend.NewQueryDataResponse() + + for _, q := range req.Queries { + respD := resp.Responses[q.RefID] + frame := randomHeatmapData(q, func(index int) float64 { + return float64(index * 10) + }) + respD.Frames = append(respD.Frames, frame) + resp.Responses[q.RefID] = respD + } + + return resp, nil +} + +func (p *testDataPlugin) handleTableStaticScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) { + resp := backend.NewQueryDataResponse() + + for _, q := range req.Queries { + timeWalkerMs := q.TimeRange.From.UnixNano() / int64(time.Millisecond) + to := q.TimeRange.To.UnixNano() / int64(time.Millisecond) + step := q.Interval.Milliseconds() + + frame := data.NewFrame(q.RefID, + data.NewField("Time", nil, []time.Time{}), + data.NewField("Message", nil, []string{}), + data.NewField("Description", nil, []string{}), + data.NewField("Value", nil, []float64{}), + ) + + for i := int64(0); i < 10 && timeWalkerMs < to; i++ { + t := time.Unix(timeWalkerMs/int64(1e+3), (timeWalkerMs%int64(1e+3))*int64(1e+6)) + frame.AppendRow(t, "This is a message", "Description", 23.1) + timeWalkerMs += step + } + + respD := resp.Responses[q.RefID] + respD.Frames = append(respD.Frames, frame) + resp.Responses[q.RefID] = respD + } + + return resp, nil +} + +func (p *testDataPlugin) handleLogsScenario(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) { + resp := backend.NewQueryDataResponse() + + for _, q := range req.Queries { + from := q.TimeRange.From.UnixNano() / int64(time.Millisecond) + to := q.TimeRange.To.UnixNano() / int64(time.Millisecond) + + model, err := simplejson.NewJson(q.JSON) + if err != nil { + continue + } + + lines := model.Get("lines").MustInt64(10) + includeLevelColumn := model.Get("levelColumn").MustBool(false) + + logLevelGenerator := newRandomStringProvider([]string{ + "emerg", + "alert", + "crit", + "critical", + "warn", + "warning", + "err", + "eror", + "error", + "info", + "notice", + "dbug", + "debug", + "trace", + "", + }) + containerIDGenerator := newRandomStringProvider([]string{ + "f36a9eaa6d34310686f2b851655212023a216de955cbcc764210cefa71179b1a", + "5a354a630364f3742c602f315132e16def594fe68b1e4a195b2fce628e24c97a", + }) + hostnameGenerator := newRandomStringProvider([]string{ + "srv-001", + "srv-002", + }) + + frame := data.NewFrame(q.RefID, + data.NewField("time", nil, []time.Time{}), + data.NewField("message", nil, []string{}), + data.NewField("container_id", nil, []string{}), + data.NewField("hostname", nil, []string{}), + ).SetMeta(&data.FrameMeta{ + PreferredVisualization: "logs", + }) + + if includeLevelColumn { + frame.Fields = append(frame.Fields, data.NewField("level", nil, []string{})) + } + + for i := int64(0); i < lines && to > from; i++ { + logLevel := logLevelGenerator.Next() + timeFormatted := time.Unix(to/1000, 0).Format(time.RFC3339) + lvlString := "" + if !includeLevelColumn { + lvlString = fmt.Sprintf("lvl=%s ", logLevel) + } + + message := fmt.Sprintf("t=%s %smsg=\"Request Completed\" logger=context userId=1 orgId=1 uname=admin method=GET path=/api/datasources/proxy/152/api/prom/label status=502 remote_addr=[::1] time_ms=1 size=0 referer=\"http://localhost:3000/explore?left=%%5B%%22now-6h%%22,%%22now%%22,%%22Prometheus%%202.x%%22,%%7B%%7D,%%7B%%22ui%%22:%%5Btrue,true,true,%%22none%%22%%5D%%7D%%5D\"", timeFormatted, lvlString) + containerID := containerIDGenerator.Next() + hostname := hostnameGenerator.Next() + + t := time.Unix(to/int64(1e+3), (to%int64(1e+3))*int64(1e+6)) + + if includeLevelColumn { + frame.AppendRow(t, message, containerID, hostname, logLevel) + } else { + frame.AppendRow(t, message, containerID, hostname) + } + + to -= q.Interval.Milliseconds() + } + + respD := resp.Responses[q.RefID] + respD.Frames = append(respD.Frames, frame) + resp.Responses[q.RefID] = respD + } + + return resp, nil +} + +func randomWalk(query backend.DataQuery, model *simplejson.Json, index int) *data.Frame { + timeWalkerMs := query.TimeRange.From.UnixNano() / int64(time.Millisecond) + to := query.TimeRange.To.UnixNano() / int64(time.Millisecond) + startValue := model.Get("startValue").MustFloat64(rand.Float64() * 100) + spread := model.Get("spread").MustFloat64(1) + noise := model.Get("noise").MustFloat64(0) + + min, err := model.Get("min").Float64() hasMin := err == nil - max, err := query.Model.Get("max").Float64() + max, err := model.Get("max").Float64() hasMax := err == nil - points := make(tsdb.TimeSeriesPoints, 0) + timeVec := make([]*time.Time, 0) + floatVec := make([]*float64, 0) + walker := startValue for i := int64(0); i < 10000 && timeWalkerMs < to; i++ { @@ -601,66 +660,35 @@ func getRandomWalk(query *tsdb.Query, tsdbQuery *tsdb.TsdbQuery, index int) *tsd walker = max } - points = append(points, tsdb.NewTimePoint(null.FloatFrom(nextValue), float64(timeWalkerMs))) + t := time.Unix(timeWalkerMs/int64(1e+3), (timeWalkerMs%int64(1e+3))*int64(1e+6)) + timeVec = append(timeVec, &t) + floatVec = append(floatVec, &nextValue) walker += (rand.Float64() - 0.5) * spread - timeWalkerMs += query.IntervalMs + timeWalkerMs += query.Interval.Milliseconds() } - series.Points = points - series.Tags = parseLabels(query) - return series + return data.NewFrame("", + data.NewField("time", nil, timeVec), + data.NewField(frameNameForQuery(query, model, 0), parseLabels(model), floatVec), + ) } -/** - * Looks for a labels request and adds them as tags - * - * '{job="foo", instance="bar"} => {job: "foo", instance: "bar"}` - */ -func parseLabels(query *tsdb.Query) map[string]string { - tags := map[string]string{} - - labelText := query.Model.Get("labels").MustString("") - if labelText == "" { - return map[string]string{} - } - - text := strings.Trim(labelText, `{}`) - if len(text) < 2 { - return tags - } - - tags = make(map[string]string) - - for _, keyval := range strings.Split(text, ",") { - idx := strings.Index(keyval, "=") - key := strings.TrimSpace(keyval[:idx]) - val := strings.TrimSpace(keyval[idx+1:]) - val = strings.Trim(val, "\"") - tags[key] = val - } - - return tags -} - -func getRandomWalkTable(query *tsdb.Query, tsdbQuery *tsdb.TsdbQuery) *tsdb.QueryResult { - timeWalkerMs := tsdbQuery.TimeRange.GetFromAsMsEpoch() - to := tsdbQuery.TimeRange.GetToAsMsEpoch() - - table := tsdb.Table{ - Columns: []tsdb.TableColumn{ - {Text: "Time"}, - {Text: "Value"}, - {Text: "Min"}, - {Text: "Max"}, - {Text: "Info"}, - }, - Rows: []tsdb.RowValues{}, - } - - withNil := query.Model.Get("withNil").MustBool(false) - walker := query.Model.Get("startValue").MustFloat64(rand.Float64() * 100) +func randomWalkTable(query backend.DataQuery, model *simplejson.Json) *data.Frame { + timeWalkerMs := query.TimeRange.From.UnixNano() / int64(time.Millisecond) + to := query.TimeRange.To.UnixNano() / int64(time.Millisecond) + withNil := model.Get("withNil").MustBool(false) + walker := model.Get("startValue").MustFloat64(rand.Float64() * 100) spread := 2.5 + + frame := data.NewFrame(query.RefID, + data.NewField("Time", nil, []*time.Time{}), + data.NewField("Value", nil, []*float64{}), + data.NewField("Min", nil, []*float64{}), + data.NewField("Max", nil, []*float64{}), + data.NewField("Info", nil, []*string{}), + ) + var info strings.Builder for i := int64(0); i < query.MaxDataPoints && timeWalkerMs < to; i++ { @@ -676,37 +704,181 @@ func getRandomWalkTable(query *tsdb.Query, tsdbQuery *tsdb.TsdbQuery) *tsdb.Quer if math.Abs(delta) > .4 { info.WriteString(" fast") } - row := tsdb.RowValues{ - float64(timeWalkerMs), - walker, - walker - ((rand.Float64() * spread) + 0.01), // Min - walker + ((rand.Float64() * spread) + 0.01), // Max - info.String(), - } + t := time.Unix(timeWalkerMs/int64(1e+3), (timeWalkerMs%int64(1e+3))*int64(1e+6)) + val := walker + min := walker - ((rand.Float64() * spread) + 0.01) + max := walker + ((rand.Float64() * spread) + 0.01) + infoString := info.String() + + vals := []*float64{&val, &min, &max} // Add some random null values if withNil && rand.Float64() > 0.8 { - for i := 1; i < 4; i++ { + for i := range vals { if rand.Float64() > .2 { - row[i] = nil + vals[i] = nil } } } - table.Rows = append(table.Rows, row) - timeWalkerMs += query.IntervalMs + frame.AppendRow(&t, vals[0], vals[1], vals[2], &infoString) + + timeWalkerMs += query.Interval.Milliseconds() } - queryRes := tsdb.NewQueryResult() - queryRes.Tables = append(queryRes.Tables, &table) - return queryRes + + return frame } -func registerScenario(scenario *Scenario) { - ScenarioRegistry[scenario.Id] = scenario +func predictableCSVWave(query backend.DataQuery, model *simplejson.Json) (*data.Frame, error) { + options := model.Get("csvWave") + + var timeStep int64 + var err error + if timeStep, err = options.Get("timeStep").Int64(); err != nil { + return nil, fmt.Errorf("failed to parse timeStep value '%v' into integer: %v", options.Get("timeStep"), err) + } + rawValues := options.Get("valuesCSV").MustString() + rawValues = strings.TrimRight(strings.TrimSpace(rawValues), ",") // Strip Trailing Comma + rawValesCSV := strings.Split(rawValues, ",") + values := make([]null.Float, len(rawValesCSV)) + for i, rawValue := range rawValesCSV { + val, err := null.FloatFromString(strings.TrimSpace(rawValue), "null") + if err != nil { + return nil, errutil.Wrapf(err, "failed to parse value '%v' into nullable float", rawValue) + } + values[i] = val + } + + timeStep *= 1000 // Seconds to Milliseconds + valuesLen := int64(len(values)) + getValue := func(mod int64) (null.Float, error) { + var i int64 + for i = 0; i < valuesLen; i++ { + if mod == i*timeStep { + return values[i], nil + } + } + return null.Float{}, fmt.Errorf("did not get value at point in waveform - should not be here") + } + fields, err := predictableSeries(query.TimeRange, timeStep, valuesLen, getValue) + if err != nil { + return nil, err + } + + frame := newSeriesForQuery(query, model, 0) + frame.Fields = fields + frame.Fields[1].Labels = parseLabels(model) + + return frame, nil } -func newSeriesForQuery(query *tsdb.Query, index int) *tsdb.TimeSeries { - alias := query.Model.Get("alias").MustString("") +func predictableSeries(timeRange backend.TimeRange, timeStep, length int64, getValue func(mod int64) (null.Float, error)) (data.Fields, error) { + from := timeRange.From.UnixNano() / int64(time.Millisecond) + to := timeRange.To.UnixNano() / int64(time.Millisecond) + + timeCursor := from - (from % timeStep) // Truncate Start + wavePeriod := timeStep * length + maxPoints := 10000 // Don't return too many points + + timeVec := make([]*time.Time, 0) + floatVec := make([]*float64, 0) + + for i := 0; i < maxPoints && timeCursor < to; i++ { + val, err := getValue(timeCursor % wavePeriod) + if err != nil { + return nil, err + } + + t := time.Unix(timeCursor/int64(1e+3), (timeCursor%int64(1e+3))*int64(1e+6)) + timeVec = append(timeVec, &t) + floatVec = append(floatVec, &val.Float64) + + timeCursor += timeStep + } + + return data.Fields{ + data.NewField("time", nil, timeVec), + data.NewField("value", nil, floatVec), + }, nil +} + +func predictablePulse(query backend.DataQuery, model *simplejson.Json) (*data.Frame, error) { + // Process Input + var timeStep int64 + var onCount int64 + var offCount int64 + var onValue null.Float + var offValue null.Float + + options := model.Get("pulseWave") + + var err error + if timeStep, err = options.Get("timeStep").Int64(); err != nil { + return nil, fmt.Errorf("failed to parse timeStep value '%v' into integer: %v", options.Get("timeStep"), err) + } + if onCount, err = options.Get("onCount").Int64(); err != nil { + return nil, fmt.Errorf("failed to parse onCount value '%v' into integer: %v", options.Get("onCount"), err) + } + if offCount, err = options.Get("offCount").Int64(); err != nil { + return nil, fmt.Errorf("failed to parse offCount value '%v' into integer: %v", options.Get("offCount"), err) + } + + onValue, err = fromStringOrNumber(options.Get("onValue")) + if err != nil { + return nil, fmt.Errorf("failed to parse onValue value '%v' into float: %v", options.Get("onValue"), err) + } + offValue, err = fromStringOrNumber(options.Get("offValue")) + if err != nil { + return nil, fmt.Errorf("failed to parse offValue value '%v' into float: %v", options.Get("offValue"), err) + } + + timeStep *= 1000 // Seconds to Milliseconds + onFor := func(mod int64) (null.Float, error) { // How many items in the cycle should get the on value + var i int64 + for i = 0; i < onCount; i++ { + if mod == i*timeStep { + return onValue, nil + } + } + return offValue, nil + } + fields, err := predictableSeries(query.TimeRange, timeStep, onCount+offCount, onFor) + if err != nil { + return nil, err + } + + frame := newSeriesForQuery(query, model, 0) + frame.Fields = fields + frame.Fields[1].Labels = parseLabels(model) + + return frame, nil +} + +func randomHeatmapData(query backend.DataQuery, fnBucketGen func(index int) float64) *data.Frame { + frame := data.NewFrame("data", data.NewField("time", nil, []*time.Time{})) + for i := 0; i < 10; i++ { + frame.Fields = append(frame.Fields, data.NewField(strconv.FormatInt(int64(fnBucketGen(i)), 10), nil, []*float64{})) + } + + timeWalkerMs := query.TimeRange.From.UnixNano() / int64(time.Millisecond) + to := query.TimeRange.To.UnixNano() / int64(time.Millisecond) + + for j := int64(0); j < 100 && timeWalkerMs < to; j++ { + t := time.Unix(timeWalkerMs/int64(1e+3), (timeWalkerMs%int64(1e+3))*int64(1e+6)) + vals := []interface{}{&t} + for n := 1; n < len(frame.Fields); n++ { + v := float64(rand.Int63n(100)) + vals = append(vals, &v) + } + frame.AppendRow(vals...) + timeWalkerMs += query.Interval.Milliseconds() * 50 + } + + return frame +} + +func newSeriesForQuery(query backend.DataQuery, model *simplejson.Json, index int) *data.Frame { + alias := model.Get("alias").MustString("") suffix := "" if index > 0 { @@ -714,7 +886,7 @@ func newSeriesForQuery(query *tsdb.Query, index int) *tsdb.TimeSeries { } if alias == "" { - alias = fmt.Sprintf("%s-series%s", query.RefId, suffix) + alias = fmt.Sprintf("%s-series%s", query.RefID, suffix) } if alias == "__server_names" && len(serverNames) > index { @@ -725,7 +897,61 @@ func newSeriesForQuery(query *tsdb.Query, index int) *tsdb.TimeSeries { alias = houseLocations[index] } - return &tsdb.TimeSeries{Name: alias} + return data.NewFrame(alias) +} + +/** + * Looks for a labels request and adds them as tags + * + * '{job="foo", instance="bar"} => {job: "foo", instance: "bar"}` + */ +func parseLabels(model *simplejson.Json) data.Labels { + tags := data.Labels{} + + labelText := model.Get("labels").MustString("") + if labelText == "" { + return data.Labels{} + } + + text := strings.Trim(labelText, `{}`) + if len(text) < 2 { + return tags + } + + tags = make(data.Labels) + + for _, keyval := range strings.Split(text, ",") { + idx := strings.Index(keyval, "=") + key := strings.TrimSpace(keyval[:idx]) + val := strings.TrimSpace(keyval[idx+1:]) + val = strings.Trim(val, "\"") + tags[key] = val + } + + return tags +} + +func frameNameForQuery(query backend.DataQuery, model *simplejson.Json, index int) string { + name := model.Get("alias").MustString("") + suffix := "" + + if index > 0 { + suffix = strconv.Itoa(index) + } + + if name == "" { + name = fmt.Sprintf("%s-series%s", query.RefID, suffix) + } + + if name == "__server_names" && len(serverNames) > index { + name = serverNames[index] + } + + if name == "__house_locations" && len(houseLocations) > index { + name = houseLocations[index] + } + + return name } func fromStringOrNumber(val *simplejson.Json) (null.Float, error) { diff --git a/pkg/tsdb/testdatasource/scenarios_test.go b/pkg/tsdb/testdatasource/scenarios_test.go index 9c6cd35604d..b4da6905c6e 100644 --- a/pkg/tsdb/testdatasource/scenarios_test.go +++ b/pkg/tsdb/testdatasource/scenarios_test.go @@ -1,55 +1,115 @@ package testdatasource import ( + "context" + "fmt" "testing" "time" + "github.com/grafana/grafana-plugin-sdk-go/backend" + "github.com/grafana/grafana-plugin-sdk-go/data" "github.com/grafana/grafana/pkg/components/simplejson" "github.com/grafana/grafana/pkg/tsdb" + "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) func TestTestdataScenarios(t *testing.T) { + p := &testDataPlugin{} + t.Run("random walk ", func(t *testing.T) { - scenario := ScenarioRegistry["random_walk"] - t.Run("Should start at the requested value", func(t *testing.T) { - req := &tsdb.TsdbQuery{ - TimeRange: tsdb.NewFakeTimeRange("5m", "now", time.Now()), - Queries: []*tsdb.Query{ - {RefId: "A", IntervalMs: 100, MaxDataPoints: 100, Model: simplejson.New()}, + timeRange := tsdb.NewFakeTimeRange("5m", "now", time.Now()) + + model := simplejson.New() + model.Set("startValue", 1.234) + modelBytes, err := model.MarshalJSON() + require.NoError(t, err) + + query := backend.DataQuery{ + RefID: "A", + TimeRange: backend.TimeRange{ + From: timeRange.MustGetFrom(), + To: timeRange.MustGetTo(), }, + Interval: 100 * time.Millisecond, + MaxDataPoints: 100, + JSON: modelBytes, } - query := req.Queries[0] - query.Model.Set("startValue", 1.234) - result := scenario.Handler(req.Queries[0], req) - require.NotNil(t, result.Series) + req := &backend.QueryDataRequest{ + PluginContext: backend.PluginContext{}, + Queries: []backend.DataQuery{query}, + } - points := result.Series[0].Points - require.Equal(t, 1.234, points[0][0].Float64) + resp, err := p.handleRandomWalkScenario(context.Background(), req) + require.NoError(t, err) + require.NotNil(t, resp) + + dResp, exists := resp.Responses[query.RefID] + require.True(t, exists) + require.NoError(t, dResp.Error) + + require.Len(t, dResp.Frames, 1) + frame := dResp.Frames[0] + require.Len(t, frame.Fields, 2) + require.Equal(t, "time", frame.Fields[0].Name) + require.Equal(t, "A-series", frame.Fields[1].Name) + val, ok := frame.Fields[1].ConcreteAt(0) + require.True(t, ok) + require.Equal(t, 1.234, val) }) }) t.Run("random walk table", func(t *testing.T) { - scenario := ScenarioRegistry["random_walk_table"] - t.Run("Should return a table that looks like value/min/max", func(t *testing.T) { - req := &tsdb.TsdbQuery{ - TimeRange: tsdb.NewFakeTimeRange("5m", "now", time.Now()), - Queries: []*tsdb.Query{ - {RefId: "A", IntervalMs: 100, MaxDataPoints: 100, Model: simplejson.New()}, + timeRange := tsdb.NewFakeTimeRange("5m", "now", time.Now()) + + model := simplejson.New() + modelBytes, err := model.MarshalJSON() + require.NoError(t, err) + + query := backend.DataQuery{ + RefID: "A", + TimeRange: backend.TimeRange{ + From: timeRange.MustGetFrom(), + To: timeRange.MustGetTo(), }, + Interval: 100 * time.Millisecond, + MaxDataPoints: 100, + JSON: modelBytes, } - result := scenario.Handler(req.Queries[0], req) - table := result.Tables[0] + req := &backend.QueryDataRequest{ + PluginContext: backend.PluginContext{}, + Queries: []backend.DataQuery{query}, + } - require.Greater(t, len(table.Rows), 50) - for _, row := range table.Rows { - value := row[1] - min := row[2] - max := row[3] + resp, err := p.handleRandomWalkTableScenario(context.Background(), req) + require.NoError(t, err) + require.NotNil(t, resp) + + dResp, exists := resp.Responses[query.RefID] + require.True(t, exists) + require.NoError(t, dResp.Error) + + require.Len(t, dResp.Frames, 1) + frame := dResp.Frames[0] + require.Greater(t, frame.Rows(), 50) + require.Len(t, frame.Fields, 5) + require.Equal(t, "Time", frame.Fields[0].Name) + require.Equal(t, "Value", frame.Fields[1].Name) + require.Equal(t, "Min", frame.Fields[2].Name) + require.Equal(t, "Max", frame.Fields[3].Name) + require.Equal(t, "Info", frame.Fields[4].Name) + + for i := 0; i < frame.Rows(); i++ { + value, ok := frame.ConcreteAt(1, i) + require.True(t, ok) + min, ok := frame.ConcreteAt(2, i) + require.True(t, ok) + max, ok := frame.ConcreteAt(3, i) + require.True(t, ok) require.Less(t, min, value) require.Greater(t, max, value) @@ -57,66 +117,98 @@ func TestTestdataScenarios(t *testing.T) { }) t.Run("Should return a table with some nil values", func(t *testing.T) { - req := &tsdb.TsdbQuery{ - TimeRange: tsdb.NewFakeTimeRange("5m", "now", time.Now()), - Queries: []*tsdb.Query{ - {RefId: "A", IntervalMs: 100, MaxDataPoints: 100, Model: simplejson.New()}, + timeRange := tsdb.NewFakeTimeRange("5m", "now", time.Now()) + + model := simplejson.New() + model.Set("withNil", true) + + modelBytes, err := model.MarshalJSON() + require.NoError(t, err) + + query := backend.DataQuery{ + RefID: "A", + TimeRange: backend.TimeRange{ + From: timeRange.MustGetFrom(), + To: timeRange.MustGetTo(), }, + Interval: 100 * time.Millisecond, + MaxDataPoints: 100, + JSON: modelBytes, } - query := req.Queries[0] - query.Model.Set("withNil", true) - result := scenario.Handler(req.Queries[0], req) - table := result.Tables[0] + req := &backend.QueryDataRequest{ + PluginContext: backend.PluginContext{}, + Queries: []backend.DataQuery{query}, + } - nil1 := false - nil2 := false - nil3 := false + resp, err := p.handleRandomWalkTableScenario(context.Background(), req) + require.NoError(t, err) + require.NotNil(t, resp) - require.Greater(t, len(table.Rows), 50) - for _, row := range table.Rows { - if row[1] == nil { - nil1 = true + dResp, exists := resp.Responses[query.RefID] + require.True(t, exists) + require.NoError(t, dResp.Error) + + require.Len(t, dResp.Frames, 1) + frame := dResp.Frames[0] + require.Greater(t, frame.Rows(), 50) + require.Len(t, frame.Fields, 5) + require.Equal(t, "Time", frame.Fields[0].Name) + require.Equal(t, "Value", frame.Fields[1].Name) + require.Equal(t, "Min", frame.Fields[2].Name) + require.Equal(t, "Max", frame.Fields[3].Name) + require.Equal(t, "Info", frame.Fields[4].Name) + + valNil := false + minNil := false + maxNil := false + + for i := 0; i < frame.Rows(); i++ { + _, ok := frame.ConcreteAt(1, i) + if !ok { + valNil = true } - if row[2] == nil { - nil2 = true + + _, ok = frame.ConcreteAt(2, i) + if !ok { + minNil = true } - if row[3] == nil { - nil3 = true + + _, ok = frame.ConcreteAt(3, i) + if !ok { + maxNil = true } } - require.True(t, nil1) - require.True(t, nil2) - require.True(t, nil3) + require.True(t, valNil) + require.True(t, minNil) + require.True(t, maxNil) }) }) } func TestParseLabels(t *testing.T) { - expectedTags := map[string]string{ + expectedTags := data.Labels{ "job": "foo", "instance": "bar", } - query1 := tsdb.Query{ - Model: simplejson.NewFromAny(map[string]interface{}{ + tcs := []struct { + model map[string]interface{} + }{ + {model: map[string]interface{}{ "labels": `{job="foo", instance="bar"}`, - }), - } - require.Equal(t, expectedTags, parseLabels(&query1)) - - query2 := tsdb.Query{ - Model: simplejson.NewFromAny(map[string]interface{}{ + }}, + {model: map[string]interface{}{ "labels": `job=foo, instance=bar`, - }), - } - require.Equal(t, expectedTags, parseLabels(&query2)) - - query3 := tsdb.Query{ - Model: simplejson.NewFromAny(map[string]interface{}{ + }}, + {model: map[string]interface{}{ "labels": `job = foo,instance = bar`, - }), + }}, + } + + for i, tc := range tcs { + model := simplejson.NewFromAny(tc.model) + assert.Equal(t, expectedTags, parseLabels(model), fmt.Sprintf("Actual tags in test case %d doesn't match expected tags", i+1)) } - require.Equal(t, expectedTags, parseLabels(&query3)) } diff --git a/pkg/tsdb/testdatasource/testdata.go b/pkg/tsdb/testdatasource/testdata.go index 26a0df51d8d..72d722c78bb 100644 --- a/pkg/tsdb/testdatasource/testdata.go +++ b/pkg/tsdb/testdatasource/testdata.go @@ -1,42 +1,42 @@ package testdatasource import ( - "context" + "net/http" + "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/pkg/infra/log" - "github.com/grafana/grafana/pkg/models" - "github.com/grafana/grafana/pkg/tsdb" + "github.com/grafana/grafana/pkg/plugins/backendplugin" + "github.com/grafana/grafana/pkg/plugins/backendplugin/coreplugin" + "github.com/grafana/grafana/pkg/registry" ) -type TestDataExecutor struct { - *models.DataSource - log log.Logger -} - -func NewTestDataExecutor(dsInfo *models.DataSource) (tsdb.TsdbQueryEndpoint, error) { - return &TestDataExecutor{ - DataSource: dsInfo, - log: log.New("tsdb.testdata"), - }, nil -} - func init() { - tsdb.RegisterTsdbQueryEndpoint("testdata", NewTestDataExecutor) + registry.RegisterService(&testDataPlugin{}) } -func (e *TestDataExecutor) Query(ctx context.Context, dsInfo *models.DataSource, tsdbQuery *tsdb.TsdbQuery) (*tsdb.Response, error) { - result := &tsdb.Response{} - result.Results = make(map[string]*tsdb.QueryResult) +type testDataPlugin struct { + BackendPluginManager backendplugin.Manager `inject:""` + logger log.Logger + scenarios map[string]*Scenario + queryMux *datasource.QueryTypeMux +} - for _, query := range tsdbQuery.Queries { - scenarioId := query.Model.Get("scenarioId").MustString("random_walk") - if scenario, exist := ScenarioRegistry[scenarioId]; exist { - result.Results[query.RefId] = scenario.Handler(query, tsdbQuery) - result.Results[query.RefId].RefId = query.RefId - } else { - e.log.Error("Scenario not found", "scenarioId", scenarioId) - } +func (p *testDataPlugin) Init() error { + p.logger = log.New("tsdb.testdata") + p.scenarios = map[string]*Scenario{} + p.queryMux = datasource.NewQueryTypeMux() + p.registerScenarios() + resourceMux := http.NewServeMux() + p.registerRoutes(resourceMux) + factory := coreplugin.New(backend.ServeOpts{ + QueryDataHandler: p.queryMux, + CallResourceHandler: httpadapter.New(resourceMux), + }) + err := p.BackendPluginManager.Register("testdata", factory) + if err != nil { + p.logger.Error("Failed to register plugin", "error", err) } - - return result, nil + return nil } diff --git a/public/app/plugins/datasource/testdata/datasource.ts b/public/app/plugins/datasource/testdata/datasource.ts index c52908e4cdf..cebd9be0dbb 100644 --- a/public/app/plugins/datasource/testdata/datasource.ts +++ b/public/app/plugins/datasource/testdata/datasource.ts @@ -1,6 +1,5 @@ -import set from 'lodash/set'; import { from, merge, Observable, of } from 'rxjs'; -import { delay, map } from 'rxjs/operators'; +import { delay } from 'rxjs/operators'; import { AnnotationEvent, @@ -8,20 +7,17 @@ import { arrowTableToDataFrame, base64StringToArrowTable, DataFrame, - DataQueryError, DataQueryRequest, DataQueryResponse, - DataSourceApi, DataSourceInstanceSettings, DataTopic, LiveChannelScope, LoadingState, - TableData, TimeRange, - TimeSeries, } from '@grafana/data'; import { Scenario, TestDataQuery } from './types'; import { + DataSourceWithBackend, getBackendSrv, getLiveMeasurementsObserver, getTemplateSrv, @@ -34,9 +30,7 @@ import { getSearchFilterScopedVar } from 'app/features/variables/utils'; import { TestDataVariableSupport } from './variables'; import { generateRandomNodes, savedNodesResponse } from './nodeGraphUtils'; -type TestData = TimeSeries | TableData; - -export class TestDataDataSource extends DataSourceApi { +export class TestDataDataSource extends DataSourceWithBackend { scenariosCache?: Promise; constructor( @@ -48,7 +42,7 @@ export class TestDataDataSource extends DataSourceApi { } query(options: DataQueryRequest): Observable { - const queries: any[] = []; + const backendQueries: TestDataQuery[] = []; const streams: Array> = []; // Start streams and prepare queries @@ -80,68 +74,21 @@ export class TestDataDataSource extends DataSourceApi { streams.push(this.nodesQuery(target, options)); break; default: - queries.push({ - ...target, - intervalMs: options.intervalMs, - maxDataPoints: options.maxDataPoints, - datasourceId: this.id, - alias: this.templateSrv.replace(target.alias || '', options.scopedVars), - }); + backendQueries.push(target); } } - if (queries.length) { - const stream = getBackendSrv() - .fetch({ - method: 'POST', - url: '/api/tsdb/query', - data: { - from: options.range.from.valueOf().toString(), - to: options.range.to.valueOf().toString(), - queries: queries, - }, - }) - .pipe(map((res) => this.processQueryResult(queries, res))); - - streams.push(stream); + if (backendQueries.length) { + const backendOpts = { + ...options, + targets: backendQueries, + }; + streams.push(super.query(backendOpts)); } return merge(...streams); } - processQueryResult(queries: any, res: any): DataQueryResponse { - const data: TestData[] = []; - let error: DataQueryError | undefined = undefined; - - for (const query of queries) { - const results = res.data.results[query.refId]; - - for (const t of results.tables || []) { - const table = t as TableData; - table.refId = query.refId; - table.name = query.alias; - - if (query.scenarioId === 'logs') { - set(table, 'meta.preferredVisualisationType', 'logs'); - } - - data.push(table); - } - - for (const series of results.series || []) { - data.push({ target: series.name, datapoints: series.points, refId: query.refId, tags: series.tags }); - } - - if (results.error) { - error = { - message: results.error, - }; - } - } - - return { data, error }; - } - annotationDataTopicTest(target: TestDataQuery, req: DataQueryRequest): Observable { return new Observable((observer) => { const events = this.buildFakeAnnotationEvents(req.range, 10); @@ -190,7 +137,7 @@ export class TestDataDataSource extends DataSourceApi { getScenarios(): Promise { if (!this.scenariosCache) { - this.scenariosCache = getBackendSrv().get('/api/tsdb/testdata/scenarios'); + this.scenariosCache = this.getResource('scenarios'); } return this.scenariosCache;