This is an automated email from the ASF dual-hosted git repository.
dominikriemer pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/streampipes.git
The following commit(s) were added to refs/heads/dev by this push:
new 644fd68559 fix: Properly cleanup subscriptions in kiosk mode (#4622)
644fd68559 is described below
commit 644fd6855922a7855784f4cf2bbd4e9fe05eebd6
Author: Dominik Riemer <[email protected]>
AuthorDate: Tue Jun 23 19:34:12 2026 +0200
fix: Properly cleanup subscriptions in kiosk mode (#4622)
---
.../base/base-data-explorer-widget.directive.ts | 16 +++--
.../components/kiosk/dashboard-kiosk.component.ts | 78 +++++++++++++---------
2 files changed, 58 insertions(+), 36 deletions(-)
diff --git
a/ui/src/app/chart-shared/components/charts/base/base-data-explorer-widget.directive.ts
b/ui/src/app/chart-shared/components/charts/base/base-data-explorer-widget.directive.ts
index 4b8d9f3f4a..5cf0f3e886 100644
---
a/ui/src/app/chart-shared/components/charts/base/base-data-explorer-widget.directive.ts
+++
b/ui/src/app/chart-shared/components/charts/base/base-data-explorer-widget.directive.ts
@@ -42,9 +42,9 @@ import {
FieldProvider,
ObservableGenerator,
} from '../../../models/dataview-dashboard.model';
-import { Observable, Subject, Subscription, zip } from 'rxjs';
+import { EMPTY, Observable, Subject, Subscription, zip } from 'rxjs';
import { ChartFieldProviderService } from
'../../../services/chart-field-provider.service';
-import { catchError, switchMap } from 'rxjs/operators';
+import { catchError, exhaustMap, finalize } from 'rxjs/operators';
import { ChartRegistry } from '../../../registry/chart-registry.service';
import { SpFieldUpdateService } from '../../../services/field-update.service';
import {
@@ -123,6 +123,7 @@ export abstract class BaseDataExplorerWidgetDirective<
requestQueue$: Subject<Observable<SpQueryResult>[]> = new Subject<
Observable<SpQueryResult>[]
>();
+ requestInProgress = false;
protected widgetConfigurationService = inject(ChartConfigurationService);
protected resizeService = inject(ResizeService);
@@ -142,14 +143,18 @@ export abstract class BaseDataExplorerWidgetDirective<
this.requestQueue$
.pipe(
- switchMap(observables => {
+ exhaustMap(observables => {
+ this.requestInProgress = true;
this.errorCallback.emit(undefined);
return zip(...observables).pipe(
catchError(err => {
this.timerCallback.emit(false);
this.errorCallback.emit(err.error);
this.dataReceivedCallback.emit([]);
- return [];
+ return EMPTY;
+ }),
+ finalize(() => {
+ this.requestInProgress = false;
}),
);
}),
@@ -254,6 +259,9 @@ export abstract class BaseDataExplorerWidgetDirective<
}
public updateData(includeTooMuchEventsParameter: boolean = true) {
+ if (this.requestInProgress) {
+ return;
+ }
this.beforeDataFetched();
this.loadData(includeTooMuchEventsParameter);
}
diff --git
a/ui/src/app/dashboard-kiosk/components/kiosk/dashboard-kiosk.component.ts
b/ui/src/app/dashboard-kiosk/components/kiosk/dashboard-kiosk.component.ts
index 45f4defe32..3e7d1f9540 100644
--- a/ui/src/app/dashboard-kiosk/components/kiosk/dashboard-kiosk.component.ts
+++ b/ui/src/app/dashboard-kiosk/components/kiosk/dashboard-kiosk.component.ts
@@ -16,7 +16,14 @@
*
*/
-import { Component, inject, OnDestroy, OnInit } from '@angular/core';
+import {
+ Component,
+ DestroyRef,
+ inject,
+ OnDestroy,
+ OnInit,
+} from '@angular/core';
+import { takeUntilDestroyed } from '@angular/core/rxjs-interop';
import {
CompositeDashboard,
Dashboard,
@@ -25,8 +32,8 @@ import {
TimeSettings,
} from '@streampipes/platform-services';
import { ActivatedRoute } from '@angular/router';
-import { of, Subscription, timer } from 'rxjs';
-import { switchMap } from 'rxjs/operators';
+import { EMPTY, Subscription, timer } from 'rxjs';
+import { catchError, exhaustMap, tap } from 'rxjs/operators';
import { TimeSelectionService } from '@streampipes/shared-ui';
import { DataExplorerDashboardService } from
'../../../dashboard-shared/services/dashboard.service';
import { ChartSharedService } from
'../../../chart-shared/services/chart-shared.service';
@@ -53,6 +60,7 @@ import {
})
export class DashboardKioskComponent implements OnInit, OnDestroy {
private route = inject(ActivatedRoute);
+ private destroyRef = inject(DestroyRef);
private dashboardService = inject(DashboardService);
private timeSelectionService = inject(TimeSelectionService);
private dataExplorerDashboardService =
inject(DataExplorerDashboardService);
@@ -62,6 +70,7 @@ export class DashboardKioskComponent implements OnInit,
OnDestroy {
dashboard: Dashboard;
widgets: DataExplorerWidgetModel[] = [];
refresh$: Subscription;
+ dashboardRefresh$: Subscription;
eTag: string;
ngOnInit() {
@@ -72,26 +81,26 @@ export class DashboardKioskComponent implements OnInit,
OnDestroy {
);
this.dashboardService
.getCompositeDashboard(dashboardId)
+ .pipe(takeUntilDestroyed(this.destroyRef))
.subscribe(res => {
if (res.ok) {
- const cd = res.body;
- cd.dashboard.widgets.forEach(w => {
- w.id ??=
-
this.dataExplorerDashboardService.makeUniqueWidgetId();
- });
const eTag = res.headers.get('ETag');
- this.initDashboard(cd, eTag);
+ this.initDashboard(res.body, eTag);
}
});
}
initDashboard(cd: CompositeDashboard, eTag: string): void {
+ cd.dashboard.widgets.forEach(w => {
+ w.id ??= this.dataExplorerDashboardService.makeUniqueWidgetId();
+ });
this.dashboard = cd.dashboard;
this.widgets = cd.widgets;
this.eTag = eTag;
+ this.refresh$?.unsubscribe();
if (this.dashboard.dashboardLiveSettings.refreshModeActive) {
this.createQuerySubscription();
- this.createRefreshListener();
+ this.createDashboardRefreshSubscription();
}
}
@@ -102,40 +111,44 @@ export class DashboardKioskComponent implements OnInit,
OnDestroy {
1000,
)
.pipe(
- switchMap(() => {
+ tap(() => {
this.timeSelectionService.updateTimeSettings(
this.timeSelectionService.defaultQuickTimeSelections,
this.dashboard.dashboardTimeSettings,
new Date(),
);
this.updateDateRange(this.dashboard.dashboardTimeSettings);
- return of(null);
}),
+ takeUntilDestroyed(this.destroyRef),
)
.subscribe();
}
- createRefreshListener(): void {
- this.dashboardService
- .getCompositeDashboard(this.dashboard.elementId, this.eTag) //
this should send If-None-Match
- .subscribe({
- next: res => {
- if (res.status === 200) {
- const newEtag = res.headers.get('ETag');
- if (newEtag) {
- this.eTag = newEtag;
- }
- this.dashboard = undefined;
- this.refresh$?.unsubscribe();
- setTimeout(() => {
- this.initDashboard(res.body, newEtag);
- });
+ createDashboardRefreshSubscription(): void {
+ if (this.dashboardRefresh$) {
+ return;
+ }
+
+ this.dashboardRefresh$ = timer(5000, 5000)
+ .pipe(
+ exhaustMap(() =>
+ this.dashboardService
+ .getCompositeDashboard(
+ this.dashboard.elementId,
+ this.eTag,
+ )
+ .pipe(catchError(() => EMPTY)),
+ ),
+ takeUntilDestroyed(this.destroyRef),
+ )
+ .subscribe(res => {
+ if (res.status === 200) {
+ const newEtag = res.headers.get('ETag');
+ if (newEtag) {
+ this.eTag = newEtag;
}
- setTimeout(() => this.createRefreshListener(), 5000);
- },
- error: _err => {
- setTimeout(() => this.createRefreshListener(), 5000);
- },
+ this.initDashboard(res.body, newEtag);
+ }
});
}
@@ -150,5 +163,6 @@ export class DashboardKioskComponent implements OnInit,
OnDestroy {
ngOnDestroy() {
this.refresh$?.unsubscribe();
+ this.dashboardRefresh$?.unsubscribe();
}
}