Browse Source

UI: Add cache for resolved deivce info to alias controller. Introduce shared telemetry subscriber to reuse telemetry subscriptions with same entityId and keys.

pull/12017/head
Igor Kulikov 2 years ago
parent
commit
344accbdb5
  1. 35
      ui-ngx/src/app/core/api/alias-controller.ts
  2. 19
      ui-ngx/src/app/core/services/utils.service.ts
  3. 9
      ui-ngx/src/app/modules/home/components/widget/lib/action/action-widget.models.ts
  4. 8
      ui-ngx/src/app/modules/home/models/datasource/attribute-datasource.ts
  5. 6
      ui-ngx/src/app/modules/home/models/widget-component.models.ts
  6. 87
      ui-ngx/src/app/shared/models/telemetry/telemetry.models.ts

35
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<AliasInfo> } = {};
resolvedDevices: { [deviceId: string]: EntityInfo } = {};
resolvedDevicesObservable: { [deviceId: string]: Observable<EntityInfo> } = {};
resolvedAliasesToStateEntities: { [aliasId: string]: StateEntityInfo } = {};
constructor(private utils: UtilsService,
@ -261,9 +264,35 @@ export class AliasController implements IAliasController {
}
resolveSingleEntityInfoForDeviceId(deviceId: string): Observable<EntityInfo> {
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<EntityInfo>();
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<EntityInfo> {

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

9
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<V> extends ValueGetter<V> {
export abstract class TelemetryValueGetter<V, S extends TelemetryValueSettings> extends ValueGetter<V> {
protected targetEntityId: EntityId;
private telemetrySubscriber: TelemetrySubscriber;
private telemetrySubscriber: SharedTelemetrySubscriber;
protected constructor(protected ctx: WidgetContext,
protected settings: GetValueSettings<V>,
@ -470,10 +469,10 @@ export abstract class TelemetryValueGetter<V, S extends TelemetryValueSettings>
private subscribeForTelemetryValue(): Observable<V> {
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);

8
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<AttributeData> {
public selection = new SelectionModel<AttributeData>(true, []);
private allAttributes: Observable<Array<AttributeData>>;
private telemetrySubscriber: TelemetrySubscriber;
private telemetrySubscriber: SharedTelemetrySubscriber;
constructor(private attributeService: AttributeService,
private telemetryWsService: TelemetryWebsocketService,
@ -99,10 +99,10 @@ export class AttributeDatasource implements DataSource<AttributeData> {
if (!this.allAttributes) {
let attributesObservable: Observable<Array<AttributeData>>;
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);
}

6
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<TelemetrySubscriber | SharedTelemetrySubscriber>;
date: DatePipe;
imagePipe: ImagePipe;
milliSecondsToTimeString: MillisecondsToTimeStringPipe;

87
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<SharedTelemetrySubscriber>;
}
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<Array<AttributeData>>(1)});
private subscriptions = new Array<Subscription>();
public attributeData$: Observable<Array<AttributeData>> = 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>()
};
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<SubscriptionUpdate>(1);

Loading…
Cancel
Save