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();
     }
 }

Reply via email to