From 63af403897269b3a58a83a394983b41f972c0ec5 Mon Sep 17 00:00:00 2001 From: Ryan McKinley Date: Mon, 7 Apr 2025 17:47:35 +0300 Subject: [PATCH] Live: Remove queryOverLive and live-service-web-worker experimental feature flags (#103518) --- .betterer.results | 3 - .../feature-toggles/index.md | 2 - .../src/types/featureToggles.gen.ts | 8 --- packages/grafana-runtime/src/services/live.ts | 9 --- .../src/utils/DataSourceWithBackend.ts | 7 -- pkg/services/featuremgmt/registry.go | 14 ---- pkg/services/featuremgmt/toggles_gen.csv | 2 - pkg/services/featuremgmt/toggles_gen.go | 8 --- pkg/services/featuremgmt/toggles_gen.json | 6 +- .../createCentrifugeServiceWorker.ts | 3 - .../live/centrifuge/remoteObservable.ts | 33 ---------- .../live/centrifuge/service.worker.ts | 64 ------------------- .../live/centrifuge/serviceWorkerProxy.ts | 54 ---------------- .../live/centrifuge/transferHandlers.ts | 27 -------- public/app/features/live/index.ts | 5 +- public/app/features/live/live.ts | 27 +------- public/test/mocks/workers.ts | 9 --- 17 files changed, 8 insertions(+), 273 deletions(-) delete mode 100644 public/app/features/live/centrifuge/createCentrifugeServiceWorker.ts delete mode 100644 public/app/features/live/centrifuge/remoteObservable.ts delete mode 100644 public/app/features/live/centrifuge/service.worker.ts delete mode 100644 public/app/features/live/centrifuge/serviceWorkerProxy.ts delete mode 100644 public/app/features/live/centrifuge/transferHandlers.ts diff --git a/.betterer.results b/.betterer.results index 3612d3f9a76..0e94061d6ec 100644 --- a/.betterer.results +++ b/.betterer.results @@ -2653,9 +2653,6 @@ exports[`better eslint`] = { "public/app/features/live/centrifuge/channel.ts:5381": [ [0, 0, 0, "Unexpected any. Specify a different type.", "0"] ], - "public/app/features/live/centrifuge/serviceWorkerProxy.ts:5381": [ - [0, 0, 0, "Do not use any type assertions.", "0"] - ], "public/app/features/live/dashboard/DashboardChangedModal.tsx:5381": [ [0, 0, 0, "No untranslated strings. Wrap text with ", "0"] ], diff --git a/docs/sources/setup-grafana/configure-grafana/feature-toggles/index.md b/docs/sources/setup-grafana/configure-grafana/feature-toggles/index.md index dcc1a40655e..769d8fb18e6 100644 --- a/docs/sources/setup-grafana/configure-grafana/feature-toggles/index.md +++ b/docs/sources/setup-grafana/configure-grafana/feature-toggles/index.md @@ -124,8 +124,6 @@ Experimental features might be changed or removed without prior notice. | Feature toggle name | Description | | ------------------------------------------- | --------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -| `live-service-web-worker` | This will use a webworker thread to processes events rather than the main thread | -| `queryOverLive` | Use Grafana Live WebSocket to execute backend queries | | `lokiExperimentalStreaming` | Support new streaming approach for loki (prototype, needs special loki build) | | `storage` | Configurable storage for dashboards, datasources, and resources | | `canvasPanelNesting` | Allow elements nesting | diff --git a/packages/grafana-data/src/types/featureToggles.gen.ts b/packages/grafana-data/src/types/featureToggles.gen.ts index bb7cf35fac7..2c806d13802 100644 --- a/packages/grafana-data/src/types/featureToggles.gen.ts +++ b/packages/grafana-data/src/types/featureToggles.gen.ts @@ -24,14 +24,6 @@ export interface FeatureToggles { */ disableEnvelopeEncryption?: boolean; /** - * This will use a webworker thread to processes events rather than the main thread - */ - ['live-service-web-worker']?: boolean; - /** - * Use Grafana Live WebSocket to execute backend queries - */ - queryOverLive?: boolean; - /** * Search for dashboards using panel title */ panelTitleSearch?: boolean; diff --git a/packages/grafana-runtime/src/services/live.ts b/packages/grafana-runtime/src/services/live.ts index 41c05ac76bd..0d8ae155cf1 100644 --- a/packages/grafana-runtime/src/services/live.ts +++ b/packages/grafana-runtime/src/services/live.ts @@ -72,15 +72,6 @@ export interface GrafanaLiveSrv { */ getDataStream(options: LiveDataStreamOptions): Observable; - /** - * Execute a query over the live websocket and potentiall subscribe to a live channel. - * - * Since the initial request and subscription are on the same socket, this will support HA setups - * - * @alpha -- this function requires the feature toggle `queryOverLive` to be set - */ - getQueryData(options: LiveQueryDataOptions): Observable; - /** * For channels that support presence, this will request the current state from the server. * diff --git a/packages/grafana-runtime/src/utils/DataSourceWithBackend.ts b/packages/grafana-runtime/src/utils/DataSourceWithBackend.ts index 7ce7e9aa52d..3d5a32dc252 100644 --- a/packages/grafana-runtime/src/utils/DataSourceWithBackend.ts +++ b/packages/grafana-runtime/src/utils/DataSourceWithBackend.ts @@ -196,13 +196,6 @@ class DataSourceWithBackend< to: range?.to.valueOf().toString(), }; - if (config.featureToggles.queryOverLive) { - return getGrafanaLiveSrv().getQueryData({ - request, - body, - }); - } - const headers: Record = request.headers ?? {}; headers[PluginRequestHeaders.PluginID] = Array.from(pluginIDs).join(', '); headers[PluginRequestHeaders.DatasourceUID] = Array.from(dsUIDs).join(', '); diff --git a/pkg/services/featuremgmt/registry.go b/pkg/services/featuremgmt/registry.go index a978456034e..140e758fe68 100644 --- a/pkg/services/featuremgmt/registry.go +++ b/pkg/services/featuremgmt/registry.go @@ -27,20 +27,6 @@ var ( AllowSelfServe: false, Expression: "false", }, - { - Name: "live-service-web-worker", - Description: "This will use a webworker thread to processes events rather than the main thread", - Stage: FeatureStageExperimental, - FrontendOnly: true, - Owner: grafanaDashboardsSquad, - }, - { - Name: "queryOverLive", - Description: "Use Grafana Live WebSocket to execute backend queries", - Stage: FeatureStageExperimental, - FrontendOnly: true, - Owner: grafanaDashboardsSquad, - }, { Name: "panelTitleSearch", Description: "Search for dashboards using panel title", diff --git a/pkg/services/featuremgmt/toggles_gen.csv b/pkg/services/featuremgmt/toggles_gen.csv index 9af21c53f96..a8f471f7422 100644 --- a/pkg/services/featuremgmt/toggles_gen.csv +++ b/pkg/services/featuremgmt/toggles_gen.csv @@ -1,7 +1,5 @@ Name,Stage,Owner,requiresDevMode,RequiresRestart,FrontendOnly disableEnvelopeEncryption,GA,@grafana/grafana-as-code,false,false,false -live-service-web-worker,experimental,@grafana/dashboards-squad,false,false,true -queryOverLive,experimental,@grafana/dashboards-squad,false,false,true panelTitleSearch,preview,@grafana/search-and-storage,false,false,false publicDashboardsEmailSharing,preview,@grafana/sharing-squad,false,false,false publicDashboardsScene,GA,@grafana/sharing-squad,false,false,true diff --git a/pkg/services/featuremgmt/toggles_gen.go b/pkg/services/featuremgmt/toggles_gen.go index 19496aea6d0..97cd485c30f 100644 --- a/pkg/services/featuremgmt/toggles_gen.go +++ b/pkg/services/featuremgmt/toggles_gen.go @@ -11,14 +11,6 @@ const ( // Disable envelope encryption (emergency only) FlagDisableEnvelopeEncryption = "disableEnvelopeEncryption" - // FlagLiveServiceWebWorker - // This will use a webworker thread to processes events rather than the main thread - FlagLiveServiceWebWorker = "live-service-web-worker" - - // FlagQueryOverLive - // Use Grafana Live WebSocket to execute backend queries - FlagQueryOverLive = "queryOverLive" - // FlagPanelTitleSearch // Search for dashboards using panel title FlagPanelTitleSearch = "panelTitleSearch" diff --git a/pkg/services/featuremgmt/toggles_gen.json b/pkg/services/featuremgmt/toggles_gen.json index 4de4884bddb..a135cd05cfa 100644 --- a/pkg/services/featuremgmt/toggles_gen.json +++ b/pkg/services/featuremgmt/toggles_gen.json @@ -1726,7 +1726,8 @@ "metadata": { "name": "live-service-web-worker", "resourceVersion": "1743693517832", - "creationTimestamp": "2022-01-26T17:44:20Z" + "creationTimestamp": "2022-01-26T17:44:20Z", + "deletionTimestamp": "2025-04-07T09:56:17Z" }, "spec": { "description": "This will use a webworker thread to processes events rather than the main thread", @@ -2588,7 +2589,8 @@ "metadata": { "name": "queryOverLive", "resourceVersion": "1743693517832", - "creationTimestamp": "2022-01-26T17:44:20Z" + "creationTimestamp": "2022-01-26T17:44:20Z", + "deletionTimestamp": "2025-04-07T09:56:17Z" }, "spec": { "description": "Use Grafana Live WebSocket to execute backend queries", diff --git a/public/app/features/live/centrifuge/createCentrifugeServiceWorker.ts b/public/app/features/live/centrifuge/createCentrifugeServiceWorker.ts deleted file mode 100644 index c7c68f5b0e7..00000000000 --- a/public/app/features/live/centrifuge/createCentrifugeServiceWorker.ts +++ /dev/null @@ -1,3 +0,0 @@ -import { CorsWorker as Worker } from 'app/core/utils/CorsWorker'; - -export const createWorker = () => new Worker(new URL('./service.worker.ts', import.meta.url)); diff --git a/public/app/features/live/centrifuge/remoteObservable.ts b/public/app/features/live/centrifuge/remoteObservable.ts deleted file mode 100644 index 20b75b96132..00000000000 --- a/public/app/features/live/centrifuge/remoteObservable.ts +++ /dev/null @@ -1,33 +0,0 @@ -import * as comlink from 'comlink'; -import { from, Observable, switchMap } from 'rxjs'; - -export const remoteObservableAsObservable = (remoteObs: comlink.RemoteObject>): Observable => - new Observable((subscriber) => { - // Passing the callbacks as 3 separate arguments is deprecated, but it's the only option for now - // - // RxJS recreates the functions via `Function.bind` https://github.com/ReactiveX/rxjs/blob/62aca850a37f598b5db6085661e0594b81ec4281/src/internal/Subscriber.ts#L169 - // and thus erases the ProxyMarker created via comlink.proxy(fN) when the callbacks - // are grouped together in a Observer object (ie. { next: (v) => ..., error: (err) => ..., complete: () => ... }) - // - // solution: TBD (autoproxy all functions?) - const remoteSubPromise = remoteObs.subscribe( - comlink.proxy((nextValueInRemoteObs: T) => { - subscriber.next(nextValueInRemoteObs); - }), - comlink.proxy((err: unknown) => { - subscriber.error(err); - }), - comlink.proxy(() => { - subscriber.complete(); - }) - ); - return { - unsubscribe: () => { - remoteSubPromise.then((remoteSub) => remoteSub.unsubscribe()); - }, - }; - }); - -export const promiseWithRemoteObservableAsObservable = ( - promiseWithProxyObservable: Promise>> -): Observable => from(promiseWithProxyObservable).pipe(switchMap((val) => remoteObservableAsObservable(val))); diff --git a/public/app/features/live/centrifuge/service.worker.ts b/public/app/features/live/centrifuge/service.worker.ts deleted file mode 100644 index 0c535bc9cb8..00000000000 --- a/public/app/features/live/centrifuge/service.worker.ts +++ /dev/null @@ -1,64 +0,0 @@ -import './transferHandlers'; - -import * as comlink from 'comlink'; - -import { LiveChannelAddress } from '@grafana/data'; -import { LiveDataStreamOptions, LivePublishOptions, LiveQueryDataOptions } from '@grafana/runtime'; - -import { remoteObservableAsObservable } from './remoteObservable'; -import { CentrifugeService, CentrifugeSrvDeps } from './service'; - -let centrifuge: CentrifugeService; - -const initialize = ( - deps: CentrifugeSrvDeps, - remoteDataStreamSubscriberReadiness: comlink.RemoteObject< - CentrifugeSrvDeps['dataStreamSubscriberReadiness'] & comlink.ProxyMarked - > -) => { - centrifuge = new CentrifugeService({ - ...deps, - dataStreamSubscriberReadiness: remoteObservableAsObservable(remoteDataStreamSubscriberReadiness), - }); -}; - -const getConnectionState = () => { - return comlink.proxy(centrifuge.getConnectionState()); -}; - -const getDataStream = (options: LiveDataStreamOptions) => { - return comlink.proxy(centrifuge.getDataStream(options)); -}; - -const getQueryData = async (options: LiveQueryDataOptions) => { - return await centrifuge.getQueryData(options); -}; - -const getStream = (address: LiveChannelAddress) => { - return comlink.proxy(centrifuge.getStream(address)); -}; - -const getPresence = async (address: LiveChannelAddress) => { - return await centrifuge.getPresence(address); -}; - -const publish = (address: LiveChannelAddress, data: unknown, options?: LivePublishOptions) => - centrifuge.publish(address, data, options); - -const workObj = { - initialize, - getConnectionState, - getDataStream, - getStream, - getQueryData, - getPresence, - publish, -}; - -export type RemoteCentrifugeService = typeof workObj; - -comlink.expose(workObj); - -export default class { - constructor() {} -} diff --git a/public/app/features/live/centrifuge/serviceWorkerProxy.ts b/public/app/features/live/centrifuge/serviceWorkerProxy.ts deleted file mode 100644 index 6f1b8fcc92d..00000000000 --- a/public/app/features/live/centrifuge/serviceWorkerProxy.ts +++ /dev/null @@ -1,54 +0,0 @@ -import './transferHandlers'; - -import * as comlink from 'comlink'; -import { asyncScheduler, Observable, observeOn } from 'rxjs'; - -import { LiveChannelAddress, LiveChannelEvent } from '@grafana/data'; - -import { createWorker } from './createCentrifugeServiceWorker'; -import { promiseWithRemoteObservableAsObservable } from './remoteObservable'; -import { CentrifugeSrv, CentrifugeSrvDeps } from './service'; -import { RemoteCentrifugeService } from './service.worker'; - -export class CentrifugeServiceWorkerProxy implements CentrifugeSrv { - private centrifugeWorker; - - constructor(deps: CentrifugeSrvDeps) { - this.centrifugeWorker = comlink.wrap(createWorker()); - this.centrifugeWorker.initialize(deps, comlink.proxy(deps.dataStreamSubscriberReadiness)); - } - - getConnectionState: CentrifugeSrv['getConnectionState'] = () => { - return promiseWithRemoteObservableAsObservable(this.centrifugeWorker.getConnectionState()); - }; - - getDataStream: CentrifugeSrv['getDataStream'] = (options) => { - return promiseWithRemoteObservableAsObservable(this.centrifugeWorker.getDataStream(options)).pipe( - // async scheduler splits the synchronous task of deserializing data from web worker and - // consuming the message (ie. updating react component) into two to avoid blocking the event loop - observeOn(asyncScheduler) - ); - }; - - /** - * Query over websocket - */ - getQueryData: CentrifugeSrv['getQueryData'] = async (options) => { - const optionsAsPlainSerializableObject = JSON.parse(JSON.stringify(options)); - return this.centrifugeWorker.getQueryData(optionsAsPlainSerializableObject); - }; - - getPresence: CentrifugeSrv['getPresence'] = (address) => { - return this.centrifugeWorker.getPresence(address); - }; - - getStream: CentrifugeSrv['getStream'] = (address: LiveChannelAddress) => { - return promiseWithRemoteObservableAsObservable( - this.centrifugeWorker.getStream(address) as Promise>>> - ); - }; - - publish: CentrifugeSrv['publish'] = (address, data, options) => { - return this.centrifugeWorker.publish(address, data, options); - }; -} diff --git a/public/app/features/live/centrifuge/transferHandlers.ts b/public/app/features/live/centrifuge/transferHandlers.ts deleted file mode 100644 index dafda35a4f6..00000000000 --- a/public/app/features/live/centrifuge/transferHandlers.ts +++ /dev/null @@ -1,27 +0,0 @@ -import * as comlink from 'comlink'; -import { Subscriber } from 'rxjs'; - -// Observers, ie. functions passed to `observable.subscribe(...)`, are converted to a subclass of `Subscriber` before they are sent to the source Observable. -// The conversion happens internally in the RxJS library - this transfer handler is catches them and wraps them with a proxy -const subscriberTransferHandler = { - canHandle(value: unknown): value is Subscriber { - return Boolean(value && value instanceof Subscriber); - }, - - serialize(value: Function): [MessagePort, Transferable[]] { - const obj = comlink.proxy(value); - - const { port1, port2 } = new MessageChannel(); - - comlink.expose(obj, port1); - - return [port2, [port2]]; - }, - - deserialize(value: MessagePort): comlink.Remote { - value.start(); - - return comlink.wrap(value); - }, -}; -comlink.transferHandlers.set('SubscriberHandler', subscriberTransferHandler); diff --git a/public/app/features/live/index.ts b/public/app/features/live/index.ts index f434a7149fd..69ca8d9d512 100644 --- a/public/app/features/live/index.ts +++ b/public/app/features/live/index.ts @@ -5,7 +5,6 @@ import { contextSrv } from '../../core/services/context_srv'; import { loadUrlToken } from '../../core/utils/urlToken'; import { CentrifugeService } from './centrifuge/service'; -import { CentrifugeServiceWorkerProxy } from './centrifuge/serviceWorkerProxy'; import { GrafanaLiveService } from './live'; export function initGrafanaLive() { @@ -18,9 +17,7 @@ export function initGrafanaLive() { grafanaAuthToken: loadUrlToken(), }; - const centrifugeSrv = config.featureToggles['live-service-web-worker'] - ? new CentrifugeServiceWorkerProxy(centrifugeServiceDeps) - : new CentrifugeService(centrifugeServiceDeps); + const centrifugeSrv = new CentrifugeService(centrifugeServiceDeps); setGrafanaLiveSrv( new GrafanaLiveService({ diff --git a/public/app/features/live/live.ts b/public/app/features/live/live.ts index 61c2224ffef..24492183714 100644 --- a/public/app/features/live/live.ts +++ b/public/app/features/live/live.ts @@ -1,8 +1,7 @@ -import { from, map, of, switchMap } from 'rxjs'; +import { map } from 'rxjs'; -import { DataFrame, toLiveChannelId, StreamingDataFrame } from '@grafana/data'; -import { BackendSrv, GrafanaLiveSrv, toDataQueryResponse } from '@grafana/runtime'; -import { standardStreamOptionsProvider, toStreamingDataResponse } from '@grafana/runtime/internal'; +import { toLiveChannelId, StreamingDataFrame } from '@grafana/data'; +import { BackendSrv, GrafanaLiveSrv } from '@grafana/runtime'; import { CentrifugeSrv, StreamingDataQueryResponse } from './centrifuge/service'; import { isStreamingResponseData, StreamingResponseDataType } from './data/utils'; @@ -65,26 +64,6 @@ export class GrafanaLiveService implements GrafanaLiveSrv { return this.deps.centrifugeSrv.getStream(address); }; - /** - * Execute a query over the live websocket and potentially subscribe to a live channel. - * - * Since the initial request and subscription are on the same socket, this will support HA setups - */ - getQueryData: GrafanaLiveSrv['getQueryData'] = (options) => { - return from(this.deps.centrifugeSrv.getQueryData(options)).pipe( - switchMap((rawResponse) => { - const parsedResponse = toDataQueryResponse(rawResponse, options.request.targets); - - const isSubscribable = - parsedResponse.data?.length && parsedResponse.data.find((f: DataFrame) => f.meta?.channel); - - return isSubscribable - ? toStreamingDataResponse(parsedResponse, options.request, standardStreamOptionsProvider) - : of(parsedResponse); - }) - ); - }; - /** * Publish into a channel * diff --git a/public/test/mocks/workers.ts b/public/test/mocks/workers.ts index 70ec661c12c..7880944953b 100644 --- a/public/test/mocks/workers.ts +++ b/public/test/mocks/workers.ts @@ -25,12 +25,3 @@ class LayoutMockWorker { jest.mock('../../app/plugins/panel/nodeGraph/createLayoutWorker', () => ({ createWorker: () => new LayoutMockWorker(), })); - -class BasicMockWorker { - postMessage() {} -} -const mockCreateWorker = { - createWorker: () => new BasicMockWorker(), -}; - -jest.mock('../../app/features/live/centrifuge/createCentrifugeServiceWorker', () => mockCreateWorker);