From d88da205f64816506a23ccdea877ff7d8c4922f5 Mon Sep 17 00:00:00 2001 From: Ivana Huckova <30407135+ivanahuckova@users.noreply.github.com> Date: Wed, 10 May 2023 09:30:57 +0200 Subject: [PATCH] Elasticsearch: Migrate annotation calls to be run trough resources (#68075) * Update * Remove comment * Add annotation to test dashboard * Update devenv dashboard to correctly use textField --- .betterer.results | 22 +- .../elasticsearch_complex.json | 17 ++ .../elasticsearch/LegacyQueryRunner.ts | 159 +------------- .../elasticsearch/datasource.test.ts | 129 ++++++++++++ .../datasource/elasticsearch/datasource.ts | 194 +++++++++++++++++- 5 files changed, 348 insertions(+), 173 deletions(-) diff --git a/.betterer.results b/.betterer.results index 4f0262e4ecb..ef1bca72e3e 100644 --- a/.betterer.results +++ b/.betterer.results @@ -3994,17 +3994,8 @@ exports[`better eslint`] = { "public/app/plugins/datasource/elasticsearch/LegacyQueryRunner.ts:5381": [ [0, 0, 0, "Unexpected any. Specify a different type.", "0"], [0, 0, 0, "Unexpected any. Specify a different type.", "1"], - [0, 0, 0, "Unexpected any. Specify a different type.", "2"], - [0, 0, 0, "Unexpected any. Specify a different type.", "3"], - [0, 0, 0, "Unexpected any. Specify a different type.", "4"], - [0, 0, 0, "Unexpected any. Specify a different type.", "5"], - [0, 0, 0, "Unexpected any. Specify a different type.", "6"], - [0, 0, 0, "Unexpected any. Specify a different type.", "7"], - [0, 0, 0, "Unexpected any. Specify a different type.", "8"], - [0, 0, 0, "Unexpected any. Specify a different type.", "9"], - [0, 0, 0, "Unexpected any. Specify a different type.", "10"], - [0, 0, 0, "Do not use any type assertions.", "11"], - [0, 0, 0, "Unexpected any. Specify a different type.", "12"] + [0, 0, 0, "Do not use any type assertions.", "2"], + [0, 0, 0, "Unexpected any. Specify a different type.", "3"] ], "public/app/plugins/datasource/elasticsearch/QueryBuilder.test.ts:5381": [ [0, 0, 0, "Unexpected any. Specify a different type.", "0"] @@ -4083,7 +4074,14 @@ exports[`better eslint`] = { [0, 0, 0, "Unexpected any. Specify a different type.", "12"], [0, 0, 0, "Unexpected any. Specify a different type.", "13"], [0, 0, 0, "Unexpected any. Specify a different type.", "14"], - [0, 0, 0, "Unexpected any. Specify a different type.", "15"] + [0, 0, 0, "Unexpected any. Specify a different type.", "15"], + [0, 0, 0, "Unexpected any. Specify a different type.", "16"], + [0, 0, 0, "Unexpected any. Specify a different type.", "17"], + [0, 0, 0, "Unexpected any. Specify a different type.", "18"], + [0, 0, 0, "Unexpected any. Specify a different type.", "19"], + [0, 0, 0, "Unexpected any. Specify a different type.", "20"], + [0, 0, 0, "Unexpected any. Specify a different type.", "21"], + [0, 0, 0, "Unexpected any. Specify a different type.", "22"] ], "public/app/plugins/datasource/elasticsearch/hooks/useStatelessReducer.ts:5381": [ [0, 0, 0, "Do not use any type assertions.", "0"] diff --git a/devenv/dev-dashboards/datasource-elasticsearch/elasticsearch_complex.json b/devenv/dev-dashboards/datasource-elasticsearch/elasticsearch_complex.json index dc4e4b7810c..34a3db9116b 100644 --- a/devenv/dev-dashboards/datasource-elasticsearch/elasticsearch_complex.json +++ b/devenv/dev-dashboards/datasource-elasticsearch/elasticsearch_complex.json @@ -12,6 +12,23 @@ "iconColor": "rgba(0, 211, 255, 1)", "name": "Annotations & Alerts", "type": "dashboard" + }, + { + "datasource": { + "type": "elasticsearch", + "uid": "gdev-elasticsearch" + }, + "enable": true, + "iconColor": "red", + "name": "errors", + "tagsField": "metric", + "target": { + "lines": 10, + "query": "level:error", + "refId": "Anno", + "scenarioId": "annotations" + }, + "textField": "line" } ] }, diff --git a/public/app/plugins/datasource/elasticsearch/LegacyQueryRunner.ts b/public/app/plugins/datasource/elasticsearch/LegacyQueryRunner.ts index 6665b6028fd..dfc24921b43 100644 --- a/public/app/plugins/datasource/elasticsearch/LegacyQueryRunner.ts +++ b/public/app/plugins/datasource/elasticsearch/LegacyQueryRunner.ts @@ -1,4 +1,4 @@ -import { isNumber, isString, first as _first, cloneDeep } from 'lodash'; +import { first as _first, cloneDeep } from 'lodash'; import { lastValueFrom, Observable, of, throwError } from 'rxjs'; import { catchError, map, tap } from 'rxjs/operators'; @@ -11,14 +11,13 @@ import { LogRowContextOptions, LogRowContextQueryDirection, LogRowModel, - toUtc, } from '@grafana/data'; import { BackendSrvRequest, getBackendSrv, TemplateSrv } from '@grafana/runtime'; import { ElasticResponse } from './ElasticResponse'; import { ElasticDatasource, enhanceDataFrame } from './datasource'; import { defaultBucketAgg, hasMetricOfType } from './queryDef'; -import { trackAnnotationQuery, trackQuery } from './tracking'; +import { trackQuery } from './tracking'; import { ElasticsearchQuery, Logs } from './types'; export class LegacyQueryRunner { @@ -86,160 +85,6 @@ export class LegacyQueryRunner { ); } - annotationQuery(options: any) { - const annotation = options.annotation; - const timeField = annotation.timeField || '@timestamp'; - const timeEndField = annotation.timeEndField || null; - - // the `target.query` is the "new" location for the query. - // normally we would write this code as - // try-the-new-place-then-try-the-old-place, - // but we had the bug at - // https://github.com/grafana/grafana/issues/61107 - // that may have stored annotations where - // both the old and the new place are set, - // and in that scenario the old place needs - // to have priority. - const queryString = annotation.query ?? annotation.target?.query; - const tagsField = annotation.tagsField || 'tags'; - const textField = annotation.textField || null; - - const dateRanges = []; - const rangeStart: any = {}; - rangeStart[timeField] = { - from: options.range.from.valueOf(), - to: options.range.to.valueOf(), - format: 'epoch_millis', - }; - dateRanges.push({ range: rangeStart }); - - if (timeEndField) { - const rangeEnd: any = {}; - rangeEnd[timeEndField] = { - from: options.range.from.valueOf(), - to: options.range.to.valueOf(), - format: 'epoch_millis', - }; - dateRanges.push({ range: rangeEnd }); - } - - const queryInterpolated = this.datasource.interpolateLuceneQuery(queryString); - const query: any = { - bool: { - filter: [ - { - bool: { - should: dateRanges, - minimum_should_match: 1, - }, - }, - ], - }, - }; - - if (queryInterpolated) { - query.bool.filter.push({ - query_string: { - query: queryInterpolated, - }, - }); - } - const data: any = { - query, - size: 10000, - }; - - const header: any = { - search_type: 'query_then_fetch', - ignore_unavailable: true, - }; - - // @deprecated - // Field annotation.index is deprecated and will be removed in the future - if (annotation.index) { - header.index = annotation.index; - } else { - header.index = this.datasource.indexPattern.getIndexList(options.range.from, options.range.to); - } - - const payload = JSON.stringify(header) + '\n' + JSON.stringify(data) + '\n'; - - trackAnnotationQuery(annotation); - return lastValueFrom( - this.request('POST', '_msearch', payload).pipe( - map((res) => { - const list = []; - const hits = res.responses[0].hits.hits; - - const getFieldFromSource = (source: any, fieldName: any) => { - if (!fieldName) { - return; - } - - const fieldNames = fieldName.split('.'); - let fieldValue = source; - - for (let i = 0; i < fieldNames.length; i++) { - fieldValue = fieldValue[fieldNames[i]]; - if (!fieldValue) { - console.log('could not find field in annotation: ', fieldName); - return ''; - } - } - - return fieldValue; - }; - - for (let i = 0; i < hits.length; i++) { - const source = hits[i]._source; - let time = getFieldFromSource(source, timeField); - if (typeof hits[i].fields !== 'undefined') { - const fields = hits[i].fields; - if (isString(fields[timeField]) || isNumber(fields[timeField])) { - time = fields[timeField]; - } - } - - const event: { - annotation: any; - time: number; - timeEnd?: number; - text: string; - tags: string | string[]; - } = { - annotation: annotation, - time: toUtc(time).valueOf(), - text: getFieldFromSource(source, textField), - tags: getFieldFromSource(source, tagsField), - }; - - if (timeEndField) { - const timeEnd = getFieldFromSource(source, timeEndField); - if (timeEnd) { - event.timeEnd = toUtc(timeEnd).valueOf(); - } - } - - // legacy support for title field - if (annotation.titleField) { - const title = getFieldFromSource(source, annotation.titleField); - if (title) { - event.text = title + '\n' + event.text; - } - } - - if (typeof event.tags === 'string') { - event.tags = event.tags.split(','); - } - - list.push(event); - } - return list; - }) - ) - ); - } - async logContextQuery(row: LogRowModel, options?: LogRowContextOptions): Promise<{ data: DataFrame[] }> { const sortField = row.dataFrame.fields.find((f) => f.name === 'sort'); const searchAfter = sortField?.values[row.rowIndex] || [row.timeEpochMs]; diff --git a/public/app/plugins/datasource/elasticsearch/datasource.test.ts b/public/app/plugins/datasource/elasticsearch/datasource.test.ts index 479e9e20852..8df01e4fa8a 100644 --- a/public/app/plugins/datasource/elasticsearch/datasource.test.ts +++ b/public/app/plugins/datasource/elasticsearch/datasource.test.ts @@ -1304,6 +1304,135 @@ describe('ElasticDatasource using backend', () => { console.error = originalConsoleError; config.featureToggles.enableElasticsearchBackendQuerying = false; }); + describe('annotationQuery', () => { + describe('results processing', () => { + it('should return simple annotations using defaults', async () => { + const { ds, timeSrv } = getTestContext(); + ds.postResourceRequest = jest.fn().mockResolvedValue({ + responses: [ + { + hits: { + hits: [ + { _source: { '@timestamp': 1, '@test_tags': 'foo', text: 'abc' } }, + { _source: { '@timestamp': 3, '@test_tags': 'bar', text: 'def' } }, + ], + }, + }, + ], + }); + + const annotations = await ds.annotationQuery({ + annotation: {}, + range: timeSrv.timeRange(), + }); + + expect(annotations).toHaveLength(2); + expect(annotations[0].time).toBe(1); + expect(annotations[1].time).toBe(3); + }); + + it('should return annotation events using options', async () => { + const { ds, timeSrv } = getTestContext(); + ds.postResourceRequest = jest.fn().mockResolvedValue({ + responses: [ + { + hits: { + hits: [ + { _source: { '@test_time': 1, '@test_tags': 'foo', text: 'abc' } }, + { _source: { '@test_time': 3, '@test_tags': 'bar', text: 'def' } }, + ], + }, + }, + ], + }); + + const annotations = await ds.annotationQuery({ + annotation: { + timeField: '@test_time', + name: 'foo', + query: 'abc', + tagsField: '@test_tags', + textField: 'text', + }, + range: timeSrv.timeRange(), + }); + expect(annotations).toHaveLength(2); + expect(annotations[0].time).toBe(1); + expect(annotations[0].tags?.[0]).toBe('foo'); + expect(annotations[0].text).toBe('abc'); + + expect(annotations[1].time).toBe(3); + expect(annotations[1].tags?.[0]).toBe('bar'); + expect(annotations[1].text).toBe('def'); + }); + }); + + describe('request processing', () => { + it('should process annotation request using options', async () => { + const { ds } = getTestContext(); + const postResourceRequestMock = jest.spyOn(ds, 'postResourceRequest').mockResolvedValue({ + responses: [ + { + hits: { + hits: [ + { _source: { '@test_time': 1, '@test_tags': 'foo', text: 'abc' } }, + { _source: { '@test_time': 3, '@test_tags': 'bar', text: 'def' } }, + ], + }, + }, + ], + }); + + await ds.annotationQuery({ + annotation: { + timeField: '@test_time', + timeEndField: '@time_end_field', + name: 'foo', + query: 'abc', + tagsField: '@test_tags', + textField: 'text', + }, + range: { + from: dateTime(1683291160012), + to: dateTime(1683291460012), + }, + }); + expect(postResourceRequestMock).toHaveBeenCalledWith( + '_msearch', + '{"search_type":"query_then_fetch","ignore_unavailable":true,"index":"[test-]YYYY.MM.DD"}\n{"query":{"bool":{"filter":[{"bool":{"should":[{"range":{"@test_time":{"from":1683291160012,"to":1683291460012,"format":"epoch_millis"}}},{"range":{"@time_end_field":{"from":1683291160012,"to":1683291460012,"format":"epoch_millis"}}}],"minimum_should_match":1}},{"query_string":{"query":"abc"}}]}},"size":10000}\n' + ); + }); + + it('should process annotation request using defaults', async () => { + const { ds } = getTestContext(); + const postResourceRequestMock = jest.spyOn(ds, 'postResourceRequest').mockResolvedValue({ + responses: [ + { + hits: { + hits: [ + { _source: { '@test_time': 1, '@test_tags': 'foo', text: 'abc' } }, + { _source: { '@test_time': 3, '@test_tags': 'bar', text: 'def' } }, + ], + }, + }, + ], + }); + + await ds.annotationQuery({ + annotation: {}, + range: { + from: dateTime(1683291160012), + to: dateTime(1683291460012), + }, + }); + expect(postResourceRequestMock).toHaveBeenCalledWith( + '_msearch', + '{"search_type":"query_then_fetch","ignore_unavailable":true,"index":"[test-]YYYY.MM.DD"}\n{"query":{"bool":{"filter":[{"bool":{"should":[{"range":{"@timestamp":{"from":1683291160012,"to":1683291460012,"format":"epoch_millis"}}}],"minimum_should_match":1}}]}},"size":10000}\n' + ); + }); + }); + }); + describe('getDatabaseVersion', () => { it('should correctly get db version', async () => { const { ds } = getTestContext(); diff --git a/public/app/plugins/datasource/elasticsearch/datasource.ts b/public/app/plugins/datasource/elasticsearch/datasource.ts index 7e2cad0728e..62ac495d825 100644 --- a/public/app/plugins/datasource/elasticsearch/datasource.ts +++ b/public/app/plugins/datasource/elasticsearch/datasource.ts @@ -1,4 +1,4 @@ -import { cloneDeep, find, first as _first, isObject, isString, map as _map } from 'lodash'; +import { cloneDeep, find, first as _first, isNumber, isObject, isString, map as _map } from 'lodash'; import { from, generate, lastValueFrom, Observable, of } from 'rxjs'; import { catchError, first, map, mergeMap, skipWhile, throwIfEmpty, tap } from 'rxjs/operators'; import { SemVer } from 'semver'; @@ -29,6 +29,8 @@ import { LogRowContextQueryDirection, LogRowContextOptions, SupplementaryQueryOptions, + toUtc, + AnnotationEvent, } from '@grafana/data'; import { DataSourceWithBackend, getDataSourceSrv, config, BackendSrvRequest } from '@grafana/runtime'; import { queryLogsVolume } from 'app/core/logsModel'; @@ -49,7 +51,7 @@ import { isPipelineAggregationWithMultipleBucketPaths, } from './components/QueryEditor/MetricAggregationsEditor/aggregations'; import { metricAggregationConfig } from './components/QueryEditor/MetricAggregationsEditor/utils'; -import { trackQuery } from './tracking'; +import { trackAnnotationQuery, trackQuery } from './tracking'; import { Logs, BucketAggregation, @@ -211,8 +213,192 @@ export class ElasticDatasource ); } - annotationQuery(options: any): Promise { - return this.legacyQueryRunner.annotationQuery(options); + annotationQuery(options: any): Promise { + const payload = this.prepareAnnotationRequest(options); + trackAnnotationQuery(options.annotation); + const annotationObservable = config.featureToggles.enableElasticsearchBackendQuerying + ? // TODO: We should migrate this to use query and not resource call + // The plan is to look at this when we start to work on raw query editor for ES + // as we will have to explore how to handle any query + from(this.postResourceRequest('_msearch', payload)) + : this.legacyQueryRunner.request('POST', '_msearch', payload); + + return lastValueFrom( + annotationObservable.pipe( + map((res) => { + const hits = res.responses[0].hits.hits; + return this.processHitsToAnnotationEvents(options.annotation, hits); + }) + ) + ); + } + + private prepareAnnotationRequest(options: { + annotation: { + target: ElasticsearchQuery; + timeField?: string; + timeEndField?: string; + titleField?: string; + query?: string; + index?: string; + }; + range: TimeRange; + }) { + const annotation = options.annotation; + const timeField = annotation.timeField || '@timestamp'; + const timeEndField = annotation.timeEndField || null; + + // the `target.query` is the "new" location for the query. + // normally we would write this code as + // try-the-new-place-then-try-the-old-place, + // but we had the bug at + // https://github.com/grafana/grafana/issues/61107 + // that may have stored annotations where + // both the old and the new place are set, + // and in that scenario the old place needs + // to have priority. + const queryString = annotation.query ?? annotation.target?.query ?? ''; + + const dateRanges = []; + const rangeStart: any = {}; + rangeStart[timeField] = { + from: options.range.from.valueOf(), + to: options.range.to.valueOf(), + format: 'epoch_millis', + }; + dateRanges.push({ range: rangeStart }); + + if (timeEndField) { + const rangeEnd: any = {}; + rangeEnd[timeEndField] = { + from: options.range.from.valueOf(), + to: options.range.to.valueOf(), + format: 'epoch_millis', + }; + dateRanges.push({ range: rangeEnd }); + } + + const queryInterpolated = this.interpolateLuceneQuery(queryString); + const query: any = { + bool: { + filter: [ + { + bool: { + should: dateRanges, + minimum_should_match: 1, + }, + }, + ], + }, + }; + + if (queryInterpolated) { + query.bool.filter.push({ + query_string: { + query: queryInterpolated, + }, + }); + } + const data: any = { + query, + size: 10000, + }; + + const header: any = { + search_type: 'query_then_fetch', + ignore_unavailable: true, + }; + + // @deprecated + // Field annotation.index is deprecated and will be removed in the future + if (annotation.index) { + header.index = annotation.index; + } else { + header.index = this.indexPattern.getIndexList(options.range.from, options.range.to); + } + + const payload = JSON.stringify(header) + '\n' + JSON.stringify(data) + '\n'; + return payload; + } + + private processHitsToAnnotationEvents( + annotation: { + target: ElasticsearchQuery; + timeField?: string; + titleField?: string; + timeEndField?: string; + query?: string; + tagsField?: string; + textField?: string; + index?: string; + }, + hits: Array<{ [key: string]: any }> + ) { + const timeField = annotation.timeField || '@timestamp'; + const timeEndField = annotation.timeEndField || null; + const textField = annotation.textField || 'tags'; + const tagsField = annotation.tagsField || null; + const list: AnnotationEvent[] = []; + + const getFieldFromSource = (source: any, fieldName: any) => { + if (!fieldName) { + return; + } + + const fieldNames = fieldName.split('.'); + let fieldValue = source; + + for (let i = 0; i < fieldNames.length; i++) { + fieldValue = fieldValue[fieldNames[i]]; + if (!fieldValue) { + return ''; + } + } + + return fieldValue; + }; + + for (let i = 0; i < hits.length; i++) { + const source = hits[i]._source; + let time = getFieldFromSource(source, timeField); + if (typeof hits[i].fields !== 'undefined') { + const fields = hits[i].fields; + if (isString(fields[timeField]) || isNumber(fields[timeField])) { + time = fields[timeField]; + } + } + + const event: AnnotationEvent = { + annotation: annotation, + time: toUtc(time).valueOf(), + text: getFieldFromSource(source, textField), + }; + + if (timeEndField) { + const timeEnd = getFieldFromSource(source, timeEndField); + if (timeEnd) { + event.timeEnd = toUtc(timeEnd).valueOf(); + } + } + + // legacy support for title field + if (annotation.titleField) { + const title = getFieldFromSource(source, annotation.titleField); + if (title) { + event.text = title + '\n' + event.text; + } + } + + const tags = getFieldFromSource(source, tagsField); + if (typeof tags === 'string') { + event.tags = tags.split(','); + } else { + event.tags = tags; + } + + list.push(event); + } + return list; } interpolateLuceneQuery(queryString: string, scopedVars?: ScopedVars) {