diff --git a/public/app/core/logs_model.test.ts b/public/app/core/logs_model.test.ts index d9a5366f290..2c1642afef5 100644 --- a/public/app/core/logs_model.test.ts +++ b/public/app/core/logs_model.test.ts @@ -1,7 +1,12 @@ import { ArrayVector, DataFrame, + DataQuery, + DataQueryRequest, + DataQueryResponse, + dateTimeParse, FieldType, + LoadingState, LogLevel, LogRowModel, LogsDedupStrategy, @@ -10,14 +15,17 @@ import { toDataFrame, } from '@grafana/data'; import { + COMMON_LABELS, dataFrameToLogsModel, dedupLogRows, - getSeriesProperties, - logSeriesToLogsModel, filterLogLevels, + getSeriesProperties, LIMIT_LABEL, - COMMON_LABELS, + logSeriesToLogsModel, + queryLogsVolume, } from './logs_model'; +import { Observable } from 'rxjs'; +import { MockObservableDataSourceApi } from '../../test/mocks/datasource_srv'; describe('dedupLogRows()', () => { test('should return rows as is when dedup is set to none', () => { @@ -939,3 +947,120 @@ describe('getSeriesProperties()', () => { expect(result.visibleRange).toMatchObject({ from: 8, to: 30 }); }); }); + +describe('logs volume', () => { + class TestDataQuery implements DataQuery { + refId = 'a'; + target = ''; + } + + let volumeProvider: Observable, + datasource: MockObservableDataSourceApi, + request: DataQueryRequest; + + function createFrame(labels: object, timestamps: number[], values: number[]) { + return toDataFrame({ + fields: [ + { name: 'Time', type: FieldType.time, values: timestamps }, + { + name: 'Number', + type: FieldType.number, + values, + labels, + }, + ], + }); + } + + function createExpectedFields(levelName: string, timestamps: number[], values: number[]) { + return [ + { name: 'Time', values: { buffer: timestamps } }, + { + name: 'Value', + config: { displayNameFromDS: levelName }, + values: { buffer: values }, + }, + ]; + } + + function setup(datasourceSetup: () => void) { + datasourceSetup(); + request = ({ + targets: [{ target: 'volume query 1' }, { target: 'volume query 2' }], + scopedVars: {}, + } as unknown) as DataQueryRequest; + volumeProvider = queryLogsVolume(datasource, request, { + extractLevel: (dataFrame: DataFrame) => { + return dataFrame.fields[1]!.labels!.level === 'error' ? LogLevel.error : LogLevel.unknown; + }, + range: { + from: dateTimeParse('2021-06-17 00:00:00', { timeZone: 'utc' }), + to: dateTimeParse('2021-06-17 00:00:00', { timeZone: 'utc' }), + raw: { from: '0', to: '1' }, + }, + targets: request.targets, + }); + } + + function setupMultipleResults() { + // level=unknown + const resultAFrame1 = createFrame({ app: 'app01' }, [100, 200, 300], [5, 5, 5]); + // level=error + const resultAFrame2 = createFrame({ app: 'app01', level: 'error' }, [100, 200, 300], [0, 1, 0]); + // level=unknown + const resultBFrame1 = createFrame({ app: 'app02' }, [100, 200, 300], [1, 2, 3]); + // level=error + const resultBFrame2 = createFrame({ app: 'app02', level: 'error' }, [100, 200, 300], [1, 1, 1]); + + datasource = new MockObservableDataSourceApi('loki', [ + { + data: [resultAFrame1, resultAFrame2], + }, + { + data: [resultBFrame1, resultBFrame2], + }, + ]); + } + + function setupErrorResponse() { + datasource = new MockObservableDataSourceApi('loki', [], undefined, 'Error message'); + } + + it('aggregates data frames by level', async () => { + setup(setupMultipleResults); + + await expect(volumeProvider).toEmitValuesWith((received) => { + expect(received).toMatchObject([ + { state: LoadingState.Loading, error: undefined, data: [] }, + { + state: LoadingState.Done, + error: undefined, + data: [ + { + fields: createExpectedFields('unknown', [100, 200, 300], [6, 7, 8]), + }, + { + fields: createExpectedFields('error', [100, 200, 300], [1, 2, 1]), + }, + ], + }, + ]); + }); + }); + + it('returns error', async () => { + setup(setupErrorResponse); + + await expect(volumeProvider).toEmitValuesWith((received) => { + expect(received).toMatchObject([ + { state: LoadingState.Loading, error: undefined, data: [] }, + { + state: LoadingState.Error, + error: 'Error message', + data: [], + }, + 'Error message', + ]); + }); + }); +}); diff --git a/public/app/core/logs_model.ts b/public/app/core/logs_model.ts index 4d29360f66d..86e5e2b1c90 100644 --- a/public/app/core/logs_model.ts +++ b/public/app/core/logs_model.ts @@ -6,11 +6,15 @@ import { AbsoluteTimeRange, DataFrame, DataQuery, + DataQueryRequest, + DataQueryResponse, + DataSourceApi, dateTime, dateTimeFormat, dateTimeFormatTimeAgo, FieldCache, FieldColorModeId, + FieldConfig, FieldType, FieldWithIndex, findCommonLabels, @@ -18,19 +22,24 @@ import { getLogLevel, getLogLevelFromKey, Labels, + LoadingState, LogLevel, LogRowModel, LogsDedupStrategy, LogsMetaItem, LogsMetaKind, LogsModel, + MutableDataFrame, rangeUtil, + ScopedVars, sortInAscendingOrder, textUtil, + TimeRange, toDataFrame, } from '@grafana/data'; import { getThemeColor } from 'app/core/utils/colors'; import { SIPrefix } from '@grafana/data/src/valueFormats/symbolFormatters'; +import { Observable, throwError, timeout } from 'rxjs'; export const LIMIT_LABEL = 'Line limit'; export const COMMON_LABELS = 'Common labels'; @@ -45,6 +54,11 @@ export const LogLevelColor = { [LogLevel.unknown]: getThemeColor('#8e8e8e', '#dde4ed'), }; +const SECOND = 1000; +const MINUTE = 60 * SECOND; +const HOUR = 60 * MINUTE; +const DAY = 24 * HOUR; + const isoDateRegexp = /\d{4}-[01]\d-[0-3]\dT[0-2]\d:[0-5]\d:[0-6]\d[,\.]\d+([+-][0-2]\d:[0-5]\d|Z)/g; function isDuplicateRow(row: LogRowModel, other: LogRowModel, strategy?: LogsDedupStrategy): boolean { switch (strategy) { @@ -510,3 +524,194 @@ function adjustMetaInfo(logsModel: LogsModel, visibleRangeMs?: number, requested return logsModelMeta; } + +/** + * Returns field configuration used to render logs volume bars + */ +function getLogVolumeFieldConfig(level: LogLevel, oneLevelDetected: boolean) { + const name = oneLevelDetected && level === LogLevel.unknown ? 'logs' : level; + const color = LogLevelColor[level]; + return { + displayNameFromDS: name, + color: { + mode: FieldColorModeId.Fixed, + fixedColor: color, + }, + custom: { + drawStyle: GraphDrawStyle.Bars, + barAlignment: BarAlignment.Center, + lineColor: color, + pointColor: color, + fillColor: color, + lineWidth: 1, + fillOpacity: 100, + stacking: { + mode: StackingMode.Normal, + group: 'A', + }, + }, + }; +} + +/** + * Take multiple data frames, sum up values and group by level. + * Return a list of data frames, each representing single level. + */ +export function aggregateRawLogsVolume( + rawLogsVolume: DataFrame[], + extractLevel: (dataFrame: DataFrame) => LogLevel +): DataFrame[] { + const logsVolumeByLevelMap: Partial> = {}; + rawLogsVolume.forEach((dataFrame) => { + const level = extractLevel(dataFrame); + if (!logsVolumeByLevelMap[level]) { + logsVolumeByLevelMap[level] = []; + } + logsVolumeByLevelMap[level]!.push(dataFrame); + }); + + return Object.keys(logsVolumeByLevelMap).map((level: string) => { + return aggregateFields( + logsVolumeByLevelMap[level as LogLevel]!, + getLogVolumeFieldConfig(level as LogLevel, Object.keys(logsVolumeByLevelMap).length === 1) + ); + }); +} + +/** + * Aggregate multiple data frames into a single data frame by adding values. + * Multiple data frames for the same level are passed here to get a single + * data frame for a given level. Aggregation by level happens in aggregateRawLogsVolume() + */ +function aggregateFields(dataFrames: DataFrame[], config: FieldConfig): DataFrame { + const aggregatedDataFrame = new MutableDataFrame(); + if (!dataFrames.length) { + return aggregatedDataFrame; + } + + const totalLength = dataFrames[0].length; + const timeField = new FieldCache(dataFrames[0]).getFirstFieldOfType(FieldType.time); + + if (!timeField) { + return aggregatedDataFrame; + } + + aggregatedDataFrame.addField({ name: 'Time', type: FieldType.time }, totalLength); + aggregatedDataFrame.addField({ name: 'Value', type: FieldType.number, config }, totalLength); + + dataFrames.forEach((dataFrame) => { + dataFrame.fields.forEach((field) => { + if (field.type === FieldType.number) { + for (let pointIndex = 0; pointIndex < totalLength; pointIndex++) { + const currentValue = aggregatedDataFrame.get(pointIndex).Value; + const valueToAdd = field.values.get(pointIndex); + const totalValue = + currentValue === null && valueToAdd === null ? null : (currentValue || 0) + (valueToAdd || 0); + aggregatedDataFrame.set(pointIndex, { Value: totalValue, Time: timeField.values.get(pointIndex) }); + } + } + }); + }); + + return aggregatedDataFrame; +} + +const LOGS_VOLUME_QUERY_DEFAULT_TIMEOUT = 60000; + +type LogsVolumeQueryOptions = { + timeout?: number; + extractLevel: (dataFrame: DataFrame) => LogLevel; + targets: T[]; + range: TimeRange; +}; + +/** + * Creates an observable, which makes requests to get logs volume and aggregates results. + */ +export function queryLogsVolume( + datasource: DataSourceApi, + logsVolumeRequest: DataQueryRequest, + options: LogsVolumeQueryOptions +): Observable { + const intervalInfo = getIntervalInfo(logsVolumeRequest.scopedVars); + logsVolumeRequest.interval = intervalInfo.interval; + logsVolumeRequest.scopedVars.__interval = { value: intervalInfo.interval, text: intervalInfo.interval }; + if (intervalInfo.intervalMs !== undefined) { + logsVolumeRequest.intervalMs = intervalInfo.intervalMs; + logsVolumeRequest.scopedVars.__interval_ms = { value: intervalInfo.intervalMs, text: intervalInfo.intervalMs }; + } + + return new Observable((observer) => { + let rawLogsVolume: DataFrame[] = []; + observer.next({ + state: LoadingState.Loading, + error: undefined, + data: [], + }); + + const subscription = (datasource.query(logsVolumeRequest) as Observable) + .pipe( + timeout({ + each: options.timeout || LOGS_VOLUME_QUERY_DEFAULT_TIMEOUT, + with: () => throwError(new Error('Request timed-out. Please make your query more specific and try again.')), + }) + ) + .subscribe({ + complete: () => { + const aggregatedLogsVolume = aggregateRawLogsVolume(rawLogsVolume, options.extractLevel); + if (aggregatedLogsVolume[0]) { + aggregatedLogsVolume[0].meta = { + custom: { + targets: options.targets, + absoluteRange: { from: options.range.from.valueOf(), to: options.range.to.valueOf() }, + }, + }; + } + observer.next({ + state: LoadingState.Done, + error: undefined, + data: aggregatedLogsVolume, + }); + observer.complete(); + }, + next: (dataQueryResponse: DataQueryResponse) => { + rawLogsVolume = rawLogsVolume.concat(dataQueryResponse.data.map(toDataFrame)); + }, + error: (error) => { + observer.next({ + state: LoadingState.Error, + error: error, + data: [], + }); + observer.error(error); + }, + }); + return () => { + subscription?.unsubscribe(); + }; + }); +} + +function getIntervalInfo(scopedVars: ScopedVars): { interval: string; intervalMs?: number } { + if (scopedVars.__interval) { + let intervalMs: number = scopedVars.__interval_ms.value; + let interval = ''; + if (intervalMs > HOUR) { + intervalMs = DAY; + interval = '1d'; + } else if (intervalMs > MINUTE) { + intervalMs = HOUR; + interval = '1h'; + } else if (intervalMs > SECOND) { + intervalMs = MINUTE; + interval = '1m'; + } else { + intervalMs = SECOND; + interval = '1s'; + } + + return { interval, intervalMs }; + } else { + return { interval: '$__interval' }; + } +} diff --git a/public/app/plugins/datasource/elasticsearch/datasource.ts b/public/app/plugins/datasource/elasticsearch/datasource.ts index d660831600d..f1aaeb2f90e 100644 --- a/public/app/plugins/datasource/elasticsearch/datasource.ts +++ b/public/app/plugins/datasource/elasticsearch/datasource.ts @@ -12,10 +12,13 @@ import { DataSourceApi, DataSourceInstanceSettings, DataSourceWithLogsContextSupport, + DataSourceWithLogsVolumeSupport, DateTime, dateTime, Field, getDefaultTimeRange, + getLogLevelFromKey, + LogLevel, LogRowModel, MetricFindValue, ScopedVars, @@ -42,6 +45,7 @@ import { isBucketAggregationWithField, } from './components/QueryEditor/BucketAggregationsEditor/aggregations'; import { coerceESVersion, getScriptValue } from './utils'; +import { queryLogsVolume } from 'app/core/logs_model'; // Those are metadata fields as defined in https://www.elastic.co/guide/en/elasticsearch/reference/current/mapping-fields.html#_identity_metadata_fields. // custom fields can start with underscores, therefore is not safe to exclude anything that starts with one. @@ -59,7 +63,7 @@ const ELASTIC_META_FIELDS = [ export class ElasticDatasource extends DataSourceApi - implements DataSourceWithLogsContextSupport { + implements DataSourceWithLogsContextSupport, DataSourceWithLogsVolumeSupport { basicAuth?: string; withCredentials?: boolean; url: string; @@ -557,6 +561,60 @@ export class ElasticDatasource return logResponse; }; + getLogsVolumeDataProvider(request: DataQueryRequest): Observable | undefined { + const isLogsVolumeAvailable = request.targets.some((target) => { + return target.metrics?.length === 1 && target.metrics[0].type === 'logs'; + }); + if (!isLogsVolumeAvailable) { + return undefined; + } + const logsVolumeRequest = cloneDeep(request); + logsVolumeRequest.targets = logsVolumeRequest.targets.map((target) => { + const bucketAggs: BucketAggregation[] = []; + const timeField = this.timeField ?? '@timestamp'; + + if (this.logLevelField) { + bucketAggs.push({ + id: '2', + type: 'terms', + settings: { + min_doc_count: '0', + size: '0', + order: 'desc', + orderBy: '_count', + missing: LogLevel.unknown, + }, + field: this.logLevelField, + }); + } + bucketAggs.push({ + id: '3', + type: 'date_histogram', + settings: { + interval: 'auto', + min_doc_count: '0', + trimEdges: '0', + }, + field: timeField, + }); + + const logsVolumeQuery: ElasticsearchQuery = { + refId: target.refId, + query: target.query, + metrics: [{ type: 'count', id: '1' }], + timeField, + bucketAggs, + }; + return logsVolumeQuery; + }); + + return queryLogsVolume(this, logsVolumeRequest, { + range: request.range, + targets: request.targets, + extractLevel: (dataFrame) => getLogLevelFromKey(dataFrame.name || ''), + }); + } + query(options: DataQueryRequest): Observable { let payload = ''; const targets = this.interpolateVariablesInQueries(cloneDeep(options.targets), options.scopedVars); diff --git a/public/app/plugins/datasource/loki/dataProviders/logsVolumeProvider.test.ts b/public/app/plugins/datasource/loki/dataProviders/logsVolumeProvider.test.ts deleted file mode 100644 index 6399e86a87e..00000000000 --- a/public/app/plugins/datasource/loki/dataProviders/logsVolumeProvider.test.ts +++ /dev/null @@ -1,113 +0,0 @@ -import { MockObservableDataSourceApi } from '../../../../../test/mocks/datasource_srv'; -import { createLokiLogsVolumeProvider } from './logsVolumeProvider'; -import LokiDatasource from '../datasource'; -import { DataQueryRequest, DataQueryResponse, FieldType, LoadingState, toDataFrame } from '@grafana/data'; -import { LokiQuery } from '../types'; -import { Observable } from 'rxjs'; - -function createFrame(labels: object, timestamps: number[], values: number[]) { - return toDataFrame({ - fields: [ - { name: 'Time', type: FieldType.time, values: timestamps }, - { - name: 'Number', - type: FieldType.number, - values, - labels, - }, - ], - }); -} - -function createExpectedFields(levelName: string, timestamps: number[], values: number[]) { - return [ - { name: 'Time', values: { buffer: timestamps } }, - { - name: 'Value', - config: { displayNameFromDS: levelName }, - values: { buffer: values }, - }, - ]; -} - -describe('LokiLogsVolumeProvider', () => { - let volumeProvider: Observable, - datasource: MockObservableDataSourceApi, - request: DataQueryRequest; - - function setup(datasourceSetup: () => void) { - datasourceSetup(); - request = ({ - targets: [{ expr: '{app="app01"}' }, { expr: '{app="app02"}' }], - range: { from: 0, to: 1 }, - scopedVars: { - __interval_ms: { - value: 1000, - }, - }, - } as unknown) as DataQueryRequest; - volumeProvider = createLokiLogsVolumeProvider((datasource as unknown) as LokiDatasource, request); - } - - function setupMultipleResults() { - // level=unknown - const resultAFrame1 = createFrame({ app: 'app01' }, [100, 200, 300], [5, 5, 5]); - // level=error - const resultAFrame2 = createFrame({ app: 'app01', level: 'error' }, [100, 200, 300], [0, 1, 0]); - // level=unknown - const resultBFrame1 = createFrame({ app: 'app02' }, [100, 200, 300], [1, 2, 3]); - // level=error - const resultBFrame2 = createFrame({ app: 'app02', level: 'error' }, [100, 200, 300], [1, 1, 1]); - - datasource = new MockObservableDataSourceApi('loki', [ - { - data: [resultAFrame1, resultAFrame2], - }, - { - data: [resultBFrame1, resultBFrame2], - }, - ]); - } - - function setupErrorResponse() { - datasource = new MockObservableDataSourceApi('loki', [], undefined, 'Error message'); - } - - it('aggregates data frames by level', async () => { - setup(setupMultipleResults); - - await expect(volumeProvider).toEmitValuesWith((received) => { - expect(received).toMatchObject([ - { state: LoadingState.Loading, error: undefined, data: [] }, - { - state: LoadingState.Done, - error: undefined, - data: [ - { - fields: createExpectedFields('unknown', [100, 200, 300], [6, 7, 8]), - }, - { - fields: createExpectedFields('error', [100, 200, 300], [1, 2, 1]), - }, - ], - }, - ]); - }); - }); - - it('returns error', async () => { - setup(setupErrorResponse); - - await expect(volumeProvider).toEmitValuesWith((received) => { - expect(received).toMatchObject([ - { state: LoadingState.Loading, error: undefined, data: [] }, - { - state: LoadingState.Error, - error: 'Error message', - data: [], - }, - 'Error message', - ]); - }); - }); -}); diff --git a/public/app/plugins/datasource/loki/dataProviders/logsVolumeProvider.ts b/public/app/plugins/datasource/loki/dataProviders/logsVolumeProvider.ts deleted file mode 100644 index 364fde378a7..00000000000 --- a/public/app/plugins/datasource/loki/dataProviders/logsVolumeProvider.ts +++ /dev/null @@ -1,236 +0,0 @@ -import { - DataFrame, - DataQueryRequest, - DataQueryResponse, - FieldCache, - FieldColorModeId, - FieldConfig, - FieldType, - getLogLevelFromKey, - Labels, - LoadingState, - LogLevel, - MutableDataFrame, - ScopedVars, - toDataFrame, -} from '@grafana/data'; -import { LokiQuery } from '../types'; -import { Observable, throwError, timeout } from 'rxjs'; -import { cloneDeep } from 'lodash'; -import LokiDatasource, { isMetricsQuery } from '../datasource'; -import { LogLevelColor } from '../../../../core/logs_model'; -import { BarAlignment, GraphDrawStyle, StackingMode } from '@grafana/schema'; - -const SECOND = 1000; -const MINUTE = 60 * SECOND; -const HOUR = 60 * MINUTE; -const DAY = 24 * HOUR; - -/** - * Logs volume query may be expensive as it requires counting all logs in the selected range. If such query - * takes too much time it may need be made more specific to limit number of logs processed under the hood. - */ -const TIMEOUT = 10 * SECOND; - -export function createLokiLogsVolumeProvider( - datasource: LokiDatasource, - dataQueryRequest: DataQueryRequest -): Observable { - const logsVolumeRequest = cloneDeep(dataQueryRequest); - const intervalInfo = getIntervalInfo(dataQueryRequest.scopedVars); - logsVolumeRequest.targets = logsVolumeRequest.targets - .filter((target) => target.expr && !isMetricsQuery(target.expr)) - .map((target) => { - return { - ...target, - instant: false, - expr: `sum by (level) (count_over_time(${target.expr}[${intervalInfo.interval}]))`, - }; - }); - logsVolumeRequest.interval = intervalInfo.interval; - if (intervalInfo.intervalMs !== undefined) { - logsVolumeRequest.intervalMs = intervalInfo.intervalMs; - } - - return new Observable((observer) => { - let rawLogsVolume: DataFrame[] = []; - observer.next({ - state: LoadingState.Loading, - error: undefined, - data: [], - }); - - const subscription = datasource - .query(logsVolumeRequest) - .pipe( - timeout({ - each: TIMEOUT, - with: () => - throwError( - new Error( - 'Request timed-out. Please try making your query more specific or narrow selected time range and try again.' - ) - ), - }) - ) - .subscribe({ - complete: () => { - const aggregatedLogsVolume = aggregateRawLogsVolume(rawLogsVolume); - if (aggregatedLogsVolume[0]) { - aggregatedLogsVolume[0].meta = { - custom: { - targets: dataQueryRequest.targets, - absoluteRange: { from: dataQueryRequest.range.from.valueOf(), to: dataQueryRequest.range.to.valueOf() }, - }, - }; - } - observer.next({ - state: LoadingState.Done, - error: undefined, - data: aggregatedLogsVolume, - }); - observer.complete(); - }, - next: (dataQueryResponse: DataQueryResponse) => { - rawLogsVolume = rawLogsVolume.concat(dataQueryResponse.data.map(toDataFrame)); - }, - error: (error) => { - observer.next({ - state: LoadingState.Error, - error: error, - data: [], - }); - observer.error(error); - }, - }); - return () => { - subscription?.unsubscribe(); - }; - }); -} - -/** - * Add up values for the same level and create a single data frame for each level - */ -function aggregateRawLogsVolume(rawLogsVolume: DataFrame[]): DataFrame[] { - const logsVolumeByLevelMap: { [level in LogLevel]?: DataFrame[] } = {}; - let levels = 0; - rawLogsVolume.forEach((dataFrame) => { - let valueField; - try { - valueField = new FieldCache(dataFrame).getFirstFieldOfType(FieldType.number); - } catch {} - // If value field doesn't exist skip the frame (it may happen with instant queries) - if (!valueField) { - return; - } - const level: LogLevel = valueField.labels ? getLogLevelFromLabels(valueField.labels) : LogLevel.unknown; - if (!logsVolumeByLevelMap[level]) { - logsVolumeByLevelMap[level] = []; - levels++; - } - logsVolumeByLevelMap[level]!.push(dataFrame); - }); - - return Object.keys(logsVolumeByLevelMap).map((level: string) => { - return aggregateFields(logsVolumeByLevelMap[level as LogLevel]!, getFieldConfig(level as LogLevel, levels)); - }); -} - -function getFieldConfig(level: LogLevel, levels: number) { - const name = levels === 1 && level === LogLevel.unknown ? 'logs' : level; - const color = LogLevelColor[level]; - return { - displayNameFromDS: name, - color: { - mode: FieldColorModeId.Fixed, - fixedColor: color, - }, - custom: { - drawStyle: GraphDrawStyle.Bars, - barAlignment: BarAlignment.Center, - lineColor: color, - pointColor: color, - fillColor: color, - lineWidth: 1, - fillOpacity: 100, - stacking: { - mode: StackingMode.Normal, - group: 'A', - }, - }, - }; -} - -/** - * Create a new data frame with a single field and values creating by adding field values - * from all provided data frames - */ -function aggregateFields(dataFrames: DataFrame[], config: FieldConfig): DataFrame { - const aggregatedDataFrame = new MutableDataFrame(); - if (!dataFrames.length) { - return aggregatedDataFrame; - } - - const totalLength = dataFrames[0].length; - const timeField = new FieldCache(dataFrames[0]).getFirstFieldOfType(FieldType.time); - - if (!timeField) { - return aggregatedDataFrame; - } - - aggregatedDataFrame.addField({ name: 'Time', type: FieldType.time }, totalLength); - aggregatedDataFrame.addField({ name: 'Value', type: FieldType.number, config }, totalLength); - - dataFrames.forEach((dataFrame) => { - dataFrame.fields.forEach((field) => { - if (field.type === FieldType.number) { - for (let pointIndex = 0; pointIndex < totalLength; pointIndex++) { - const currentValue = aggregatedDataFrame.get(pointIndex).Value; - const valueToAdd = field.values.get(pointIndex); - const totalValue = - currentValue === null && valueToAdd === null ? null : (currentValue || 0) + (valueToAdd || 0); - aggregatedDataFrame.set(pointIndex, { Value: totalValue, Time: timeField.values.get(pointIndex) }); - } - } - }); - }); - - return aggregatedDataFrame; -} - -function getLogLevelFromLabels(labels: Labels): LogLevel { - const labelNames = ['level', 'lvl', 'loglevel']; - let levelLabel; - for (let labelName of labelNames) { - if (labelName in labels) { - levelLabel = labelName; - break; - } - } - return levelLabel ? getLogLevelFromKey(labels[levelLabel]) : LogLevel.unknown; -} - -function getIntervalInfo(scopedVars: ScopedVars): { interval: string; intervalMs?: number } { - if (scopedVars.__interval) { - let intervalMs: number = scopedVars.__interval_ms.value; - let interval = ''; - if (intervalMs > HOUR) { - intervalMs = DAY; - interval = '1d'; - } else if (intervalMs > MINUTE) { - intervalMs = HOUR; - interval = '1h'; - } else if (intervalMs > SECOND) { - intervalMs = MINUTE; - interval = '1m'; - } else { - intervalMs = SECOND; - interval = '1s'; - } - - return { interval, intervalMs }; - } else { - return { interval: '$__interval' }; - } -} diff --git a/public/app/plugins/datasource/loki/datasource.ts b/public/app/plugins/datasource/loki/datasource.ts index 9d28a4e4df8..7f00511671f 100644 --- a/public/app/plugins/datasource/loki/datasource.ts +++ b/public/app/plugins/datasource/loki/datasource.ts @@ -21,7 +21,11 @@ import { dateMath, DateTime, FieldCache, + FieldType, + getLogLevelFromKey, + Labels, LoadingState, + LogLevel, LogRowModel, QueryResultMeta, ScopedVars, @@ -54,13 +58,19 @@ import { serializeParams } from '../../../core/utils/fetch'; import { RowContextOptions } from '@grafana/ui/src/components/Logs/LogRowContextProvider'; import syntax from './syntax'; import { DEFAULT_RESOLUTION } from './components/LokiOptionFields'; -import { createLokiLogsVolumeProvider } from './dataProviders/logsVolumeProvider'; +import { queryLogsVolume } from 'app/core/logs_model'; export type RangeQueryOptions = DataQueryRequest | AnnotationQueryRequest; export const DEFAULT_MAX_LINES = 1000; export const LOKI_ENDPOINT = '/loki/api/v1'; const NS_IN_MS = 1000000; +/** + * Loki's logs volume query may be expensive as it requires counting all logs in the selected range. If such query + * takes too much time it may need be made more specific to limit number of logs processed under the hood. + */ +const LOGS_VOLUME_TIMEOUT = 10000; + const RANGE_QUERY_ENDPOINT = `${LOKI_ENDPOINT}/query_range`; const INSTANT_QUERY_ENDPOINT = `${LOKI_ENDPOINT}/query`; @@ -109,7 +119,27 @@ export class LokiDatasource getLogsVolumeDataProvider(request: DataQueryRequest): Observable | undefined { const isLogsVolumeAvailable = request.targets.some((target) => target.expr && !isMetricsQuery(target.expr)); - return isLogsVolumeAvailable ? createLokiLogsVolumeProvider(this, request) : undefined; + if (!isLogsVolumeAvailable) { + return undefined; + } + + const logsVolumeRequest = cloneDeep(request); + logsVolumeRequest.targets = logsVolumeRequest.targets + .filter((target) => target.expr && !isMetricsQuery(target.expr)) + .map((target) => { + return { + ...target, + instant: false, + expr: `sum by (level) (count_over_time(${target.expr}[$__interval]))`, + }; + }); + + return queryLogsVolume(this, logsVolumeRequest, { + timeout: LOGS_VOLUME_TIMEOUT, + extractLevel, + range: request.range, + targets: request.targets, + }); } query(options: DataQueryRequest): Observable { @@ -721,4 +751,24 @@ export function isMetricsQuery(query: string): boolean { }); } +function extractLevel(dataFrame: DataFrame): LogLevel { + let valueField; + try { + valueField = new FieldCache(dataFrame).getFirstFieldOfType(FieldType.number); + } catch {} + return valueField?.labels ? getLogLevelFromLabels(valueField.labels) : LogLevel.unknown; +} + +function getLogLevelFromLabels(labels: Labels): LogLevel { + const labelNames = ['level', 'lvl', 'loglevel']; + let levelLabel; + for (let labelName of labelNames) { + if (labelName in labels) { + levelLabel = labelName; + break; + } + } + return levelLabel ? getLogLevelFromKey(labels[levelLabel]) : LogLevel.unknown; +} + export default LokiDatasource;