diff --git a/ui-ngx/src/app/core/api/alias-controller.ts b/ui-ngx/src/app/core/api/alias-controller.ts index 83dd4e102f..e22aeb5248 100644 --- a/ui-ngx/src/app/core/api/alias-controller.ts +++ b/ui-ngx/src/app/core/api/alias-controller.ts @@ -61,6 +61,9 @@ export class AliasController implements IAliasController { resolvedAliases: { [aliasId: string]: AliasInfo } = {}; resolvedAliasesObservable: { [aliasId: string]: Observable } = {}; + resolvedDevices: { [deviceId: string]: EntityInfo } = {}; + resolvedDevicesObservable: { [deviceId: string]: Observable } = {}; + resolvedAliasesToStateEntities: { [aliasId: string]: StateEntityInfo } = {}; constructor(private utils: UtilsService, @@ -261,9 +264,35 @@ export class AliasController implements IAliasController { } resolveSingleEntityInfoForDeviceId(deviceId: string): Observable { - const entityFilter = singleEntityFilterFromDeviceId(deviceId); - return this.entityService.findSingleEntityInfoByEntityFilter(entityFilter, - {ignoreLoading: true, ignoreErrors: true}); + let entityInfo = this.resolvedDevices[deviceId]; + if (entityInfo) { + return of(entityInfo); + } else if (this.resolvedDevicesObservable[deviceId]) { + return this.resolvedDevicesObservable[deviceId]; + } else { + const resolvedDeviceSubject = new ReplaySubject(); + this.resolvedDevicesObservable[deviceId] = resolvedDeviceSubject.asObservable(); + const entityFilter = singleEntityFilterFromDeviceId(deviceId); + this.entityService.findSingleEntityInfoByEntityFilter(entityFilter, + {ignoreLoading: true, ignoreErrors: true}).subscribe( + (resolvedEntityInfo) => { + this.resolvedDevices[deviceId] = resolvedEntityInfo; + delete this.resolvedDevicesObservable[deviceId]; + resolvedDeviceSubject.next(resolvedEntityInfo); + resolvedDeviceSubject.complete(); + }, + () => { + resolvedDeviceSubject.error(null); + delete this.resolvedDevicesObservable[deviceId]; + } + ); + entityInfo = this.resolvedDevices[deviceId]; + if (entityInfo) { + return of(entityInfo); + } else { + return this.resolvedDevicesObservable[deviceId]; + } + } } resolveSingleEntityInfoForTargetDevice(targetDevice: TargetDevice): Observable { diff --git a/ui-ngx/src/app/core/services/utils.service.ts b/ui-ngx/src/app/core/services/utils.service.ts index 4e78336b32..f2f2ba39e8 100644 --- a/ui-ngx/src/app/core/services/utils.service.ts +++ b/ui-ngx/src/app/core/services/utils.service.ts @@ -40,24 +40,15 @@ import { WindowMessage } from '@shared/models/window-message.model'; import { TranslateService } from '@ngx-translate/core'; import { customTranslationsPrefix, i18nPrefix } from '@app/shared/models/constants'; import { DataKey, Datasource, DatasourceType, KeyInfo } from '@shared/models/widget.models'; -import { DataKeyType } from '@app/shared/models/telemetry/telemetry.models'; -import { - alarmFields, - alarmSeverityTranslations, - alarmStatusTranslations -} from '@shared/models/alarm.models'; +import { DataKeyType, SharedTelemetrySubscriber } from '@app/shared/models/telemetry/telemetry.models'; +import { alarmFields, alarmSeverityTranslations, alarmStatusTranslations } from '@shared/models/alarm.models'; import { materialColors } from '@app/shared/models/material.models'; import { WidgetInfo } from '@home/models/widget-component.models'; import jsonSchemaDefaults from 'json-schema-defaults'; import { Observable } from 'rxjs'; import { publishReplay, refCount } from 'rxjs/operators'; import { WidgetContext } from '@app/modules/home/models/widget-component.models'; -import { - AttributeData, - LatestTelemetry, - TelemetrySubscriber, - TelemetryType -} from '@shared/models/telemetry/telemetry.models'; +import { AttributeData, LatestTelemetry, TelemetryType } from '@shared/models/telemetry/telemetry.models'; import { EntityId } from '@shared/models/id/entity-id'; import { DatePipe, DOCUMENT } from '@angular/common'; import { entityTypeTranslations } from '@shared/models/entity-type.models'; @@ -483,13 +474,13 @@ export class UtilsService { if (!entityId && ctx.datasources.length > 0) { entityId = this.getEntityIdFromDatasource(ctx.datasources[0]); } - const subscription = TelemetrySubscriber.createEntityAttributesSubscription(ctx.telemetryWsService, entityId, type, ctx.ngZone, keys); + const subscription = SharedTelemetrySubscriber.createEntityAttributesSubscription(ctx.telemetryWsService, entityId, type, ctx.ngZone, keys); if (!ctx.telemetrySubscribers) { ctx.telemetrySubscribers = []; } ctx.telemetrySubscribers.push(subscription); subscription.subscribe(); - return subscription.attributeData$().pipe( + return subscription.attributeData$.pipe( publishReplay(1), refCount() ); diff --git a/ui-ngx/src/app/modules/home/components/widget/lib/action/action-widget.models.ts b/ui-ngx/src/app/modules/home/components/widget/lib/action/action-widget.models.ts index b1eb79f2bb..20613b3d4f 100644 --- a/ui-ngx/src/app/modules/home/components/widget/lib/action/action-widget.models.ts +++ b/ui-ngx/src/app/modules/home/components/widget/lib/action/action-widget.models.ts @@ -17,8 +17,7 @@ import { AttributeData, AttributeScope, - LatestTelemetry, - TelemetrySubscriber, + LatestTelemetry, SharedTelemetrySubscriber, TelemetryType, telemetryTypeTranslationsShort } from '@shared/models/telemetry/telemetry.models'; @@ -436,7 +435,7 @@ export class ExecuteRpcValueGetter extends ValueGetter { export abstract class TelemetryValueGetter extends ValueGetter { protected targetEntityId: EntityId; - private telemetrySubscriber: TelemetrySubscriber; + private telemetrySubscriber: SharedTelemetrySubscriber; protected constructor(protected ctx: WidgetContext, protected settings: GetValueSettings, @@ -470,10 +469,10 @@ export abstract class TelemetryValueGetter private subscribeForTelemetryValue(): Observable { this.telemetrySubscriber = - TelemetrySubscriber.createEntityAttributesSubscription(this.ctx.telemetryWsService, this.targetEntityId, + SharedTelemetrySubscriber.createEntityAttributesSubscription(this.ctx.telemetryWsService, this.targetEntityId, this.scope(), this.ctx.ngZone, [this.getTelemetryValueSettings().key]); this.telemetrySubscriber.subscribe(); - return this.telemetrySubscriber.attributeData$().pipe( + return this.telemetrySubscriber.attributeData$.pipe( map((data) => { let value: V = null; const entry = data.find(attr => attr.key === this.getTelemetryValueSettings().key); diff --git a/ui-ngx/src/app/modules/home/models/datasource/attribute-datasource.ts b/ui-ngx/src/app/modules/home/models/datasource/attribute-datasource.ts index 00853d868e..30756a8677 100644 --- a/ui-ngx/src/app/modules/home/models/datasource/attribute-datasource.ts +++ b/ui-ngx/src/app/modules/home/models/datasource/attribute-datasource.ts @@ -25,7 +25,7 @@ import { AttributeData, AttributeScope, isClientSideTelemetryType, - TelemetrySubscriber, + SharedTelemetrySubscriber, TelemetryType } from '@shared/models/telemetry/telemetry.models'; import { AttributeService } from '@core/http/attribute.service'; @@ -42,7 +42,7 @@ export class AttributeDatasource implements DataSource { public selection = new SelectionModel(true, []); private allAttributes: Observable>; - private telemetrySubscriber: TelemetrySubscriber; + private telemetrySubscriber: SharedTelemetrySubscriber; constructor(private attributeService: AttributeService, private telemetryWsService: TelemetryWebsocketService, @@ -99,10 +99,10 @@ export class AttributeDatasource implements DataSource { if (!this.allAttributes) { let attributesObservable: Observable>; if (isClientSideTelemetryType.get(attributesScope)) { - this.telemetrySubscriber = TelemetrySubscriber.createEntityAttributesSubscription( + this.telemetrySubscriber = SharedTelemetrySubscriber.createEntityAttributesSubscription( this.telemetryWsService, entityId, attributesScope, this.zone); this.telemetrySubscriber.subscribe(); - attributesObservable = this.telemetrySubscriber.attributeData$(); + attributesObservable = this.telemetrySubscriber.attributeData$; } else { attributesObservable = this.attributeService.getEntityAttributes(entityId, attributesScope as AttributeScope); } diff --git a/ui-ngx/src/app/modules/home/models/widget-component.models.ts b/ui-ngx/src/app/modules/home/models/widget-component.models.ts index 9e739cafb4..cfc706bccc 100644 --- a/ui-ngx/src/app/modules/home/models/widget-component.models.ts +++ b/ui-ngx/src/app/modules/home/models/widget-component.models.ts @@ -98,7 +98,9 @@ import * as RxJSOperators from 'rxjs/operators'; import { TbPopoverComponent } from '@shared/components/popover.component'; import { EntityId } from '@shared/models/id/entity-id'; import { AlarmQuery, AlarmSearchStatus, AlarmStatus } from '@app/shared/models/alarm.models'; -import { ImagePipe, MillisecondsToTimeStringPipe, TelemetrySubscriber } from '@app/shared/public-api'; +import { ImagePipe } from '@shared/pipe/image.pipe'; +import { MillisecondsToTimeStringPipe } from '@shared/pipe/milliseconds-to-time-string.pipe'; +import { SharedTelemetrySubscriber, TelemetrySubscriber } from '@shared/models/telemetry/telemetry.models'; import { UserId } from '@shared/models/id/user-id'; import { UserSettingsService } from '@core/http/user-settings.service'; import { DataKeySettingsFunction } from '@home/components/widget/config/data-keys.component.models'; @@ -204,7 +206,7 @@ export class WidgetContext { userSettingsService: UserSettingsService; utilsService: UtilsService; telemetryWsService: TelemetryWebsocketService; - telemetrySubscribers?: TelemetrySubscriber[]; + telemetrySubscribers?: Array; date: DatePipe; imagePipe: ImagePipe; milliSecondsToTimeString: MillisecondsToTimeStringPipe; diff --git a/ui-ngx/src/app/shared/models/telemetry/telemetry.models.ts b/ui-ngx/src/app/shared/models/telemetry/telemetry.models.ts index 9f0d244e9e..2db1983b0b 100644 --- a/ui-ngx/src/app/shared/models/telemetry/telemetry.models.ts +++ b/ui-ngx/src/app/shared/models/telemetry/telemetry.models.ts @@ -17,7 +17,7 @@ import { EntityType } from '@shared/models/entity-type.models'; import { AggregationType } from '../time/time.models'; -import { BehaviorSubject, Observable, ReplaySubject } from 'rxjs'; +import { BehaviorSubject, connectable, Observable, ReplaySubject, Subscription } from 'rxjs'; import { EntityId } from '@shared/models/id/entity-id'; import { map } from 'rxjs/operators'; import { NgZone } from '@angular/core'; @@ -731,6 +731,91 @@ export class NotificationsUpdate extends CmdUpdate { } } +interface SharedSubscriptionInfo { + key: string; + subscriber: TelemetrySubscriber; + subscribed: boolean; + sharedSubscribers: Set; +} + +export class SharedTelemetrySubscriber { + + private static subscribersCache: {[key: string]: SharedSubscriptionInfo} = {}; + + private static createTelemetrySubscriberKey (entityId: EntityId, attributeScope: TelemetryType, keys: string[] = null): string { + let key = entityId.entityType + '_' + entityId.id + '_' + attributeScope; + if (keys) { + key += '_' + keys.sort().join('_'); + } + return key; + } + + private subscribed = false; + + private attributeDataSubject = connectable(this.sharedSubscriptionInfo.subscriber.attributeData$(), + { connector: () => new ReplaySubject>(1)}); + + private subscriptions = new Array(); + + public attributeData$: Observable> = this.attributeDataSubject; //this.attributeDataSubject.asObservable(); + + public static createEntityAttributesSubscription(telemetryService: TelemetryWebsocketService, + entityId: EntityId, attributeScope: TelemetryType, + zone: NgZone, keys: string[] = null): SharedTelemetrySubscriber { + const key = SharedTelemetrySubscriber.createTelemetrySubscriberKey(entityId, attributeScope, keys); + let info = SharedTelemetrySubscriber.subscribersCache[key]; + if (!info) { + const subscriber = TelemetrySubscriber.createEntityAttributesSubscription( + telemetryService, entityId, attributeScope, zone, keys + ); + info = { + key, + subscriber, + subscribed: false, + sharedSubscribers: new Set() + }; + SharedTelemetrySubscriber.subscribersCache[key] = info; + } + const sharedSubscriber = new SharedTelemetrySubscriber(info); + info.sharedSubscribers.add(sharedSubscriber); + return sharedSubscriber; + } + + private constructor(private sharedSubscriptionInfo: SharedSubscriptionInfo) { + } + + public subscribe() { + if (!this.subscribed) { + this.subscribed = true; + this.subscriptions.push(this.attributeDataSubject.connect()); + if (!this.sharedSubscriptionInfo.subscribed) { + this.sharedSubscriptionInfo.subscriber.subscribe(); + this.sharedSubscriptionInfo.subscribed = true; + } + } + } + + public unsubscribe() { + if (this.subscribed) { + this.complete(); + } + this.sharedSubscriptionInfo.sharedSubscribers.delete(this); + if (!this.sharedSubscriptionInfo.sharedSubscribers.size) { + if (this.sharedSubscriptionInfo.subscribed) { + this.sharedSubscriptionInfo.subscriber.unsubscribe(); + this.sharedSubscriptionInfo.subscribed = false; + } + delete SharedTelemetrySubscriber.subscribersCache[this.sharedSubscriptionInfo.key]; + } + } + + private complete() { + this.subscriptions.forEach(subscription => subscription.unsubscribe()); + this.subscriptions.length = 0; + } + +} + export class TelemetrySubscriber extends WsSubscriber { private dataSubject = new ReplaySubject(1);