diff --git a/go.mod b/go.mod index a1160eac204..66599f00d24 100644 --- a/go.mod +++ b/go.mod @@ -45,7 +45,7 @@ require ( github.com/grafana/grafana-aws-sdk v0.4.0 github.com/grafana/grafana-live-sdk v0.0.4 github.com/grafana/grafana-plugin-model v0.0.0-20190930120109-1fc953a61fb4 - github.com/grafana/grafana-plugin-sdk-go v0.91.0 + github.com/grafana/grafana-plugin-sdk-go v0.92.0 github.com/grafana/loki v1.6.2-0.20201026154740-6978ee5d7387 github.com/grpc-ecosystem/go-grpc-middleware v1.2.2 github.com/hashicorp/go-hclog v0.15.0 @@ -103,4 +103,4 @@ require ( gopkg.in/yaml.v3 v3.0.0-20210107192922-496545a6307b xorm.io/core v0.7.3 xorm.io/xorm v0.8.2 -) \ No newline at end of file +) diff --git a/go.sum b/go.sum index 539f59ec842..31f88347e01 100644 --- a/go.sum +++ b/go.sum @@ -816,16 +816,6 @@ github.com/gorilla/websocket v1.4.2 h1:+/TMaTYc4QFitKJxsQ7Yye35DkWvkdLcvGKqM+x0U github.com/gorilla/websocket v1.4.2/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE= github.com/gosimple/slug v1.9.0 h1:r5vDcYrFz9BmfIAMC829un9hq7hKM4cHUrsv36LbEqs= github.com/gosimple/slug v1.9.0/go.mod h1:AMZ+sOVe65uByN3kgEyf9WEBKBCSS+dJjMX9x4vDJbg= -github.com/grafana/alerting-api v0.0.0-20210331135037-3294563b51bb h1:Hj25Whc/TRv0hSLm5VN0FJ5R4yZ6M4ycRcBgu7bsEAc= -github.com/grafana/alerting-api v0.0.0-20210331135037-3294563b51bb/go.mod h1:5IppnPguSHcCbVLGCVzVjBvuQZNbYgVJ4KyXXjhCyWY= -github.com/grafana/alerting-api v0.0.0-20210405171311-97906879c771 h1:CTmKHUu2n0O9fPTSXb+s5FO8Em9Atw57Z7mvw7lt6IM= -github.com/grafana/alerting-api v0.0.0-20210405171311-97906879c771/go.mod h1:5IppnPguSHcCbVLGCVzVjBvuQZNbYgVJ4KyXXjhCyWY= -github.com/grafana/alerting-api v0.0.0-20210407150830-64bd267999d1 h1:pbG8BsRHezUvUjMxwq+uZsx1ZMEQsfSj26KSd/H3A9g= -github.com/grafana/alerting-api v0.0.0-20210407150830-64bd267999d1/go.mod h1:Ce2PwraBlFMa+P0ArBzubfB/BXZV35mfYWQjM8C/BSE= -github.com/grafana/alerting-api v0.0.0-20210408130544-2969b61275a5 h1:r33Ruhf3TNvT4Y+acR8R0zBtNULnThnJRlAje2AFt6c= -github.com/grafana/alerting-api v0.0.0-20210408130544-2969b61275a5/go.mod h1:Ce2PwraBlFMa+P0ArBzubfB/BXZV35mfYWQjM8C/BSE= -github.com/grafana/alerting-api v0.0.0-20210408135630-bc10f737ceff h1:3l/P+TIOXJy1zX/kZB6dFAMPQrpwpiN3cLMVr8B8u4Q= -github.com/grafana/alerting-api v0.0.0-20210408135630-bc10f737ceff/go.mod h1:Ce2PwraBlFMa+P0ArBzubfB/BXZV35mfYWQjM8C/BSE= github.com/grafana/alerting-api v0.0.0-20210409134845-c36ac1eae41b h1:QG52Et3EVCxPoYZifm91bRPVknccfjQURcpi7zXVut8= github.com/grafana/alerting-api v0.0.0-20210409134845-c36ac1eae41b/go.mod h1:Ce2PwraBlFMa+P0ArBzubfB/BXZV35mfYWQjM8C/BSE= github.com/grafana/go-mssqldb v0.0.0-20210326084033-d0ce3c521036 h1:GplhUk6Xes5JIhUUrggPcPBhOn+eT8+WsHiebvq7GgA= @@ -840,8 +830,9 @@ github.com/grafana/grafana-plugin-model v0.0.0-20190930120109-1fc953a61fb4 h1:SP github.com/grafana/grafana-plugin-model v0.0.0-20190930120109-1fc953a61fb4/go.mod h1:nc0XxBzjeGcrMltCDw269LoWF9S8ibhgxolCdA1R8To= github.com/grafana/grafana-plugin-sdk-go v0.79.0/go.mod h1:NvxLzGkVhnoBKwzkst6CFfpMFKwAdIUZ1q8ssuLeF60= github.com/grafana/grafana-plugin-sdk-go v0.88.0/go.mod h1:PTALh0lz+Y7k0+OMczAABTpeocL63aw6FVOBptp5GVo= -github.com/grafana/grafana-plugin-sdk-go v0.91.0 h1:kiPS3NqR+KOvHrc32EkX7D40JON3+GYZ6Nm2WOtCElQ= github.com/grafana/grafana-plugin-sdk-go v0.91.0/go.mod h1:Ot3k7nY7P6DXmUsDgKvNB7oG1v7PRyTdmnYVoS554bU= +github.com/grafana/grafana-plugin-sdk-go v0.92.0 h1:Wim+Ey7BaA0BO7+wBHeTPrGotKPRznKBXYArSAJL3W8= +github.com/grafana/grafana-plugin-sdk-go v0.92.0/go.mod h1:UwW5HV5HUzRKCieDfC4J31H4PBQngn2wXjJXJmn2zjw= github.com/grafana/loki v1.6.2-0.20201026154740-6978ee5d7387 h1:iwcM8lkYJ3EhytGLJ2BvRSwutb0QWoI7EWbYv3yJRsY= github.com/grafana/loki v1.6.2-0.20201026154740-6978ee5d7387/go.mod h1:jHA1OHnPsuj3LLgMXmFopsKDt4ARHHUhrmT3JrGf71g= github.com/gregjones/httpcache v0.0.0-20180305231024-9cad4c3443a7/go.mod h1:FecbI9+v66THATjSRHfNgh1IVFe/9kFxbXtjV0ctIMA= @@ -1051,7 +1042,6 @@ github.com/jung-kurt/gofpdf v1.0.3-0.20190309125859-24315acbbda5/go.mod h1:7Id9E github.com/jung-kurt/gofpdf v1.16.2 h1:jgbatWHfRlPYiK85qgevsZTHviWXKwB1TTiKdz5PtRc= github.com/jung-kurt/gofpdf v1.16.2/go.mod h1:1hl7y57EsiPAkLbOwzpzqgx1A30nQCk/YmFV8S2vmK0= github.com/jwilder/encoding v0.0.0-20170811194829-b4e1701a28ef/go.mod h1:Ct9fl0F6iIOGgxJ5npU/IUOhOhqlVrGjyIZc8/MagT0= -github.com/k0kubun/colorstring v0.0.0-20150214042306-9440f1994b88 h1:uC1QfSlInpQF+M0ao65imhwqKnz3Q2z/d8PWZRMQvDM= github.com/k0kubun/colorstring v0.0.0-20150214042306-9440f1994b88/go.mod h1:3w7q1U84EfirKl04SVQ/s7nPm1ZPhiXd34z40TNz36k= github.com/k0kubun/go-ansi v0.0.0-20180517002512-3bf9e2903213/go.mod h1:vNUNkEQ1e29fT/6vq2aBdFsgNPmy8qMdSay1npru+Sw= github.com/kardianos/osext v0.0.0-20190222173326-2bc1f35cddc0/go.mod h1:1NbS8ALrpOvjt0rHPNLyCIeMtbizbir8U//inJ+zuB8= diff --git a/packages/grafana-data/src/types/data.ts b/packages/grafana-data/src/types/data.ts index 73ac7b82b07..92d3ff43172 100644 --- a/packages/grafana-data/src/types/data.ts +++ b/packages/grafana-data/src/types/data.ts @@ -42,6 +42,9 @@ export interface QueryResultMeta { /** Currently used to show results in Explore only in preferred visualisation option */ preferredVisualisationType?: PreferredVisualisationType; + /** The path for live stream updates for this frame */ + channel?: string; + /** * Optionally identify which topic the frame should be assigned to. * A value specified in the response will override what the request asked for. diff --git a/packages/grafana-runtime/src/index.ts b/packages/grafana-runtime/src/index.ts index 987e14f06d1..64bac2b85c6 100644 --- a/packages/grafana-runtime/src/index.ts +++ b/packages/grafana-runtime/src/index.ts @@ -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'; diff --git a/packages/grafana-runtime/src/measurement/index.ts b/packages/grafana-runtime/src/measurement/index.ts deleted file mode 100644 index a38910e42af..00000000000 --- a/packages/grafana-runtime/src/measurement/index.ts +++ /dev/null @@ -1 +0,0 @@ -export * from './query'; diff --git a/packages/grafana-runtime/src/utils/DataSourceWithBackend.test.ts b/packages/grafana-runtime/src/utils/DataSourceWithBackend.test.ts index 9cc83b45edb..236d90ffaee 100644 --- a/packages/grafana-runtime/src/utils/DataSourceWithBackend.test.ts +++ b/packages/grafana-runtime/src/utils/DataSourceWithBackend.test.ts @@ -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 { @@ -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(); + }); }); diff --git a/packages/grafana-runtime/src/utils/DataSourceWithBackend.ts b/packages/grafana-runtime/src/utils/DataSourceWithBackend.ts index 60a8d6fe27e..3608bfb7979 100644 --- a/packages/grafana-runtime/src/utils/DataSourceWithBackend.ts +++ b/packages/grafana-runtime/src/utils/DataSourceWithBackend.ts @@ -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 { + 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> = []; + 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); diff --git a/packages/grafana-runtime/src/measurement/query.ts b/packages/grafana-runtime/src/utils/liveQuery.ts similarity index 65% rename from packages/grafana-runtime/src/measurement/query.ts rename to packages/grafana-runtime/src/utils/liveQuery.ts index 30847540ec5..68e1190a8ea 100644 --- a/packages/grafana-runtime/src/measurement/query.ts +++ b/packages/grafana-runtime/src/utils/liveQuery.ts @@ -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 { 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 { if (!data) { @@ -61,14 +76,17 @@ export function getLiveDataStream(options: LiveDataStreamOptions): Observable 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 { + 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 0 { - addr.Scope = parts[0] - } - if length > 1 { - addr.Namespace = parts[1] - } - if length > 2 { - addr.Path = parts[2] - } - return addr -} - -func (c Channel) String() string { - ch := c.Scope - if c.Namespace != "" { - ch += "/" + c.Namespace - } - if c.Path != "" { - ch += "/" + c.Path - } - return ch -} - -// IsValid checks if all parts of the address are valid. -func (c *Channel) IsValid() bool { - if c.Scope == ScopePush { - // Push scope channels supposed to be like push/{$stream_id}. - return c.Namespace != "" && c.Path == "" - } - return c.Scope != "" && c.Namespace != "" && c.Path != "" -} diff --git a/pkg/services/live/channel_test.go b/pkg/services/live/channel_test.go deleted file mode 100644 index 580b0fe6bf9..00000000000 --- a/pkg/services/live/channel_test.go +++ /dev/null @@ -1,110 +0,0 @@ -package live - -import ( - "testing" - - "github.com/google/go-cmp/cmp" - "github.com/stretchr/testify/require" -) - -func TestParseChannel(t *testing.T) { - addr := ParseChannel("aaa/bbb/ccc/ddd") - require.True(t, addr.IsValid()) - - ex := Channel{ - Scope: "aaa", - Namespace: "bbb", - Path: "ccc/ddd", - } - - if diff := cmp.Diff(addr, ex); diff != "" { - t.Fatalf("Result mismatch (-want +got):\n%s", diff) - } -} - -func TestParseChannel_IsValid(t *testing.T) { - tests := []struct { - name string - id string - isValid bool - }{ - { - name: "valid", - id: "stream/cpu/test", - isValid: true, - }, - { - name: "valid_long_path", - id: "stream/cpu/test/other", - isValid: true, - }, - { - name: "invalid_no_path", - id: "grafana/bbb", - isValid: false, - }, - { - name: "invalid_only_scope", - id: "grafana", - isValid: false, - }, - { - name: "push_scope_no_path_valid", - id: "push/telegraf", - isValid: true, - }, - { - name: "push_scope_with_path_invalid", - id: "push/telegraf/test", - isValid: false, - }, - } - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - if got := ParseChannel(tt.id); got.IsValid() != tt.isValid { - t.Errorf("unexpected isValid result for %s", tt.id) - } - }) - } -} - -func TestChannel_String(t *testing.T) { - type fields struct { - Scope string - Namespace string - Path string - } - tests := []struct { - name string - fields fields - want string - }{ - { - "with_all_parts", - fields{Scope: ScopeStream, Namespace: "telegraf", Path: "test"}, - "stream/telegraf/test", - }, - { - "with_scope_and_namespace", - fields{Scope: ScopeStream, Namespace: "telegraf"}, - "stream/telegraf", - }, - { - "with_scope_only", - fields{Scope: ScopeStream}, - "stream", - }, - } - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - got := Channel{ - Scope: tt.fields.Scope, - Namespace: tt.fields.Namespace, - Path: tt.fields.Path, - }.String() - if got != tt.want { - t.Errorf("String() = %v, want %v", got, tt.want) - } - }) - } -} diff --git a/pkg/services/live/features/plugin.go b/pkg/services/live/features/plugin.go index cd16d6e0108..b5d48662900 100644 --- a/pkg/services/live/features/plugin.go +++ b/pkg/services/live/features/plugin.go @@ -100,7 +100,7 @@ func (r *PluginPathRunner) OnSubscribe(ctx context.Context, user *models.SignedI Path: r.path, }) if err != nil { - logger.Error("Plugin CanSubscribeToStream call error", "error", err, "path", r.path) + logger.Error("Plugin OnSubscribe call error", "error", err, "path", r.path) return models.SubscribeReply{}, 0, err } if resp.Status != backend.SubscribeStreamStatusOK { @@ -144,7 +144,7 @@ func (r *PluginPathRunner) OnPublish(ctx context.Context, user *models.SignedInU Data: e.Data, }) if err != nil { - logger.Error("Plugin CanSubscribeToStream call error", "error", err, "path", r.path) + logger.Error("Plugin OnPublish call error", "error", err, "path", r.path) return models.PublishReply{}, 0, err } if resp.Status != backend.PublishStreamStatusOK { diff --git a/pkg/services/live/live.go b/pkg/services/live/live.go index b4e33128214..c5a59ec6bd1 100644 --- a/pkg/services/live/live.go +++ b/pkg/services/live/live.go @@ -4,12 +4,14 @@ import ( "context" "fmt" "net/http" + "strconv" "sync" "time" "github.com/centrifugal/centrifuge" "github.com/grafana/grafana-plugin-sdk-go/backend" "github.com/grafana/grafana-plugin-sdk-go/data" + "github.com/grafana/grafana-plugin-sdk-go/live" "github.com/grafana/grafana/pkg/api/dtos" "github.com/grafana/grafana/pkg/api/response" @@ -306,11 +308,11 @@ func publishStatusToHTTPError(status backend.PublishStreamStatus) (int, string) } // GetChannelHandler gives thread-safe access to the channel. -func (g *GrafanaLive) GetChannelHandler(user *models.SignedInUser, channel string) (models.ChannelHandler, Channel, error) { +func (g *GrafanaLive) GetChannelHandler(user *models.SignedInUser, channel string) (models.ChannelHandler, live.Channel, error) { // Parse the identifier ${scope}/${namespace}/${path} - addr := ParseChannel(channel) + addr := live.ParseChannel(channel) if !addr.IsValid() { - return nil, Channel{}, fmt.Errorf("invalid channel: %q", channel) + return nil, live.Channel{}, fmt.Errorf("invalid channel: %q", channel) } g.channelsMu.RLock() @@ -349,15 +351,15 @@ func (g *GrafanaLive) GetChannelHandler(user *models.SignedInUser, channel strin // It gives thread-safe access to the channel. func (g *GrafanaLive) GetChannelHandlerFactory(user *models.SignedInUser, scope string, namespace string) (models.ChannelHandlerFactory, error) { switch scope { - case ScopeGrafana: + case live.ScopeGrafana: return g.handleGrafanaScope(user, namespace) - case ScopePlugin: + case live.ScopePlugin: return g.handlePluginScope(user, namespace) - case ScopeDatasource: + case live.ScopeDatasource: return g.handleDatasourceScope(user, namespace) - case ScopeStream: + case live.ScopeStream: return g.handleStreamScope(user, namespace) - case ScopePush: + case live.ScopePush: return g.handlePushScope(user, namespace) default: return nil, fmt.Errorf("invalid scope: %q", scope) @@ -403,7 +405,14 @@ func (g *GrafanaLive) handlePushScope(_ *models.SignedInUser, namespace string) func (g *GrafanaLive) handleDatasourceScope(user *models.SignedInUser, namespace string) (models.ChannelHandlerFactory, error) { ds, err := g.DatasourceCache.GetDatasourceByUID(namespace, user, false) if err != nil { - return nil, fmt.Errorf("error getting datasource: %w", err) + // the namespace may be an ID + id, _ := strconv.ParseInt(namespace, 10, 64) + if id > 0 { + ds, err = g.DatasourceCache.GetDatasource(id, user, false) + } + if err != nil { + return nil, fmt.Errorf("error getting datasource: %w", err) + } } streamHandler, err := g.getStreamPlugin(ds.Type) if err != nil { @@ -430,7 +439,7 @@ func (g *GrafanaLive) IsEnabled() bool { } func (g *GrafanaLive) HandleHTTPPublish(ctx *models.ReqContext, cmd dtos.LivePublishCmd) response.Response { - addr := ParseChannel(cmd.Channel) + addr := live.ParseChannel(cmd.Channel) if !addr.IsValid() { return response.Error(http.StatusBadRequest, "Bad channel address", nil) } diff --git a/pkg/services/live/managed.go b/pkg/services/live/managed.go index fd089f89744..edee04b1641 100644 --- a/pkg/services/live/managed.go +++ b/pkg/services/live/managed.go @@ -8,6 +8,7 @@ import ( "github.com/grafana/grafana-plugin-sdk-go/backend" "github.com/grafana/grafana-plugin-sdk-go/data" + "github.com/grafana/grafana-plugin-sdk-go/live" "github.com/grafana/grafana/pkg/models" "github.com/grafana/grafana/pkg/util" ) @@ -111,7 +112,7 @@ func (s *ManagedStream) Push(path string, frame *data.Frame) error { } // The channel this will be posted into. - channel := Channel{Scope: ScopeStream, Namespace: s.id, Path: path}.String() + channel := live.Channel{Scope: live.ScopeStream, Namespace: s.id, Path: path}.String() logger.Debug("Publish data to channel", "channel", channel, "dataLength", len(frameJSON)) return s.publisher(channel, frameJSON) } diff --git a/pkg/services/live/scope.go b/pkg/services/live/scope.go deleted file mode 100644 index 7c64a3ad73f..00000000000 --- a/pkg/services/live/scope.go +++ /dev/null @@ -1,14 +0,0 @@ -package live - -const ( - // ScopeGrafana contains builtin features of Grafana Core. - ScopeGrafana = "grafana" - // ScopePlugin passes control to a plugin. - ScopePlugin = "plugin" - // ScopeDatasource passes control to a datasource plugin. - ScopeDatasource = "ds" - // ScopeStream is a managed data frame stream. - ScopeStream = "stream" - // ScopePush allows sending data into managed streams. It does not support subscriptions. - ScopePush = "push" -) diff --git a/public/app/features/live/channel.ts b/public/app/features/live/channel.ts index 83171adba69..762c026ed5d 100644 --- a/public/app/features/live/channel.ts +++ b/public/app/features/live/channel.ts @@ -178,6 +178,7 @@ export class CentrifugeLiveChannel implements Li shutdownWithError(err: string) { this.currentStatus.error = err; + this.sendStatus(); this.disconnect(); } } diff --git a/public/app/features/live/live.ts b/public/app/features/live/live.ts index e7897d85464..020b5883641 100644 --- a/public/app/features/live/live.ts +++ b/public/app/features/live/live.ts @@ -1,7 +1,7 @@ import Centrifuge from 'centrifuge/dist/centrifuge'; import { GrafanaLiveSrv, setGrafanaLiveSrv, getGrafanaLiveSrv, config } from '@grafana/runtime'; import { BehaviorSubject } from 'rxjs'; -import { LiveChannel, LiveChannelScope, LiveChannelAddress } from '@grafana/data'; +import { LiveChannel, LiveChannelScope, LiveChannelAddress, LiveChannelConnectionState } from '@grafana/data'; import { CentrifugeLiveChannel, getErrorChannel } from './channel'; import { GrafanaLiveScope, @@ -104,7 +104,10 @@ export class CentrifugeSrv implements GrafanaLiveSrv { // Initialize the channel in the background this.initChannel(scope, channel).catch((err) => { - channel?.shutdownWithError(err); + if (channel) { + channel.currentStatus.state = LiveChannelConnectionState.Invalid; + channel.shutdownWithError(err); + } this.open.delete(id); }); @@ -116,7 +119,7 @@ export class CentrifugeSrv implements GrafanaLiveSrv { const { addr } = channel; const support = await scope.getChannelSupport(addr.namespace); if (!support) { - throw new Error(channel.addr.namespace + 'does not support streaming'); + throw new Error(channel.addr.namespace + ' does not support streaming'); } const config = support.getChannelConfig(addr.path); if (!config) { diff --git a/public/app/features/live/scopes.ts b/public/app/features/live/scopes.ts index ebda62d7305..900389d0f03 100644 --- a/public/app/features/live/scopes.ts +++ b/public/app/features/live/scopes.ts @@ -73,7 +73,10 @@ export class GrafanaLiveDataSourceScope extends GrafanaLiveScope { */ async getChannelSupport(namespace: string) { const ds = await getDataSourceSrv().get(namespace); - return ds.channelSupport; + if (ds.channelSupport) { + return ds.channelSupport; + } + return new LiveMeasurementsSupport(); // default support? } /** @@ -119,10 +122,13 @@ export class GrafanaLivePluginScope extends GrafanaLiveScope { */ async getChannelSupport(namespace: string) { const plugin = await loadPlugin(namespace); - if (!plugin.channelSupport) { - throw new Error('Unknown plugin: ' + namespace); + if (!plugin) { + throw new Error('Unknown streaming plugin: ' + namespace); } - return plugin.channelSupport; + if (plugin.channelSupport) { + return plugin.channelSupport; // explicit + } + throw new Error('Plugin does not support streaming: ' + namespace); } /** diff --git a/public/app/features/plugins/datasource_srv.ts b/public/app/features/plugins/datasource_srv.ts index c09318c59a6..eec08f7b721 100644 --- a/public/app/features/plugins/datasource_srv.ts +++ b/public/app/features/plugins/datasource_srv.ts @@ -21,6 +21,7 @@ export class DatasourceSrv implements DataSourceService { private datasources: Record = {}; private settingsMapByName: Record = {}; private settingsMapByUid: Record = {}; + private settingsMapById: Record = {}; private defaultName = ''; /** @ngInject */ @@ -38,6 +39,7 @@ export class DatasourceSrv implements DataSourceService { for (const dsSettings of Object.values(settingsMapByName)) { this.settingsMapByUid[dsSettings.uid] = dsSettings; + this.settingsMapById[dsSettings.id] = dsSettings; } } @@ -115,9 +117,12 @@ export class DatasourceSrv implements DataSourceService { return Promise.resolve(expressionDatasource); } - const dsConfig = this.settingsMapByName[name]; + let dsConfig = this.settingsMapByName[name]; if (!dsConfig) { - return Promise.reject({ message: `Datasource named ${name} was not found` }); + dsConfig = this.settingsMapById[name]; + if (!dsConfig) { + return Promise.reject({ message: `Datasource named ${name} was not found` }); + } } try { diff --git a/public/app/plugins/datasource/testdata/runStreams.ts b/public/app/plugins/datasource/testdata/runStreams.ts index 1100aa74593..7dca5aa1bb5 100644 --- a/public/app/plugins/datasource/testdata/runStreams.ts +++ b/public/app/plugins/datasource/testdata/runStreams.ts @@ -16,7 +16,7 @@ import { import { TestDataQuery, StreamingQuery } from './types'; import { getRandomLine } from './LogIpsum'; -import { perf } from '@grafana/runtime/src/measurement/perf'; // not exported +import { perf } from '@grafana/runtime/src/utils/perf'; // not exported export const defaultStreamQuery: StreamingQuery = { type: 'signal',