Live: support streaming results out-of-the-box (#32821)

This commit is contained in:
Ryan McKinley
2021-04-09 21:17:22 +02:00
committed by GitHub
parent 2d7e980da7
commit b96e45299d
20 changed files with 179 additions and 242 deletions
+1 -1
View File
@@ -6,7 +6,7 @@
export * from './services';
export * from './config';
export * from './types';
export * from './measurement';
export * from './utils/liveQuery';
export { loadPluginCss, SystemJS, PluginCssOptions } from './utils/plugin';
export { reportMetaAnalytics } from './utils/analytics';
export { logInfo, logDebug, logWarning, logError } from './utils/logging';
@@ -1 +0,0 @@
export * from './query';
@@ -1,6 +1,13 @@
import { BackendSrv, BackendSrvRequest } from 'src/services';
import { DataSourceWithBackend } from './DataSourceWithBackend';
import { DataSourceJsonData, DataQuery, DataSourceInstanceSettings, DataQueryRequest } from '@grafana/data';
import { DataSourceWithBackend, toStreamingDataResponse } from './DataSourceWithBackend';
import {
DataSourceJsonData,
DataQuery,
DataSourceInstanceSettings,
DataQueryRequest,
DataQueryResponseData,
MutableDataFrame,
} from '@grafana/data';
import { of } from 'rxjs';
class MyDataSource extends DataSourceWithBackend<DataQuery, DataSourceJsonData> {
@@ -73,4 +80,26 @@ describe('DataSourceWithBackend', () => {
}
`);
});
test('it converts results with channels to streaming queries', () => {
const request: DataQueryRequest = {
intervalMs: 100,
} as DataQueryRequest;
const rsp: DataQueryResponseData = {
data: [],
};
// Simple empty query
let obs = toStreamingDataResponse(request, rsp);
expect(obs).toBeDefined();
let frame = new MutableDataFrame();
frame.meta = {
channel: 'a/b/c',
};
rsp.data = [frame];
obs = toStreamingDataResponse(request, rsp);
expect(obs).toBeDefined();
});
});
@@ -7,11 +7,15 @@ import {
DataSourceJsonData,
ScopedVars,
makeClassES5Compatible,
DataFrame,
parseLiveChannelAddress,
StreamingFrameOptions,
} from '@grafana/data';
import { Observable, of } from 'rxjs';
import { map, catchError } from 'rxjs/operators';
import { merge, Observable, of } from 'rxjs';
import { catchError, switchMap } from 'rxjs/operators';
import { getBackendSrv, getDataSourceSrv } from '../services';
import { BackendDataSourceResponse, toDataQueryResponse } from './queryResponse';
import { getLiveDataStream } from './liveQuery';
const ExpressionDatasourceID = '__expr__';
@@ -132,8 +136,13 @@ class DataSourceWithBackend<
requestId,
})
.pipe(
map((rsp) => {
return toDataQueryResponse(rsp, queries as DataQuery[]);
switchMap((raw) => {
const rsp = toDataQueryResponse(raw, queries as DataQuery[]);
// Check if any response should subscribe to a live stream
if (rsp.data?.length && rsp.data.find((f: DataFrame) => f.meta?.channel)) {
return toStreamingDataResponse(request, rsp);
}
return of(rsp);
}),
catchError((err) => {
return of(toDataQueryResponse(err));
@@ -209,6 +218,44 @@ class DataSourceWithBackend<
}
}
export function toStreamingDataResponse(
request: DataQueryRequest,
rsp: DataQueryResponse
): Observable<DataQueryResponse> {
const buffer: StreamingFrameOptions = {
maxLength: request.maxDataPoints ?? 500,
};
// For recent queries, clamp to the current time range
if (request.rangeRaw?.to === 'now') {
buffer.maxDelta = request.range.to.valueOf() - request.range.from.valueOf();
}
const staticdata: DataFrame[] = [];
const streams: Array<Observable<DataQueryResponse>> = [];
for (const frame of rsp.data) {
const addr = parseLiveChannelAddress(frame.meta?.channel);
if (addr) {
streams.push(
getLiveDataStream({
addr,
buffer,
frame: frame as DataFrame,
})
);
} else {
staticdata.push(frame);
}
}
if (staticdata.length) {
streams.push(of({ ...rsp, data: staticdata }));
}
if (streams.length === 1) {
return streams[0]; // avoid merge wrapper
}
return merge(...streams);
}
//@ts-ignore
DataSourceWithBackend = makeClassES5Compatible(DataSourceWithBackend);
@@ -1,6 +1,7 @@
import {
DataFrame,
DataFrameJSON,
dataFrameToJSON,
DataQueryResponse,
isLiveChannelMessageEvent,
isLiveChannelStatusEvent,
@@ -15,7 +16,7 @@ import {
import { getGrafanaLiveSrv } from '../services/live';
import { Observable, of } from 'rxjs';
import { toDataQueryError } from '../utils/queryResponse';
import { toDataQueryError } from './queryResponse';
import { perf } from './perf';
export interface LiveDataFilter {
@@ -28,6 +29,7 @@ export interface LiveDataFilter {
export interface LiveDataStreamOptions {
key?: string;
addr: LiveChannelAddress;
frame?: DataFrame; // initial results
buffer?: StreamingFrameOptions;
filter?: LiveDataFilter;
}
@@ -39,8 +41,13 @@ export interface LiveDataStreamOptions {
*/
export function getLiveDataStream(options: LiveDataStreamOptions): Observable<DataQueryResponse> {
if (!isValidLiveChannelAddress(options.addr)) {
return of({ error: toDataQueryError('invalid address'), data: [] });
return of({
error: toDataQueryError(`invalid channel address: ${JSON.stringify(options.addr)}`),
state: LoadingState.Error,
data: options.frame ? [options.frame] : [],
});
}
const live = getGrafanaLiveSrv();
if (!live) {
return of({ error: toDataQueryError('grafana live is not initalized'), data: [] });
@@ -50,8 +57,16 @@ export function getLiveDataStream(options: LiveDataStreamOptions): Observable<Da
let data: StreamingDataFrame | undefined = undefined;
let filtered: DataFrame | undefined = undefined;
let state = LoadingState.Loading;
const { key, filter } = options;
let { key } = options;
let last = perf.last;
if (options.frame) {
const msg = dataFrameToJSON(options.frame);
data = new StreamingDataFrame(msg, options.buffer);
state = LoadingState.Streaming;
}
if (!key) {
key = `xstr/${streamCounter++}`;
}
const process = (msg: DataFrameJSON) => {
if (!data) {
@@ -61,14 +76,17 @@ export function getLiveDataStream(options: LiveDataStreamOptions): Observable<Da
}
state = LoadingState.Streaming;
// Select the fields we are actually looking at
// Filter out fields
if (!filtered || msg.schema) {
filtered = data;
if (filter?.fields?.length) {
filtered = {
...data,
fields: data.fields.filter((f) => filter.fields!.includes(f.name)),
};
if (options.filter) {
const { fields } = options.filter;
if (fields?.length) {
filtered = {
...data,
fields: data.fields.filter((f) => fields.includes(f.name)),
};
}
}
}
@@ -85,15 +103,17 @@ export function getLiveDataStream(options: LiveDataStreamOptions): Observable<Da
.getStream()
.subscribe({
error: (err: any) => {
console.log('LiveQuery [error]', { err }, options.addr);
state = LoadingState.Error;
subscriber.next({ state, data: [data], key });
subscriber.next({ state, data: [data], key, error: toDataQueryError(err) });
sub.unsubscribe(); // close after error
},
complete: () => {
console.log('LiveQuery [complete]', options.addr);
if (state !== LoadingState.Error) {
state = LoadingState.Done;
}
subscriber.next({ state, data: [data], key });
// or track errors? subscriber.next({ state, data: [data], key });
subscriber.complete();
sub.unsubscribe();
},
@@ -103,14 +123,19 @@ export function getLiveDataStream(options: LiveDataStreamOptions): Observable<Da
return;
}
if (isLiveChannelStatusEvent(evt)) {
if (
if (evt.error) {
let error = toDataQueryError(evt.error);
error.message = `Streaming channel error: ${error.message}`;
state = LoadingState.Error;
subscriber.next({ state, data: [data], key, error });
return;
} else if (
evt.state === LiveChannelConnectionState.Connected ||
evt.state === LiveChannelConnectionState.Pending
) {
if (evt.message) {
process(evt.message);
}
return;
}
console.log('ignore state', evt);
}
@@ -122,3 +147,6 @@ export function getLiveDataStream(options: LiveDataStreamOptions): Observable<Da
};
});
}
// incremet the stream ids
let streamCounter = 10;