diff --git a/ui-ngx/src/app/core/api/alias-controller.ts b/ui-ngx/src/app/core/api/alias-controller.ts index 095335839c..e0afe96c42 100644 --- a/ui-ngx/src/app/core/api/alias-controller.ts +++ b/ui-ngx/src/app/core/api/alias-controller.ts @@ -17,12 +17,13 @@ import { AliasInfo, IAliasController, StateControllerHolder, StateEntityInfo } from '@core/api/widget-api.models'; import { forkJoin, Observable, of, ReplaySubject, Subject } from 'rxjs'; import { DataKey, Datasource, DatasourceType } from '@app/shared/models/widget.models'; -import { deepClone, isEqual, createLabelFromDatasource } from '@core/utils'; +import { createLabelFromDatasource, deepClone, isEqual } from '@core/utils'; import { EntityService } from '@core/http/entity.service'; import { UtilsService } from '@core/services/utils.service'; -import { EntityAliases } from '@shared/models/alias.models'; +import { AliasFilterType, EntityAliases } from '@shared/models/alias.models'; import { EntityInfo } from '@shared/models/entity.models'; import { map } from 'rxjs/operators'; +import { defaultEntityDataPageLink } from '@shared/models/query/query.models'; export class AliasController implements IAliasController { @@ -169,6 +170,92 @@ export class AliasController implements IAliasController { } private resolveDatasource(datasource: Datasource, isSingle?: boolean): Observable> { + if (datasource.type === DatasourceType.entity) { + if (datasource.entityAliasId) { + return this.getAliasInfo(datasource.entityAliasId).pipe( + map((aliasInfo) => { + datasource.aliasName = aliasInfo.alias; + if (aliasInfo.resolveMultiple && !isSingle) { + let newDatasource: Datasource; + // const resolvedEntities = aliasInfo.resolvedEntities; + if (aliasInfo.entityFilter) { + newDatasource = deepClone(datasource); + newDatasource.entityFilter = aliasInfo.entityFilter; + /*const datasources: Array = []; + for (let i = 0; i < resolvedEntities.length; i++) { + const resolvedEntity = resolvedEntities[i]; + newDatasource = deepClone(datasource); + if (resolvedEntity.origEntity) { + newDatasource.entity = deepClone(resolvedEntity.origEntity); + } else { + newDatasource.entity = {}; + } + newDatasource.entityId = resolvedEntity.id; + newDatasource.entityType = resolvedEntity.entityType; + newDatasource.entityName = resolvedEntity.name; + newDatasource.entityLabel = resolvedEntity.label; + newDatasource.entityDescription = resolvedEntity.entityDescription; + newDatasource.name = resolvedEntity.name; + newDatasource.generated = i > 0 ? true : false; + datasources.push(newDatasource); + } + return datasources;*/ + return [newDatasource]; + } else { + if (aliasInfo.stateEntity) { + newDatasource = deepClone(datasource); + newDatasource.unresolvedStateEntity = true; + return [newDatasource]; + } else { + return []; + // throw new Error('Unable to resolve datasource.'); + } + } + } else { + const entity = aliasInfo.currentEntity; + if (entity) { + if (entity.origEntity) { + datasource.entity = deepClone(entity.origEntity); + } else { + datasource.entity = {}; + } + datasource.entityId = entity.id; + datasource.entityType = entity.entityType; + datasource.entityName = entity.name; + datasource.entityLabel = entity.label; + datasource.name = entity.name; + datasource.entityDescription = entity.entityDescription; + datasource.entityFilter = { + type: AliasFilterType.singleEntity, + singleEntity: { + id: entity.id, + entityType: entity.entityType + } + }; + return [datasource]; + } else { + if (aliasInfo.stateEntity) { + datasource.unresolvedStateEntity = true; + return [datasource]; + } else { + return []; + // throw new Error('Unable to resolve datasource.'); + } + } + } + }) + ); + } else { + datasource.aliasName = datasource.entityName; + datasource.name = datasource.entityName; + return of([datasource]); + } + } else { + return of([datasource]); + } + } + + /* private resolveDatasourceOld(datasource: Datasource, isSingle?: boolean): Observable> { if (datasource.type === DatasourceType.entity) { if (datasource.entityAliasId) { return this.getAliasInfo(datasource.entityAliasId).pipe( @@ -242,7 +329,7 @@ export class AliasController implements IAliasController { } else { return of([datasource]); } - } + } */ resolveAlarmSource(alarmSource: Datasource): Observable { return this.resolveDatasource(alarmSource, true).pipe( @@ -268,6 +355,48 @@ export class AliasController implements IAliasController { } resolveDatasources(datasources: Array): Observable> { + const newDatasources = deepClone(datasources); + const observables = new Array>>(); + newDatasources.forEach((datasource) => { + observables.push(this.resolveDatasource(datasource)); + }); + return forkJoin(observables).pipe( + map((arrayOfDatasources) => { + const result = new Array(); + arrayOfDatasources.forEach((datasourcesArray) => { + result.push(...datasourcesArray); + }); + let functionIndex = 0; + result.forEach((datasource) => { + if (datasource.type === DatasourceType.function) { + let name: string; + if (datasource.name && datasource.name.length) { + name = datasource.name; + } else { + functionIndex++; + name = DatasourceType.function; + if (functionIndex > 1) { + name += ' ' + functionIndex; + } + } + datasource.name = name; + datasource.aliasName = name; + datasource.entityName = name; + } else if (datasource.unresolvedStateEntity) { + datasource.name = 'Unresolved'; + datasource.entityName = 'Unresolved'; + } else if (datasource.type === DatasourceType.entity) { + if (!datasource.pageLink) { + datasource.pageLink = deepClone(defaultEntityDataPageLink); + } + } + }); + return result; + }) + ); + } + + /*resolveDatasourcesOld(datasources: Array): Observable> { const newDatasources = deepClone(datasources); const observables = new Array>>(); newDatasources.forEach((datasource) => { @@ -317,9 +446,9 @@ export class AliasController implements IAliasController { return result; }) ); - } + }*/ - private updateDatasourceKeyLabels(datasource: Datasource) { + /* private updateDatasourceKeyLabels(datasource: Datasource) { datasource.dataKeys.forEach((dataKey) => { this.updateDataKeyLabel(dataKey, datasource); }); @@ -330,7 +459,7 @@ export class AliasController implements IAliasController { dataKey.pattern = deepClone(dataKey.label); } dataKey.label = createLabelFromDatasource(datasource, dataKey.pattern); - } + }*/ getInstantAliasInfo(aliasId: string): AliasInfo { return this.resolvedAliases[aliasId]; diff --git a/ui-ngx/src/app/core/api/entity-data-subscription.ts b/ui-ngx/src/app/core/api/entity-data-subscription.ts new file mode 100644 index 0000000000..a0135fe434 --- /dev/null +++ b/ui-ngx/src/app/core/api/entity-data-subscription.ts @@ -0,0 +1,670 @@ +/// +/// Copyright © 2016-2020 The Thingsboard Authors +/// +/// Licensed under the Apache License, Version 2.0 (the "License"); +/// you may not use this file except in compliance with the License. +/// You may obtain a copy of the License at +/// +/// http://www.apache.org/licenses/LICENSE-2.0 +/// +/// Unless required by applicable law or agreed to in writing, software +/// distributed under the License is distributed on an "AS IS" BASIS, +/// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +/// See the License for the specific language governing permissions and +/// limitations under the License. +/// + +import { DataSet, DataSetHolder, DatasourceType, widgetType } from '@shared/models/widget.models'; +import { AggregationType, SubscriptionTimewindow, YEAR } from '@shared/models/time/time.models'; +import { SubscriptionDataKey } from '@core/api/datasource-subcription'; +import { + EntityData, + EntityDataPageLink, + EntityFilter, + EntityKey, + EntityKeyType, + KeyFilter, + TsValue +} from '@shared/models/query/query.models'; +import { + DataKeyType, + EntityDataCmd, + SubscriptionData, + SubscriptionDataHolder, + TelemetryService, + TelemetrySubscriber +} from '@shared/models/telemetry/telemetry.models'; +import { UtilsService } from '@core/services/utils.service'; +import { EntityDataListener } from '@core/api/entity-data.service'; +import { deepClone, isDefinedAndNotNull, isObject, objectHashCode } from '@core/utils'; +import { PageData } from '@shared/models/page/page-data'; +import { DataAggregator } from '@core/api/data-aggregator'; +import { NULL_UUID } from '@shared/models/id/has-uuid'; +import { EntityType } from '@shared/models/entity-type.models'; +import Timeout = NodeJS.Timeout; + +export interface EntityDataSubscriptionOptions { + datasourceType: DatasourceType; + dataKeys: Array; + type: widgetType; + entityFilter?: EntityFilter; + pageLink?: EntityDataPageLink; + keyFilters?: Array; + subscriptionTimewindow?: SubscriptionTimewindow; +} + +declare type DataKeyFunction = (time: number, prevValue: any) => any; +declare type DataKeyPostFunction = (time: number, value: any, prevValue: any, timePrev: number, prevOrigValue: any) => any; +declare type DataUpdatedCb = (data: DataSetHolder, dataIndex: number, dataKeyIndex: number, detectChanges: boolean) => void; + +export class EntityDataSubscription { + + private listeners: Array = []; + private datasourceType: DatasourceType = this.entityDataSubscriptionOptions.datasourceType; + + private history = this.entityDataSubscriptionOptions.subscriptionTimewindow && + isObject(this.entityDataSubscriptionOptions.subscriptionTimewindow.fixedWindow); + + private realtime = this.entityDataSubscriptionOptions.subscriptionTimewindow && + isDefinedAndNotNull(this.entityDataSubscriptionOptions.subscriptionTimewindow.realtimeWindowMs); + + private subscriber: TelemetrySubscriber; + + private attrFields: Array; + private tsFields: Array; + private latestValues: Array; + + private pageData: PageData; + private subsTw: SubscriptionTimewindow; + private dataAggregators: Array; + private dataKeys: {[key: string]: Array | SubscriptionDataKey} = {} + private datasourceData: {[index: number]: {[key: string]: DataSetHolder}}; + private datasourceOrigData: {[index: number]: {[key: string]: DataSetHolder}}; + private entityIdToDataIndex: {[id: string]: number}; + + private frequency: number; + private tickScheduledTime = 0; + private tickElapsed = 0; + private timer: Timeout; + + constructor(private entityDataSubscriptionOptions: EntityDataSubscriptionOptions, + private telemetryService: TelemetryService, + private utils: UtilsService) { + this.initializeSubscription(); + } + + private initializeSubscription() { + for (let i = 0; i < this.entityDataSubscriptionOptions.dataKeys.length; i++) { + const dataKey = deepClone(this.entityDataSubscriptionOptions.dataKeys[i]); + dataKey.index = i; + if (this.datasourceType === DatasourceType.function) { + if (!dataKey.func) { + dataKey.func = new Function('time', 'prevValue', dataKey.funcBody) as DataKeyFunction; + } + } else { + if (dataKey.postFuncBody && !dataKey.postFunc) { + dataKey.postFunc = new Function('time', 'value', 'prevValue', 'timePrev', 'prevOrigValue', + dataKey.postFuncBody) as DataKeyPostFunction; + } + } + let key: string; + if (this.datasourceType === DatasourceType.entity || this.entityDataSubscriptionOptions.type === widgetType.timeseries) { + if (this.datasourceType === DatasourceType.function) { + key = `${dataKey.name}_${dataKey.index}_${dataKey.type}`; + } else { + key = `${dataKey.name}_${dataKey.type}`; + } + let dataKeysList = this.dataKeys[key] as Array; + if (!dataKeysList) { + dataKeysList = []; + this.dataKeys[key] = dataKeysList; + } + dataKeysList.push(dataKey); + } else { + key = String(objectHashCode(dataKey)); + this.dataKeys[key] = dataKey; + } + dataKey.key = key; + } + if (this.datasourceType === DatasourceType.function) { + this.frequency = 1000; + if (this.entityDataSubscriptionOptions.type === widgetType.timeseries) { + this.frequency = Math.min(this.entityDataSubscriptionOptions.subscriptionTimewindow.aggregation.interval, 5000); + } + } + } + + public addListener(listener: EntityDataListener) { + this.listeners.push(listener); + if (this.history) { + this.start(); + } + } + + public hasListeners(): boolean { + return this.listeners.length > 0; + } + + public removeListener(listener: EntityDataListener) { + this.listeners.splice(this.listeners.indexOf(listener), 1); + } + + public syncListener(listener: EntityDataListener) { + if (this.pageData) { + let key: string; + let dataKey: SubscriptionDataKey; + const data: Array> = []; + for (let dataIndex = 0; dataIndex < this.pageData.data.length; dataIndex++) { + data[dataIndex] = []; + for (key of Object.keys(this.dataKeys)) { + if (this.datasourceType === DatasourceType.entity || this.entityDataSubscriptionOptions.type === widgetType.timeseries) { + const dataKeysList = this.dataKeys[key] as Array; + for (let i = 0; i < dataKeysList.length; i++) { + dataKey = dataKeysList[i]; + const datasourceKey = `${key}_${i}`; + data[dataIndex][dataKey.index] = this.datasourceData[dataIndex][datasourceKey]; + } + } else { + dataKey = this.dataKeys[key] as SubscriptionDataKey; + data[dataIndex][dataKey.index] = this.datasourceData[dataIndex][key]; + } + } + } + listener.dataLoaded(this.pageData, data, listener.configDatasourceIndex); + } + } + + public unsubscribe() { + if (this.timer) { + clearTimeout(this.timer); + this.timer = null; + } + if (this.datasourceType === DatasourceType.entity) { + if (this.subscriber) { + this.subscriber.unsubscribe(); + this.subscriber = null; + } + } + if (this.dataAggregators) { + this.dataAggregators.forEach((aggregator) => { + aggregator.destroy(); + }) + this.dataAggregators = null; + } + this.pageData = null; + } + + public start() { + if (this.history && !this.hasListeners()) { + return; + } + this.subsTw = this.entityDataSubscriptionOptions.subscriptionTimewindow; + if (this.datasourceType === DatasourceType.entity) { + const entityFields: Array = + this.entityDataSubscriptionOptions.dataKeys.filter(dataKey => dataKey.type === DataKeyType.entityField).map( + dataKey => ({ type: EntityKeyType.ENTITY_FIELD, key: dataKey.name }) + ); + if (!entityFields.find(key => key.key === 'name')) { + entityFields.push({ + type: EntityKeyType.ENTITY_FIELD, + key: 'name' + }); + } + + this.attrFields = this.entityDataSubscriptionOptions.dataKeys.filter(dataKey => dataKey.type === DataKeyType.attribute).map( + dataKey => ({ type: EntityKeyType.ATTRIBUTE, key: dataKey.name }) + ); + + this.tsFields = this.entityDataSubscriptionOptions.dataKeys.filter(dataKey => dataKey.type === DataKeyType.timeseries).map( + dataKey => ({ type: EntityKeyType.TIME_SERIES, key: dataKey.name }) + ); + + this.latestValues = this.attrFields.concat(this.tsFields); + + this.subscriber = new TelemetrySubscriber(this.telemetryService); + const command = new EntityDataCmd(); + + command.query = { + entityFilter: this.entityDataSubscriptionOptions.entityFilter, + pageLink: this.entityDataSubscriptionOptions.pageLink, + keyFilters: this.entityDataSubscriptionOptions.keyFilters, + entityFields, + latestValues: this.latestValues + }; + + if (this.entityDataSubscriptionOptions.type === widgetType.timeseries) { + if (this.tsFields.length > 0) { + if (this.history) { + command.historyCmd = { + keys: this.tsFields.map(key => key.key), + startTs: this.subsTw.fixedWindow.startTimeMs, + endTs: this.subsTw.fixedWindow.endTimeMs, + interval: this.subsTw.aggregation.interval, + limit: this.subsTw.aggregation.limit, + agg: this.subsTw.aggregation.type + }; + if (this.subsTw.aggregation.stateData) { + command.historyCmd.startTs -= YEAR; + } + } else { + command.tsCmd = { + keys: this.tsFields.map(key => key.key), + startTs: this.subsTw.startTs, + timeWindow: this.subsTw.aggregation.timeWindow, + interval: this.subsTw.aggregation.interval, + limit: this.subsTw.aggregation.limit, + agg: this.subsTw.aggregation.type + } + if (this.subsTw.aggregation.stateData) { + command.historyCmd = { + keys: this.tsFields.map(key => key.key), + startTs: this.subsTw.startTs - YEAR, + endTs: this.subsTw.startTs, + interval: this.subsTw.aggregation.interval, + limit: this.subsTw.aggregation.limit, + agg: this.subsTw.aggregation.type + }; + } + this.subscriber.reconnect$.subscribe(() => { + let newSubsTw: SubscriptionTimewindow = null; + this.listeners.forEach((listener) => { + if (!newSubsTw) { + newSubsTw = listener.updateRealtimeSubscription(); + } else { + listener.setRealtimeSubscription(newSubsTw); + } + }); + this.subsTw = newSubsTw; + command.tsCmd.startTs = this.subsTw.startTs; + command.tsCmd.timeWindow = this.subsTw.aggregation.timeWindow; + command.tsCmd.interval = this.subsTw.aggregation.interval; + command.tsCmd.limit = this.subsTw.aggregation.limit; + command.tsCmd.agg = this.subsTw.aggregation.type; + if (this.subsTw.aggregation.stateData) { + command.historyCmd.startTs = this.subsTw.startTs - YEAR; + command.historyCmd.endTs = this.subsTw.startTs; + command.historyCmd.interval = this.subsTw.aggregation.interval; + command.historyCmd.limit = this.subsTw.aggregation.limit; + command.historyCmd.agg = this.subsTw.aggregation.type; + } + }); + } + } + } else if (this.entityDataSubscriptionOptions.type === widgetType.latest) { + if (this.latestValues.length > 0) { + command.latestCmd = { + keys: this.latestValues.map(key => key.key) + }; + } + } + this.subscriber.subscriptionCommands.push(command); + + this.subscriber.entityData$.subscribe( + (entityDataUpdate) => { + if (entityDataUpdate.data) { + this.onPageData(entityDataUpdate.data); + } else if (entityDataUpdate.update) { + this.onDataUpdate(entityDataUpdate.update); + } + } + ); + + this.subscriber.subscribe(); + } else if (this.datasourceType === DatasourceType.function) { + const entityData: EntityData = { + entityId: { + id: NULL_UUID, + entityType: EntityType.DEVICE + }, + timeseries: {}, + latest: {} + }; + const name = DatasourceType.function; + entityData.latest[EntityKeyType.ENTITY_FIELD] = { + name: {ts: Date.now(), value: name} + }; + const pageData: PageData = { + data: [entityData], + hasNext: false, + totalElements: 1, + totalPages: 1 + }; + this.onPageData(pageData); + this.tickScheduledTime = this.utils.currentPerfTime(); + if (this.history) { + this.onTick(true); + } else { + this.timer = setTimeout(this.onTick.bind(this, true), 0); + } + } + } + + private onPageData(pageData: PageData) { + if (this.dataAggregators) { + this.dataAggregators.forEach((aggregator) => { + aggregator.destroy(); + }) + this.dataAggregators = null; + } + this.datasourceData = []; + this.dataAggregators = []; + this.entityIdToDataIndex = {}; + let tsKeyNames; + if (this.datasourceType === DatasourceType.function) { + tsKeyNames = []; + for (const key of Object.keys(this.dataKeys)) { + const dataKeysList = this.dataKeys[key] as Array; + dataKeysList.forEach((subscriptionDataKey) => { + tsKeyNames.push(`${subscriptionDataKey.name}_${subscriptionDataKey.index}`); + }); + } + } else { + tsKeyNames = this.tsFields.map(field => field.key); + } + for (let dataIndex = 0; dataIndex < pageData.data.length; dataIndex++) { + const entityData = pageData.data[dataIndex]; + this.entityIdToDataIndex[entityData.entityId.id] = dataIndex; + this.datasourceData[dataIndex] = {}; + if (this.entityDataSubscriptionOptions.type === widgetType.timeseries) { + if (this.datasourceType === DatasourceType.function) { + this.dataAggregators[dataIndex] = this.createRealtimeDataAggregator(this.subsTw, tsKeyNames, + DataKeyType.function, dataIndex, this.notifyListeners.bind(this)); + } else if (!this.history && tsKeyNames.length) { + this.dataAggregators[dataIndex] = this.createRealtimeDataAggregator(this.subsTw, tsKeyNames, + DataKeyType.timeseries, dataIndex, this.notifyListeners.bind(this)); + } + } + for (const key of Object.keys(this.dataKeys)) { + const dataKey = this.dataKeys[key]; + if (this.datasourceType === DatasourceType.entity || this.entityDataSubscriptionOptions.type === widgetType.timeseries) { + const dataKeysList = dataKey as Array; + for (let index = 0; index < dataKeysList.length; index++) { + this.datasourceData[dataIndex][key + '_' + index] = { + data: [] + }; + } + } else { + this.datasourceData[dataIndex][key] = { + data: [] + }; + } + } + } + this.datasourceOrigData = deepClone(this.datasourceData); + + const data: Array> = []; + for (let dataIndex = 0; dataIndex < pageData.data.length; dataIndex++) { + const entityData = pageData.data[dataIndex]; + this.processEntityData(entityData, dataIndex, false, + (data1, dataIndex1, dataKeyIndex) => { + if (!data[dataIndex1]) { + data[dataIndex1] = []; + } + data[dataIndex1][dataKeyIndex] = data1; + } + ); + } + + this.listeners.forEach((listener) => { + listener.dataLoaded(pageData, data, + listener.configDatasourceIndex); + }); + } + + private onDataUpdate(update: Array) { + for (const entityData of update) { + const dataIndex = this.entityIdToDataIndex[entityData.entityId.id]; + this.processEntityData(entityData, dataIndex, true, this.notifyListeners.bind(this)); + } + } + + private notifyListeners(data: DataSetHolder, dataIndex: number, dataKeyIndex: number, detectChanges: boolean) { + this.listeners.forEach((listener) => { + listener.dataUpdated(data, + listener.configDatasourceIndex, + dataIndex, dataKeyIndex, detectChanges); + }); + } + + private processEntityData(entityData: EntityData, dataIndex: number, aggregate: boolean, + dataUpdatedCb: DataUpdatedCb) { + if (this.entityDataSubscriptionOptions.type === widgetType.latest && entityData.latest) { + for (const type of Object.keys(entityData.latest)) { + const subscriptionData = this.toSubscriptionData(entityData.latest[type], false); + this.onData(subscriptionData, type, dataIndex, true, dataUpdatedCb); + } + } + if (this.entityDataSubscriptionOptions.type === widgetType.timeseries && entityData.timeseries) { + const subscriptionData = this.toSubscriptionData(entityData.timeseries, true); + if (aggregate) { + this.dataAggregators[dataIndex].onData({data: subscriptionData}, false, false, true); + } else { + this.onData(subscriptionData, DataKeyType.timeseries, dataIndex, true, dataUpdatedCb); + } + } + } + + private onData(sourceData: SubscriptionData, type: string, dataIndex: number, detectChanges: boolean, + dataUpdatedCb: DataUpdatedCb) { + for (const keyName of Object.keys(sourceData)) { + const keyData = sourceData[keyName]; + const key = `${keyName}_${type}`; + const dataKeyList = this.dataKeys[key] as Array; + for (let keyIndex = 0; dataKeyList && keyIndex < dataKeyList.length; keyIndex++) { + const datasourceKey = `${key}_${keyIndex}`; + if (this.datasourceData[dataIndex][datasourceKey].data) { + const dataKey = dataKeyList[keyIndex]; + const data: DataSet = []; + let prevSeries: [number, any]; + let prevOrigSeries: [number, any]; + let datasourceKeyData: DataSet; + let datasourceOrigKeyData: DataSet; + let update = false; + if (this.realtime) { + datasourceKeyData = []; + datasourceOrigKeyData = []; + } else { + datasourceKeyData = this.datasourceData[dataIndex][datasourceKey].data; + datasourceOrigKeyData = this.datasourceOrigData[dataIndex][datasourceKey].data; + } + if (datasourceKeyData.length > 0) { + prevSeries = datasourceKeyData[datasourceKeyData.length - 1]; + prevOrigSeries = datasourceOrigKeyData[datasourceOrigKeyData.length - 1]; + } else { + prevSeries = [0, 0]; + prevOrigSeries = [0, 0]; + } + this.datasourceOrigData[dataIndex][datasourceKey].data = []; + if (this.entityDataSubscriptionOptions.type === widgetType.timeseries) { + keyData.forEach((keySeries) => { + let series = keySeries; + const time = series[0]; + this.datasourceOrigData[dataIndex][datasourceKey].data.push(series); + let value = this.convertValue(series[1]); + if (dataKey.postFunc) { + value = dataKey.postFunc(time, value, prevSeries[1], prevOrigSeries[0], prevOrigSeries[1]); + } + prevOrigSeries = series; + series = [time, value]; + data.push(series); + prevSeries = series; + }); + update = true; + } else if (this.entityDataSubscriptionOptions.type === widgetType.latest) { + if (keyData.length > 0) { + let series = keyData[0]; + const time = series[0]; + this.datasourceOrigData[dataIndex][datasourceKey].data.push(series); + let value = this.convertValue(series[1]); + if (dataKey.postFunc) { + value = dataKey.postFunc(time, value, prevSeries[1], prevOrigSeries[0], prevOrigSeries[1]); + } + series = [time, value]; + data.push(series); + } + update = true; + } + if (update) { + this.datasourceData[datasourceKey].data = data; + dataUpdatedCb(this.datasourceData[dataIndex][datasourceKey], dataIndex, dataKey.index, detectChanges); + } + } + } + } + } + + private isNumeric(val: any): boolean { + return (val - parseFloat( val ) + 1) >= 0; + } + + private convertValue(val: string): any { + if (val && this.isNumeric(val)) { + return Number(val); + } else { + return val; + } + } + + private toSubscriptionData(sourceData: {[key: string]: TsValue | TsValue[]}, isTs: boolean): SubscriptionData { + const subsData: SubscriptionData = {}; + for (const keyName of Object.keys(sourceData)) { + const values = sourceData[keyName]; + const dataSet: [number, any][] = []; + if (isTs) { + (values as TsValue[]).forEach((keySeries) => { + dataSet.push([keySeries.ts, keySeries.value]); + }); + } else { + const tsValue = values as TsValue; + dataSet.push([tsValue.ts, tsValue.value]); + } + subsData[keyName] = dataSet; + } + return subsData; + } + + private createRealtimeDataAggregator(subsTw: SubscriptionTimewindow, + tsKeyNames: Array, + dataKeyType: DataKeyType, + dataIndex: number, + dataUpdatedCb: DataUpdatedCb): DataAggregator { + return new DataAggregator( + (data, detectChanges) => { + this.onData(data, dataKeyType, dataIndex, detectChanges, dataUpdatedCb); + }, + tsKeyNames, + subsTw.startTs, + subsTw.aggregation.limit, + subsTw.aggregation.type, + subsTw.aggregation.timeWindow, + subsTw.aggregation.interval, + subsTw.aggregation.stateData, + this.utils + ); + } + + private generateSeries(dataKey: SubscriptionDataKey, index: number, startTime: number, endTime: number): [number, any][] { + const data: [number, any][] = []; + let prevSeries: [number, any]; + const datasourceDataKey = `${dataKey.key}_${index}`; + const datasourceKeyData = this.datasourceData[0][datasourceDataKey].data; + if (datasourceKeyData.length > 0) { + prevSeries = datasourceKeyData[datasourceKeyData.length - 1]; + } else { + prevSeries = [0, 0]; + } + for (let time = startTime; time <= endTime && (this.timer || this.history); time += this.frequency) { + const value = dataKey.func(time, prevSeries[1]); + const series: [number, any] = [time, value]; + data.push(series); + prevSeries = series; + } + if (data.length > 0) { + dataKey.lastUpdateTime = data[data.length - 1][0]; + } + return data; + } + + private generateLatest(dataKey: SubscriptionDataKey, detectChanges: boolean) { + let prevSeries: [number, any]; + const datasourceKeyData = this.datasourceData[0][dataKey.key].data; + if (datasourceKeyData.length > 0) { + prevSeries = datasourceKeyData[datasourceKeyData.length - 1]; + } else { + prevSeries = [0, 0]; + } + const time = Date.now(); + const value = dataKey.func(time, prevSeries[1]); + const series: [number, any] = [time, value]; + this.datasourceData[0][dataKey.key].data = [series]; + this.listeners.forEach( + (listener) => { + listener.dataUpdated(this.datasourceData[0][dataKey.key], + listener.configDatasourceIndex, + 0, + dataKey.index, detectChanges); + } + ); + } + + private onTick(detectChanges: boolean) { + const now = this.utils.currentPerfTime(); + this.tickElapsed += now - this.tickScheduledTime; + this.tickScheduledTime = now; + + if (this.timer) { + clearTimeout(this.timer); + } + let key: string; + if (this.entityDataSubscriptionOptions.type === widgetType.timeseries) { + let startTime: number; + let endTime: number; + let delta: number; + const generatedData: SubscriptionDataHolder = { + data: {} + }; + if (!this.history) { + delta = Math.floor(this.tickElapsed / this.frequency); + } + const deltaElapsed = this.history ? this.frequency : delta * this.frequency; + this.tickElapsed = this.tickElapsed - deltaElapsed; + for (key of Object.keys(this.dataKeys)) { + const dataKeyList = this.dataKeys[key] as Array; + for (let index = 0; index < dataKeyList.length && (this.timer || this.history); index ++) { + const dataKey = dataKeyList[index]; + if (!startTime) { + if (this.realtime) { + if (dataKey.lastUpdateTime) { + startTime = dataKey.lastUpdateTime + this.frequency; + endTime = dataKey.lastUpdateTime + deltaElapsed; + } else { + startTime = this.entityDataSubscriptionOptions.subscriptionTimewindow.startTs; + endTime = startTime + this.entityDataSubscriptionOptions.subscriptionTimewindow.realtimeWindowMs + this.frequency; + if (this.entityDataSubscriptionOptions.subscriptionTimewindow.aggregation.type === AggregationType.NONE) { + const time = endTime - this.frequency * this.entityDataSubscriptionOptions.subscriptionTimewindow.aggregation.limit; + startTime = Math.max(time, startTime); + } + } + } else { + startTime = this.entityDataSubscriptionOptions.subscriptionTimewindow.fixedWindow.startTimeMs; + endTime = this.entityDataSubscriptionOptions.subscriptionTimewindow.fixedWindow.endTimeMs; + } + } + generatedData.data[`${dataKey.name}_${dataKey.index}`] = this.generateSeries(dataKey, index, startTime, endTime); + } + } + if (this.dataAggregators && this.dataAggregators.length) { + this.dataAggregators[0].onData(generatedData, true, this.history, detectChanges); + } + } else if (this.entityDataSubscriptionOptions.type === widgetType.latest) { + for (key of Object.keys(this.dataKeys)) { + this.generateLatest(this.dataKeys[key] as SubscriptionDataKey, detectChanges); + } + } + + if (!this.history) { + this.timer = setTimeout(this.onTick.bind(this, true), this.frequency); + } + } + +} diff --git a/ui-ngx/src/app/core/api/entity-data.service.ts b/ui-ngx/src/app/core/api/entity-data.service.ts new file mode 100644 index 0000000000..b562b3f5cf --- /dev/null +++ b/ui-ngx/src/app/core/api/entity-data.service.ts @@ -0,0 +1,108 @@ +/// +/// Copyright © 2016-2020 The Thingsboard Authors +/// +/// Licensed under the Apache License, Version 2.0 (the "License"); +/// you may not use this file except in compliance with the License. +/// You may obtain a copy of the License at +/// +/// http://www.apache.org/licenses/LICENSE-2.0 +/// +/// Unless required by applicable law or agreed to in writing, software +/// distributed under the License is distributed on an "AS IS" BASIS, +/// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +/// See the License for the specific language governing permissions and +/// limitations under the License. +/// + +import { DataSetHolder, Datasource, DatasourceType, widgetType } from '@shared/models/widget.models'; +import { SubscriptionTimewindow } from '@shared/models/time/time.models'; +import { EntityData, EntityDataPageLink, EntityFilter, KeyFilter } from '@shared/models/query/query.models'; +import { PageData } from '@shared/models/page/page-data'; +import { Injectable } from '@angular/core'; +import { TelemetryWebsocketService } from '@core/ws/telemetry-websocket.service'; +import { UtilsService } from '@core/services/utils.service'; +import { SubscriptionDataKey } from '@core/api/datasource-subcription'; +import { deepClone, objectHashCode } from '@core/utils'; +import { EntityDataSubscription, EntityDataSubscriptionOptions } from '@core/api/entity-data-subscription'; + +export interface EntityDataListener { + subscriptionType: widgetType; + subscriptionTimewindow: SubscriptionTimewindow; + configDatasource: Datasource; + configDatasourceIndex: number; + dataLoaded: (pageData: PageData, data: Array>, datasourceIndex: number) => void; + dataUpdated: (data: DataSetHolder, datasourceIndex: number, dataIndex: number, dataKeyIndex: number, detectChanges: boolean) => void; + updateRealtimeSubscription: () => SubscriptionTimewindow; + setRealtimeSubscription: (subscriptionTimewindow: SubscriptionTimewindow) => void; + entityDataSubscriptionKey?: number; +} + +@Injectable({ + providedIn: 'root' +}) +export class EntityDataService { + + private subscriptions: {[entityDataSubscriptionKey: string]: EntityDataSubscription} = {}; + + constructor(private telemetryService: TelemetryWebsocketService, + private utils: UtilsService) {} + + public subscribeToEntityData(listener: EntityDataListener) { + const datasource = listener.configDatasource; + if (datasource.type === DatasourceType.entity && (!datasource.entityFilter || !datasource.pageLink)) { + return; + } + const subscriptionDataKeys: Array = []; + datasource.dataKeys.forEach((dataKey) => { + const subscriptionDataKey: SubscriptionDataKey = { + name: dataKey.name, + type: dataKey.type, + funcBody: dataKey.funcBody, + postFuncBody: dataKey.postFuncBody + }; + subscriptionDataKeys.push(subscriptionDataKey); + }); + + const entityDataSubscriptionOptions: EntityDataSubscriptionOptions = { + datasourceType: datasource.type, + dataKeys: subscriptionDataKeys, + type: listener.subscriptionType + }; + + if (listener.subscriptionType === widgetType.timeseries) { + entityDataSubscriptionOptions.subscriptionTimewindow = deepClone(listener.subscriptionTimewindow); + } + if (entityDataSubscriptionOptions.datasourceType === DatasourceType.entity) { + entityDataSubscriptionOptions.entityFilter = datasource.entityFilter; + entityDataSubscriptionOptions.pageLink = datasource.pageLink; + entityDataSubscriptionOptions.keyFilters = datasource.keyFilters; + } + listener.entityDataSubscriptionKey = objectHashCode(entityDataSubscriptionOptions); + let subscription: EntityDataSubscription; + if (this.subscriptions[listener.entityDataSubscriptionKey]) { + subscription = this.subscriptions[listener.entityDataSubscriptionKey]; + subscription.syncListener(listener); + } else { + subscription = new EntityDataSubscription(entityDataSubscriptionOptions, + this.telemetryService, this.utils); + this.subscriptions[listener.entityDataSubscriptionKey] = subscription; + subscription.start(); + } + subscription.addListener(listener); + } + + public unsubscribeFromDatasource(listener: EntityDataListener) { + if (listener.entityDataSubscriptionKey) { + const subscription = this.subscriptions[listener.entityDataSubscriptionKey]; + if (subscription) { + subscription.removeListener(listener); + if (!subscription.hasListeners()) { + subscription.unsubscribe(); + delete this.subscriptions[listener.entityDataSubscriptionKey]; + } + } + listener.entityDataSubscriptionKey = null; + } + } + +} diff --git a/ui-ngx/src/app/core/api/widget-api.models.ts b/ui-ngx/src/app/core/api/widget-api.models.ts index 9a7d855d6e..74e0716669 100644 --- a/ui-ngx/src/app/core/api/widget-api.models.ts +++ b/ui-ngx/src/app/core/api/widget-api.models.ts @@ -41,6 +41,9 @@ import { EntityAliases } from '@shared/models/alias.models'; import { EntityInfo } from '@app/shared/models/entity.models'; import { IDashboardComponent } from '@home/models/dashboard-component.models'; import * as moment_ from 'moment'; +import { EntityDataPageLink, EntityFilter, KeyFilter } from '@shared/models/query/query.models'; +import { EntityDataService } from '@core/api/entity-data.service'; +import { PageData } from '@shared/models/page/page-data'; export interface TimewindowFunctions { onUpdateTimewindow: (startTimeMs: number, endTimeMs: number, interval?: number) => void; @@ -76,6 +79,7 @@ export interface WidgetActionsApi { export interface AliasInfo { alias?: string; stateEntity?: boolean; + entityFilter?: EntityFilter; currentEntity?: EntityInfo; selectedId?: string; resolvedEntities?: Array; @@ -169,7 +173,8 @@ export class WidgetSubscriptionContext { timeService: TimeService; deviceService: DeviceService; alarmService: AlarmService; - datasourceService: DatasourceService; + // datasourceService: DatasourceService; + entityDataService: EntityDataService; utils: UtilsService; raf: RafService; widgetUtils: IWidgetUtils; @@ -197,6 +202,8 @@ export interface WidgetSubscriptionOptions { alarmsMaxCountLoad?: number; alarmsFetchSize?: number; datasources?: Array; + keyFilters?: Array; + pageLink?: EntityDataPageLink; targetDeviceAliasIds?: Array; targetDeviceIds?: Array; useDashboardTimewindow?: boolean; @@ -230,6 +237,9 @@ export interface IWidgetSubscription { useDashboardTimewindow: boolean; legendData: LegendData; + + datasourcePages?: PageData[]; + dataPages?: PageData>[]; datasources?: Array; data?: Array; hiddenData?: Array<{data: DataSet}>; diff --git a/ui-ngx/src/app/core/api/widget-subscription.ts b/ui-ngx/src/app/core/api/widget-subscription.ts index bc90c626cf..1cbafd6483 100644 --- a/ui-ngx/src/app/core/api/widget-subscription.ts +++ b/ui-ngx/src/app/core/api/widget-subscription.ts @@ -47,13 +47,16 @@ import { Observable, ReplaySubject, Subject, throwError } from 'rxjs'; import { CancelAnimationFrame } from '@core/services/raf.service'; import { EntityType } from '@shared/models/entity-type.models'; import { AlarmInfo, AlarmSearchStatus } from '@shared/models/alarm.models'; -import { deepClone, isDefined, isEqual } from '@core/utils'; +import { createLabelFromDatasource, deepClone, isDefined, isEqual } from '@core/utils'; import { AlarmSourceListener } from '@core/http/alarm.service'; import { DatasourceListener } from '@core/api/datasource.service'; import { EntityId } from '@app/shared/models/id/entity-id'; import { DataKeyType } from '@shared/models/telemetry/telemetry.models'; import { entityFields } from '@shared/models/entity.models'; import * as moment_ from 'moment'; +import { PageData } from '@shared/models/page/page-data'; +import { EntityDataListener } from '@core/api/entity-data.service'; +import { EntityData, EntityDataPageLink, EntityKeyType } from '@shared/models/query/query.models'; const moment = moment_; @@ -70,9 +73,14 @@ export class WidgetSubscription implements IWidgetSubscription { subscriptionTimewindow: SubscriptionTimewindow; useDashboardTimewindow: boolean; + datasourcePages: PageData[]; + dataPages: PageData>[]; + entityDataListeners: Array; + configuredDatasources: Array; + data: Array; datasources: Array; - datasourceListeners: Array; + // datasourceListeners: Array; hiddenData: Array; legendData: LegendData; legendConfig: LegendConfig; @@ -197,8 +205,13 @@ export class WidgetSubscription implements IWidgetSubscription { this.callbacks.legendDataUpdated = this.callbacks.legendDataUpdated || (() => {}); this.callbacks.timeWindowUpdated = this.callbacks.timeWindowUpdated || (() => {}); - this.datasources = this.ctx.utils.validateDatasources(options.datasources); - this.datasourceListeners = []; + // this.datasources = this.ctx.utils.validateDatasources(options.datasources); + this.configuredDatasources = this.ctx.utils.validateDatasources(options.datasources); + this.entityDataListeners = []; + // this.datasourceListeners = []; + this.datasourcePages = []; + this.datasources = []; + this.dataPages = []; this.data = []; this.hiddenData = []; this.originalTimewindow = null; @@ -332,6 +345,35 @@ export class WidgetSubscription implements IWidgetSubscription { } private initDataSubscription(): Observable { + const initDataSubscriptionSubject = new ReplaySubject(1); + this.loadStDiff().subscribe(() => { + if (!this.ctx.aliasController) { + this.hasResolvedData = true; + // this.configureData(); + initDataSubscriptionSubject.next(); + initDataSubscriptionSubject.complete(); + } else { + this.ctx.aliasController.resolveDatasources(this.configuredDatasources).subscribe( + (datasources) => { + this.configuredDatasources = datasources; + if (datasources && datasources.length) { + this.hasResolvedData = true; + } + // this.configureData(); + initDataSubscriptionSubject.next(); + initDataSubscriptionSubject.complete(); + }, + (err) => { + this.notifyDataLoaded(); + initDataSubscriptionSubject.error(err); + } + ); + } + }); + return initDataSubscriptionSubject.asObservable(); + } + +/* private initDataSubscriptionOld(): Observable { const initDataSubscriptionSubject = new ReplaySubject(1); this.loadStDiff().subscribe(() => { if (!this.ctx.aliasController) { @@ -358,9 +400,9 @@ export class WidgetSubscription implements IWidgetSubscription { } }); return initDataSubscriptionSubject.asObservable(); - } + } */ - private configureData() { + /* private configureData() { const additionalDatasources: Datasource[] = []; let dataIndex = 0; let additionalKeysNumber = 0; @@ -448,9 +490,19 @@ export class WidgetSubscription implements IWidgetSubscription { if (this.displayLegend) { this.legendData.keys = this.legendData.keys.sort((key1, key2) => key1.dataKey.label.localeCompare(key2.dataKey.label)); } - } + } */ private resetData() { + this.data = []; + this.hiddenData = []; + if (this.displayLegend) { + this.legendData.keys = []; + this.legendData.data = []; + } + this.onDataUpdated(); + } + +/* private resetDataOld() { for (let i = 0; i < this.data.length; i++) { this.data[i].data = []; this.hiddenData[i].data = []; @@ -463,7 +515,7 @@ export class WidgetSubscription implements IWidgetSubscription { } } this.onDataUpdated(); - } + }*/ getFirstEntityInfo(): SubscriptionEntityInfo { let entityId: EntityId; @@ -740,6 +792,87 @@ export class WidgetSubscription implements IWidgetSubscription { } private doSubscribe() { + if (this.type === widgetType.rpc) { + return; + } + if (this.type === widgetType.alarm) { + this.alarmsSubscribe(); + } else { + this.notifyDataLoading(); + if (this.type === widgetType.timeseries && this.timeWindowConfig) { + this.updateRealtimeSubscription(); + if (this.comparisonEnabled) { + this.updateSubscriptionForComparison(); + } + if (this.subscriptionTimewindow.fixedWindow) { + this.onDataUpdated(); + } + } + // let index = 0; + const forceUpdate = !this.datasources.length; + this.configuredDatasources.forEach((datasource, index) => { + const listener: EntityDataListener = { + subscriptionType: this.type, + subscriptionTimewindow: this.subscriptionTimewindow, + configDatasource: datasource, + configDatasourceIndex: index, + dataLoaded: this.dataLoaded.bind(this), + dataUpdated: this.dataUpdated.bind(this), + updateRealtimeSubscription: () => { + this.subscriptionTimewindow = this.updateRealtimeSubscription(); + return this.subscriptionTimewindow; + }, + setRealtimeSubscription: (subscriptionTimewindow) => { + this.updateRealtimeSubscription(deepClone(subscriptionTimewindow)); + } + }; + + /*if (this.comparisonEnabled && datasource.isAdditional) { + listener.subscriptionTimewindow = this.timewindowForComparison; + listener.updateRealtimeSubscription = () => { + this.subscriptionTimewindow = this.updateSubscriptionForComparison(); + return this.subscriptionTimewindow; + }; + listener.setRealtimeSubscription = () => { + this.updateSubscriptionForComparison(); + }; + }*/ + +/* let entityFieldKey = false; + + for (let a = 0; a < datasource.dataKeys.length; a++) { + if (datasource.dataKeys[a].type !== DataKeyType.entityField) { + this.data[index + a].data = []; + } else { + entityFieldKey = true; + } + } + index += datasource.dataKeys.length;*/ + + this.entityDataListeners.push(listener); + // this.datasourceListeners.push(listener); + + // if (datasource.dataKeys.length) { + // this.ctx.datasourceService.subscribeToDatasource(listener); + // } + + this.ctx.entityDataService.subscribeToEntityData(listener); + + /* if (datasource.unresolvedStateEntity || entityFieldKey || + !datasource.dataKeys.length || + (datasource.type === DatasourceType.entity && !datasource.entityId) + ) { + forceUpdate = true; + }*/ + }); + if (forceUpdate) { + this.notifyDataLoaded(); + this.onDataUpdated(); + } + } + } + + /* private doSubscribeOld() { if (this.type === widgetType.rpc) { return; } @@ -814,7 +947,7 @@ export class WidgetSubscription implements IWidgetSubscription { this.onDataUpdated(); } } - } + } */ private alarmsSubscribe() { this.notifyDataLoading(); @@ -851,6 +984,20 @@ export class WidgetSubscription implements IWidgetSubscription { unsubscribe() { + if (this.type !== widgetType.rpc) { + if (this.type === widgetType.alarm) { + this.alarmsUnsubscribe(); + } else { + this.entityDataListeners.forEach((listener) => { + this.ctx.entityDataService.unsubscribeFromDatasource(listener); + }); + this.entityDataListeners.length = 0; + this.resetData(); + } + } + } + +/* unsubscribeOld() { if (this.type !== widgetType.rpc) { if (this.type === widgetType.alarm) { this.alarmsUnsubscribe(); @@ -862,7 +1009,7 @@ export class WidgetSubscription implements IWidgetSubscription { this.resetData(); } } - } + } */ private alarmsUnsubscribe() { if (this.alarmSourceListener) { @@ -970,7 +1117,180 @@ export class WidgetSubscription implements IWidgetSubscription { return this.timewindowForComparison; } - private dataUpdated(sourceData: DataSetHolder, datasourceIndex: number, dataKeyIndex: number, detectChanges: boolean) { + private dataLoaded(pageData: PageData, data: Array>, datasourceIndex: number) { + const datasource = this.configuredDatasources[datasourceIndex]; + const datasources = pageData.data.map((entityData, index) => + this.entityDataToDatasource(datasource, entityData, index) + ); + const datasourcesPage: PageData = { + data: datasources, + hasNext: pageData.hasNext, + totalElements: pageData.totalElements, + totalPages: pageData.totalPages + }; + this.datasourcePages[datasourceIndex] = datasourcesPage; + const datasourceData = datasources.map((datasourceElement, index) => + this.entityDataToDatasourceData(datasourceElement, data[index]) + ); + const datasourceDataPage: PageData> = { + data: datasourceData, + hasNext: pageData.hasNext, + totalElements: pageData.totalElements, + totalPages: pageData.totalPages + }; + this.dataPages[datasourceIndex] = datasourceDataPage; + this.configureLoadedData(); + this.notifyDataLoaded(); + } + + private configureLoadedData() { + this.datasources.length = 0; + this.data.length = 0; + this.hiddenData.length = 0; + if (this.displayLegend) { + this.legendData.keys.length = 0; + this.legendData.data.length = 0; + } + + let dataKeyIndex = 0; + this.configuredDatasources.forEach((configuredDatasource, datasourceIndex) => { + configuredDatasource.dataKeyStartIndex = dataKeyIndex; + const datasourcesPage = this.datasourcePages[datasourceIndex]; + const datasourceDataPage = this.dataPages[datasourceIndex]; + if (datasourcesPage) { + datasourcesPage.data.forEach((datasource, currentDatasourceIndex) => { + datasource.dataKeys.forEach((dataKey, currentDataKeyIndex) => { + const datasourceData = datasourceDataPage.data[currentDatasourceIndex][currentDataKeyIndex]; + this.data.push(datasourceData); + this.hiddenData.push({data: []}); + if (this.displayLegend) { + const legendKey: LegendKey = { + dataKey, + dataIndex: dataKeyIndex + }; + this.legendData.keys.push(legendKey); + const legendKeyData: LegendKeyData = { + min: null, + max: null, + avg: null, + total: null, + hidden: false + }; + this.legendData.data.push(legendKeyData); + } + dataKeyIndex++; + }); + this.datasources.push(datasource); + }); + } + } + ); + let index = 0; + this.datasources.forEach((datasource) => { + datasource.dataKeys.forEach((dataKey) => { + if (datasource.generated) { + dataKey._hash = Math.random(); + dataKey.color = this.ctx.utils.getMaterialColor(index); + } + index++; + }); + }); + if (this.displayLegend) { + this.legendData.keys = this.legendData.keys.sort((key1, key2) => key1.dataKey.label.localeCompare(key2.dataKey.label)); + } + if (this.caulculateLegendData) { + this.data.forEach((dataSetHolder, keyIndex) => { + this.updateLegend(keyIndex, dataSetHolder.data, false); + }); + this.callbacks.legendDataUpdated(this, true); + } + this.onDataUpdated(true); + } + + private entityDataToDatasourceData(datasource: Datasource, data: Array): Array { + return datasource.dataKeys.map((dataKey, keyIndex) => { + dataKey.hidden = dataKey.settings.hideDataByDefault ? true : false; + dataKey.inLegend = dataKey.settings.removeFromLegend ? false : true; + dataKey.pattern = dataKey.label; + dataKey.label = createLabelFromDatasource(datasource, dataKey.pattern); + const datasourceData: DatasourceData = { + datasource, + dataKey, + data: [] + }; + return datasourceData; + }); + } + + private entityDataToDatasource(configDatasource: Datasource, entityData: EntityData, index: number): Datasource { + const newDatasource = deepClone(configDatasource); + newDatasource.dataReceived = true; + newDatasource.entity = {}; + newDatasource.entityId = entityData.entityId.id; + newDatasource.entityType = entityData.entityId.entityType as EntityType; + if (configDatasource.type === DatasourceType.entity) { + let name; + let label; + if (entityData.latest && entityData.latest[EntityKeyType.ENTITY_FIELD]) { + const fields = entityData.latest[EntityKeyType.ENTITY_FIELD]; + if (fields.name) { + name = fields.name.value; + } + if (fields.label) { + label = fields.label.value; + } + } + name = name || 'TODO'; + label = label || 'TODO'; + newDatasource.name = name; + newDatasource.entityName = name; + newDatasource.entityLabel = label; + newDatasource.entityDescription = 'TODO'; + } + newDatasource.generated = index > 0 ? true : false; + return newDatasource; + } + + private dataUpdated(data: DataSetHolder, datasourceIndex: number, dataIndex: number, dataKeyIndex: number, detectChanges: boolean) { + const configuredDatasource = this.configuredDatasources[datasourceIndex]; + const startIndex = configuredDatasource.dataKeyStartIndex; + const dataKeysCount = configuredDatasource.dataKeys.length; + const index = startIndex + dataIndex*dataKeysCount + dataKeyIndex; + let update = true; + let currentData: DataSetHolder; + if (this.displayLegend && this.legendData.keys[index].dataKey.hidden) { + currentData = this.hiddenData[index]; + } else { + currentData = this.data[index]; + } + if (this.type === widgetType.latest) { + const prevData = currentData.data; + if (!data.data.length) { + update = false; + } else if (prevData && prevData[0] && prevData[0].length > 1 && data.data.length > 0) { + const prevTs = prevData[0][0]; + const prevValue = prevData[0][1]; + if (prevTs === data.data[0][0] && prevValue === data.data[0][1]) { + update = false; + } + } + } + if (update) { + if (this.subscriptionTimewindow && this.subscriptionTimewindow.realtimeWindowMs) { + this.updateTimewindow(); + if (this.timewindowForComparison && this.timewindowForComparison.realtimeWindowMs) { + this.updateComparisonTimewindow(); + } + } + currentData.data = data.data; + if (this.caulculateLegendData) { + this.updateLegend(index, data.data, detectChanges); + } + this.onDataUpdated(detectChanges); + } + } + +/* private dataUpdatedOld(sourceData: DataSetHolder, datasourceIndex: number, dataKeyIndex: number, detectChanges: boolean) { for (let x = 0; x < this.datasourceListeners.length; x++) { this.datasources[x].dataReceived = this.datasources[x].dataReceived === true; if (this.datasourceListeners[x].datasourceIndex === datasourceIndex && sourceData.data.length > 0) { @@ -1010,7 +1330,7 @@ export class WidgetSubscription implements IWidgetSubscription { } this.onDataUpdated(detectChanges); } - } + } */ private alarmsUpdated(alarms: Array) { this.notifyDataLoaded(); diff --git a/ui-ngx/src/app/core/http/entity.service.ts b/ui-ngx/src/app/core/http/entity.service.ts index 8437704faa..076f132cb8 100644 --- a/ui-ngx/src/app/core/http/entity.service.ts +++ b/ui-ngx/src/app/core/http/entity.service.ts @@ -52,7 +52,7 @@ import { EntitySearchQuery } from '@shared/models/relation.models'; import { EntityRelationService } from '@core/http/entity-relation.service'; -import { isDefined } from '@core/utils'; +import { deepClone, isDefined, isDefinedAndNotNull } from '@core/utils'; import { Asset, AssetSearchQuery } from '@shared/models/asset.models'; import { Device, DeviceCredentialsType, DeviceSearchQuery } from '@shared/models/device.models'; import { EntityViewSearchQuery } from '@shared/models/entity-view.models'; @@ -604,10 +604,11 @@ export class EntityService { public resolveAlias(entityAlias: EntityAlias, stateParams: StateParams): Observable { const filter = entityAlias.filter; - return this.resolveAliasFilter(filter, stateParams, -1, false).pipe( + return this.resolveAliasFilter(filter, stateParams).pipe( map((result) => { const aliasInfo: AliasInfo = { alias: entityAlias.alias, + entityFilter: result.entityFilter, stateEntity: result.stateEntity, entityParamName: result.entityParamName, resolveMultiple: filter.resolveMultiple @@ -621,11 +622,28 @@ export class EntityService { }) ); } +/* + public resolveEntityFilter(filter: EntityAliasFilter, stateParams: StateParams): EntityFilter { + const stateEntityInfo = this.getStateEntityInfo(filter, stateParams); + let result: EntityFilter = filter; + const stateEntityId = stateEntityInfo.entityId; + if (filter.type === AliasFilterType.stateEntity) { + result = { + singleEntity: stateEntityId, + type: AliasFilterType.singleEntity + }; + } else if (filter.rootStateEntity) { + let rootEntityType; + let rootEntityId; + + } + return result; + }*/ - public resolveAliasFilter(filter: EntityAliasFilter, stateParams: StateParams, - maxItems: number, failOnEmpty: boolean): Observable { + public resolveAliasFilter(filter: EntityAliasFilter, stateParams: StateParams): Observable { const result: EntityAliasFilterResult = { entities: [], + entityFilter: null, stateEntity: false }; if (filter.stateEntityParamName && filter.stateEntityParamName.length) { @@ -636,14 +654,21 @@ export class EntityService { switch (filter.type) { case AliasFilterType.singleEntity: const aliasEntityId = this.resolveAliasEntityId(filter.singleEntity.entityType, filter.singleEntity.id); - return this.getEntity(aliasEntityId.entityType as EntityType, aliasEntityId.id, {ignoreLoading: true, ignoreErrors: true}).pipe( + result.entityFilter = { + type: AliasFilterType.singleEntity, + singleEntity: aliasEntityId + }; + return of(result); + /*return this.getEntity(aliasEntityId.entityType as EntityType, aliasEntityId.id, {ignoreLoading: true, ignoreErrors: true}).pipe( map((entity) => { result.entities = this.entitiesToEntitiesInfo([entity]); return result; } - )); + ));*/ case AliasFilterType.entityList: - return this.getEntities(filter.entityType, filter.entityList, {ignoreLoading: true, ignoreErrors: true}).pipe( + result.entityFilter = deepClone(filter); + return of(result); + /*return this.getEntities(filter.entityType, filter.entityList, {ignoreLoading: true, ignoreErrors: true}).pipe( map((entities) => { if (entities && entities.length || !failOnEmpty) { result.entities = this.entitiesToEntitiesInfo(entities); @@ -652,9 +677,11 @@ export class EntityService { throw new Error(); } } - )); + ));*/ case AliasFilterType.entityName: - return this.getEntitiesByNameFilter(filter.entityType, filter.entityNameFilter, maxItems, + result.entityFilter = deepClone(filter); + return of(result); + /*return this.getEntitiesByNameFilter(filter.entityType, filter.entityNameFilter, maxItems, '', {ignoreLoading: true, ignoreErrors: true}).pipe( map((entities) => { if (entities && entities.length || !failOnEmpty) { @@ -665,11 +692,17 @@ export class EntityService { } } ) - ); + );*/ case AliasFilterType.stateEntity: result.stateEntity = true; if (stateEntityId) { - return this.getEntity(stateEntityId.entityType as EntityType, stateEntityId.id, {ignoreLoading: true, ignoreErrors: true}).pipe( + result.entityFilter = { + type: AliasFilterType.singleEntity, + singleEntity: stateEntityId + }; + } + return of(result); + /*return this.getEntity(stateEntityId.entityType as EntityType, stateEntityId.id, {ignoreLoading: true, ignoreErrors: true}).pipe( map((entity) => { result.entities = this.entitiesToEntitiesInfo([entity]); return result; @@ -677,9 +710,11 @@ export class EntityService { )); } else { return of(result); - } + }*/ case AliasFilterType.assetType: - return this.getEntitiesByNameFilter(EntityType.ASSET, filter.assetNameFilter, maxItems, + result.entityFilter = deepClone(filter); + return of(result); + /*return this.getEntitiesByNameFilter(EntityType.ASSET, filter.assetNameFilter, maxItems, filter.assetType, {ignoreLoading: true, ignoreErrors: true}).pipe( map((entities) => { if (entities && entities.length || !failOnEmpty) { @@ -690,9 +725,11 @@ export class EntityService { } } ) - ); + );*/ case AliasFilterType.deviceType: - return this.getEntitiesByNameFilter(EntityType.DEVICE, filter.deviceNameFilter, maxItems, + result.entityFilter = deepClone(filter); + return of(result); + /*return this.getEntitiesByNameFilter(EntityType.DEVICE, filter.deviceNameFilter, maxItems, filter.deviceType, {ignoreLoading: true, ignoreErrors: true}).pipe( map((entities) => { if (entities && entities.length || !failOnEmpty) { @@ -703,9 +740,11 @@ export class EntityService { } } ) - ); + );*/ case AliasFilterType.entityViewType: - return this.getEntitiesByNameFilter(EntityType.ENTITY_VIEW, filter.entityViewNameFilter, maxItems, + result.entityFilter = deepClone(filter); + return of(result); + /*return this.getEntitiesByNameFilter(EntityType.ENTITY_VIEW, filter.entityViewNameFilter, maxItems, filter.entityViewType, {ignoreLoading: true, ignoreErrors: true}).pipe( map((entities) => { if (entities && entities.length || !failOnEmpty) { @@ -716,7 +755,7 @@ export class EntityService { } } ) - ); + );*/ case AliasFilterType.relationsQuery: result.stateEntity = filter.rootStateEntity; let rootEntityType; @@ -730,7 +769,10 @@ export class EntityService { } if (rootEntityType && rootEntityId) { const relationQueryRootEntityId = this.resolveAliasEntityId(rootEntityType, rootEntityId); - const searchQuery: EntityRelationsQuery = { + result.entityFilter = deepClone(filter); + result.entityFilter.rootEntity = relationQueryRootEntityId; + return of(result); + /*const searchQuery: EntityRelationsQuery = { parameters: { rootId: relationQueryRootEntityId.id, rootType: relationQueryRootEntityId.entityType as EntityType, @@ -757,7 +799,7 @@ export class EntityService { return throwError(null); } }) - ); + );*/ } else { return of(result); } @@ -774,7 +816,10 @@ export class EntityService { } if (rootEntityType && rootEntityId) { const searchQueryRootEntityId = this.resolveAliasEntityId(rootEntityType, rootEntityId); - const searchQuery: EntitySearchQuery = { + result.entityFilter = deepClone(filter); + result.entityFilter.rootEntity = searchQueryRootEntityId; + return of(result); + /* const searchQuery: EntitySearchQuery = { parameters: { rootId: searchQueryRootEntityId.id, rootType: searchQueryRootEntityId.entityType as EntityType, @@ -811,7 +856,7 @@ export class EntityService { throw Error(); } }) - ); + );*/ } else { return of(result); } @@ -819,17 +864,18 @@ export class EntityService { } public checkEntityAlias(entityAlias: EntityAlias): Observable { - return this.resolveAliasFilter(entityAlias.filter, null, 1, true).pipe( + return this.resolveAliasFilter(entityAlias.filter, null).pipe( map((result) => { if (result.stateEntity) { return true; } else { - const entities = result.entities; + return isDefinedAndNotNull(result.entityFilter); + /*const entities = result.entities; if (entities && entities.length) { return true; } else { return false; - } + }*/ } }), catchError(err => of(false)) diff --git a/ui-ngx/src/app/core/ws/telemetry-websocket.service.ts b/ui-ngx/src/app/core/ws/telemetry-websocket.service.ts index 09d05c3906..5d3daa9467 100644 --- a/ui-ngx/src/app/core/ws/telemetry-websocket.service.ts +++ b/ui-ngx/src/app/core/ws/telemetry-websocket.service.ts @@ -16,8 +16,8 @@ import { Inject, Injectable, NgZone } from '@angular/core'; import { - AttributesSubscriptionCmd, - GetHistoryCmd, + AttributesSubscriptionCmd, EntityDataCmd, EntityDataUnsubscribeCmd, EntityDataUpdate, + GetHistoryCmd, isEntityDataUpdateMsg, SubscriptionCmd, SubscriptionUpdate, SubscriptionUpdateMsg, @@ -25,7 +25,7 @@ import { TelemetryPluginCmdsWrapper, TelemetryService, TelemetrySubscriber, - TimeseriesSubscriptionCmd + TimeseriesSubscriptionCmd, WebsocketDataMsg } from '@app/shared/models/telemetry/telemetry.models'; import { select, Store } from '@ngrx/store'; import { AppState } from '@core/core.state'; @@ -63,7 +63,7 @@ export class TelemetryWebsocketService implements TelemetryService { cmdsWrapper = new TelemetryPluginCmdsWrapper(); telemetryUri: string; - dataStream: WebSocketSubject; + dataStream: WebSocketSubject; constructor(private store: Store, private authService: AuthService, @@ -105,6 +105,8 @@ export class TelemetryWebsocketService implements TelemetryService { } } else if (subscriptionCommand instanceof GetHistoryCmd) { this.cmdsWrapper.historyCmds.push(subscriptionCommand); + } else if (subscriptionCommand instanceof EntityDataCmd) { + this.cmdsWrapper.entityDataCmds.push(subscriptionCommand); } } ); @@ -123,6 +125,10 @@ export class TelemetryWebsocketService implements TelemetryService { } else { this.cmdsWrapper.attrSubCmds.push(subscriptionCommand as AttributesSubscriptionCmd); } + } else if (subscriptionCommand instanceof EntityDataCmd) { + const entityDataUnsubscribeCmd = new EntityDataUnsubscribeCmd(); + entityDataUnsubscribeCmd.cmdId = subscriptionCommand.cmdId; + this.cmdsWrapper.entityDataUnsubscribeCmds.push(entityDataUnsubscribeCmd); } const cmdId = subscriptionCommand.cmdId; if (cmdId) { @@ -223,7 +229,7 @@ export class TelemetryWebsocketService implements TelemetryService { this.dataStream.subscribe((message) => { this.ngZone.runOutsideAngular(() => { - this.onMessage(message as SubscriptionUpdateMsg); + this.onMessage(message as WebsocketDataMsg); }); }, (error) => { @@ -252,13 +258,21 @@ export class TelemetryWebsocketService implements TelemetryService { } } - private onMessage(message: SubscriptionUpdateMsg) { + private onMessage(message: WebsocketDataMsg) { if (message.errorCode) { this.showWsError(message.errorCode, message.errorMsg); - } else if (message.subscriptionId) { - const subscriber = this.subscribersMap.get(message.subscriptionId); - if (subscriber) { - subscriber.onData(new SubscriptionUpdate(message)); + } else { + let subscriber: TelemetrySubscriber; + if (isEntityDataUpdateMsg(message)) { + subscriber = this.subscribersMap.get(message.cmdId); + if (subscriber) { + subscriber.onEntityData(new EntityDataUpdate(message)); + } + } else if (message.subscriptionId) { + subscriber = this.subscribersMap.get(message.subscriptionId); + if (subscriber) { + subscriber.onData(new SubscriptionUpdate(message)); + } } } this.checkToClose(); diff --git a/ui-ngx/src/app/modules/home/components/widget/widget.component.ts b/ui-ngx/src/app/modules/home/components/widget/widget.component.ts index a07e94d843..6cc72505a5 100644 --- a/ui-ngx/src/app/modules/home/components/widget/widget.component.ts +++ b/ui-ngx/src/app/modules/home/components/widget/widget.component.ts @@ -92,6 +92,7 @@ import { WidgetSubscription } from '@core/api/widget-subscription'; import { EntityService } from '@core/http/entity.service'; import { ServicesMap } from '@home/models/services.map'; import { ResizeObserver } from '@juggle/resize-observer'; +import { EntityDataService } from '@core/api/entity-data.service'; @Component({ selector: 'tb-widget', @@ -161,7 +162,8 @@ export class WidgetComponent extends PageComponent implements OnInit, AfterViewI private entityService: EntityService, private alarmService: AlarmService, private dashboardService: DashboardService, - private datasourceService: DatasourceService, + // private datasourceService: DatasourceService, + private entityDataService: EntityDataService, private utils: UtilsService, private raf: RafService, private ngZone: NgZone, @@ -292,7 +294,8 @@ export class WidgetComponent extends PageComponent implements OnInit, AfterViewI this.subscriptionContext.timeService = this.timeService; this.subscriptionContext.deviceService = this.deviceService; this.subscriptionContext.alarmService = this.alarmService; - this.subscriptionContext.datasourceService = this.datasourceService; + // this.subscriptionContext.datasourceService = this.datasourceService; + this.subscriptionContext.entityDataService = this.entityDataService; this.subscriptionContext.utils = this.utils; this.subscriptionContext.raf = this.raf; this.subscriptionContext.widgetUtils = this.widgetContext.utils; diff --git a/ui-ngx/src/app/shared/models/alias.models.ts b/ui-ngx/src/app/shared/models/alias.models.ts index 3009685321..0ccde879a0 100644 --- a/ui-ngx/src/app/shared/models/alias.models.ts +++ b/ui-ngx/src/app/shared/models/alias.models.ts @@ -18,6 +18,7 @@ import { EntityType } from '@shared/models/entity-type.models'; import { EntityId } from '@shared/models/id/entity-id'; import { EntitySearchDirection, EntityTypeFilter } from '@shared/models/relation.models'; import { EntityInfo } from './entity.models'; +import { EntityFilter } from '@shared/models/query/query.models'; export enum AliasFilterType { singleEntity = 'singleEntity', @@ -157,5 +158,6 @@ export interface EntityAliases { export interface EntityAliasFilterResult { entities: Array; stateEntity: boolean; + entityFilter: EntityFilter; entityParamName?: string; } diff --git a/ui-ngx/src/app/shared/models/query/query.models.ts b/ui-ngx/src/app/shared/models/query/query.models.ts new file mode 100644 index 0000000000..304968841e --- /dev/null +++ b/ui-ngx/src/app/shared/models/query/query.models.ts @@ -0,0 +1,157 @@ +/// +/// Copyright © 2016-2020 The Thingsboard Authors +/// +/// Licensed under the Apache License, Version 2.0 (the "License"); +/// you may not use this file except in compliance with the License. +/// You may obtain a copy of the License at +/// +/// http://www.apache.org/licenses/LICENSE-2.0 +/// +/// Unless required by applicable law or agreed to in writing, software +/// distributed under the License is distributed on an "AS IS" BASIS, +/// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +/// See the License for the specific language governing permissions and +/// limitations under the License. +/// + +import { AliasFilterType, EntityFilters } from '@shared/models/alias.models'; +import { EntityId } from '@shared/models/id/entity-id'; + +export enum EntityKeyType { + ATTRIBUTE = 'ATTRIBUTE', + CLIENT_ATTRIBUTE = 'CLIENT_ATTRIBUTE', + SHARED_ATTRIBUTE = 'SHARED_ATTRIBUTE', + SERVER_ATTRIBUTE = 'SERVER_ATTRIBUTE', + TIME_SERIES = 'TIME_SERIES', + ENTITY_FIELD = 'ENTITY_FIELD' +} + +export interface EntityKey { + type: EntityKeyType; + key: string; +} + +export enum FilterPredicateType { + STRING = 'STRING', + NUMERIC = 'NUMERIC', + BOOLEAN = 'BOOLEAN', + COMPLEX = 'COMPLEX' +} + +export enum StringOperation { + EQUAL = 'EQUAL', + NOT_EQUAL = 'NOT_EQUAL', + STARTS_WITH = 'STARTS_WITH', + ENDS_WITH = 'ENDS_WITH', + CONTAINS = 'CONTAINS', + NOT_CONTAIN = 'NOT_CONTAIN' +} + +export enum NumericOperation { + EQUAL = 'EQUAL', + NOT_EQUAL = 'NOT_EQUAL', + GREATER = 'GREATER', + LESS = 'LESS', + GREATER_OR_EQUAL = 'GREATER_OR_EQUAL', + LESS_OR_EQUAL = 'LESS_OR_EQUAL' +} + +export enum BooleanOperation { + EQUAL = 'EQUAL', + NOT_EQUAL = 'NOT_EQUAL' +} + +export enum ComplexOperation { + AND = 'AND', + OR = 'OR' +} + +export interface StringFilterPredicate { + operation: StringOperation; + value: string; + ignoreCase: boolean; +} + +export interface NumericFilterPredicate { + operation: NumericOperation; + value: number; +} + +export interface BooleanFilterPredicate { + operation: BooleanOperation; + value: boolean; +} + +export interface ComplexFilterPredicate { + operation: ComplexOperation; + predicates: Array; +} + +export type KeyFilterPredicates = StringFilterPredicate & + NumericFilterPredicate & + BooleanFilterPredicate & + ComplexFilterPredicate; + +export interface KeyFilterPredicate extends KeyFilterPredicates { + type?: FilterPredicateType; +} + +export interface KeyFilter { + key: EntityKey; + predicate: KeyFilterPredicate; +} + +export interface EntityFilter extends EntityFilters { + type?: AliasFilterType; +} + +export enum Direction { + ASC = 'ASC', + DESC = 'DESC' +} + +export interface EntityDataSortOrder { + key: EntityKey; + direction: Direction; +} + +export interface EntityDataPageLink { + pageSize: number; + page: number; + textSearch?: string; + sortOrder?: EntityDataSortOrder; +} + +export const defaultEntityDataPageLink: EntityDataPageLink = { + pageSize: 1024, + page: 0, + sortOrder: { + key: { + type: EntityKeyType.ENTITY_FIELD, + key: 'createdTime' + }, + direction: Direction.DESC + } +} + +export interface EntityCountQuery { + entityFilter: EntityFilter; +} + +export interface EntityDataQuery extends EntityCountQuery { + pageLink: EntityDataPageLink; + entityFields?: Array; + latestValues?: Array; + keyFilters?: Array; +} + +export interface TsValue { + ts: number; + value: string; +} + +export interface EntityData { + entityId: EntityId; + latest: {[entityKeyType: string]: {[key: string]: TsValue}}; + timeseries: {[key: string]: Array}; +} 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 d75239ec68..e3108004c8 100644 --- a/ui-ngx/src/app/shared/models/telemetry/telemetry.models.ts +++ b/ui-ngx/src/app/shared/models/telemetry/telemetry.models.ts @@ -21,6 +21,8 @@ import { Observable, ReplaySubject, Subject } from 'rxjs'; import { EntityId } from '@shared/models/id/entity-id'; import { map } from 'rxjs/operators'; import { NgZone } from '@angular/core'; +import { EntityData, EntityDataQuery } from '@shared/models/query/query.models'; +import { PageData } from '@shared/models/page/page-data'; export enum DataKeyType { timeseries = 'timeseries', @@ -79,8 +81,11 @@ export interface AttributeData { value: any; } -export interface TelemetryPluginCmd { +export interface WebsocketCmd { cmdId: number; +} + +export interface TelemetryPluginCmd extends WebsocketCmd { keys: string; } @@ -124,27 +129,69 @@ export class GetHistoryCmd implements TelemetryPluginCmd { agg: AggregationType; } +export interface EntityHistoryCmd { + keys: Array; + startTs: number; + endTs: number; + interval: number; + limit: number; + agg: AggregationType; +} + +export interface LatestValueCmd { + keys: Array; +} + +export interface TimeSeriesCmd { + keys: Array; + startTs: number; + timeWindow: number; + interval: number; + limit: number; + agg: AggregationType; +} + +export class EntityDataCmd implements WebsocketCmd { + cmdId: number; + query: EntityDataQuery; + historyCmd?: EntityHistoryCmd; + latestCmd?: LatestValueCmd; + tsCmd?: TimeSeriesCmd; +} + +export class EntityDataUnsubscribeCmd implements WebsocketCmd { + cmdId: number; +} + export class TelemetryPluginCmdsWrapper { attrSubCmds: Array; tsSubCmds: Array; historyCmds: Array; + entityDataCmds: Array; + entityDataUnsubscribeCmds: Array; constructor() { this.attrSubCmds = []; this.tsSubCmds = []; this.historyCmds = []; + this.entityDataCmds = []; + this.entityDataUnsubscribeCmds = []; } public hasCommands(): boolean { return this.tsSubCmds.length > 0 || this.historyCmds.length > 0 || - this.attrSubCmds.length > 0; + this.attrSubCmds.length > 0 || + this.entityDataCmds.length > 0 || + this.entityDataUnsubscribeCmds.length > 0; } public clear() { this.attrSubCmds.length = 0; this.tsSubCmds.length = 0; this.historyCmds.length = 0; + this.entityDataCmds.length = 0; + this.entityDataUnsubscribeCmds.length = 0; } public preparePublishCommands(maxCommands: number): TelemetryPluginCmdsWrapper { @@ -155,10 +202,14 @@ export class TelemetryPluginCmdsWrapper { preparedWrapper.historyCmds = this.popCmds(this.historyCmds, leftCount); leftCount -= preparedWrapper.historyCmds.length; preparedWrapper.attrSubCmds = this.popCmds(this.attrSubCmds, leftCount); + leftCount -= preparedWrapper.attrSubCmds.length; + preparedWrapper.entityDataCmds = this.popCmds(this.entityDataCmds, leftCount); + leftCount -= preparedWrapper.entityDataCmds.length; + preparedWrapper.entityDataUnsubscribeCmds = this.popCmds(this.entityDataUnsubscribeCmds, leftCount); return preparedWrapper; } - private popCmds(cmds: Array, leftCount: number): Array { + private popCmds(cmds: Array, leftCount: number): Array { const toPublish = Math.min(cmds.length, leftCount); if (toPublish > 0) { return cmds.splice(0, toPublish); @@ -182,6 +233,20 @@ export interface SubscriptionUpdateMsg extends SubscriptionDataHolder { errorMsg: string; } +export interface EntityDataUpdateMsg { + cmdId: number; + data?: PageData; + update?: Array; + errorCode: number; + errorMsg: string; +} + +export type WebsocketDataMsg = EntityDataUpdateMsg | SubscriptionUpdateMsg; + +export function isEntityDataUpdateMsg(message: WebsocketDataMsg): message is EntityDataUpdateMsg { + return (message as EntityDataUpdateMsg).cmdId !== undefined; +} + export class SubscriptionUpdate implements SubscriptionUpdateMsg { subscriptionId: number; errorCode: number; @@ -231,6 +296,22 @@ export class SubscriptionUpdate implements SubscriptionUpdateMsg { } } +export class EntityDataUpdate implements EntityDataUpdateMsg { + cmdId: number; + errorCode: number; + errorMsg: string; + data?: PageData; + update?: Array; + + constructor(msg: EntityDataUpdateMsg) { + this.cmdId = msg.cmdId; + this.errorCode = msg.errorCode; + this.errorMsg = msg.errorMsg; + this.data = msg.data; + this.update = msg.update; + } +} + export interface TelemetryService { subscribe(subscriber: TelemetrySubscriber); unsubscribe(subscriber: TelemetrySubscriber); @@ -239,13 +320,15 @@ export interface TelemetryService { export class TelemetrySubscriber { private dataSubject = new ReplaySubject(1); + private entityDataSubject = new ReplaySubject(1); private reconnectSubject = new Subject(); private zone: NgZone; - public subscriptionCommands: Array; + public subscriptionCommands: Array; public data$ = this.dataSubject.asObservable(); + public entityData$ = this.entityDataSubject.asObservable(); public reconnect$ = this.reconnectSubject.asObservable(); public static createEntityAttributesSubscription(telemetryService: TelemetryService, @@ -284,6 +367,7 @@ export class TelemetrySubscriber { public complete() { this.dataSubject.complete(); + this.entityDataSubject.complete(); this.reconnectSubject.complete(); } @@ -292,8 +376,9 @@ export class TelemetrySubscriber { let keys: string[]; const cmd = this.subscriptionCommands.find((command) => command.cmdId === cmdId); if (cmd) { - if (cmd.keys && cmd.keys.length) { - keys = cmd.keys.split(','); + const telemetryPluginCmd = cmd as TelemetryPluginCmd; + if (telemetryPluginCmd.keys && telemetryPluginCmd.keys.length) { + keys = telemetryPluginCmd.keys.split(','); } } message.prepareData(keys); @@ -308,6 +393,18 @@ export class TelemetrySubscriber { } } + public onEntityData(message: EntityDataUpdate) { + if (this.zone) { + this.zone.run( + () => { + this.entityDataSubject.next(message); + } + ); + } else { + this.entityDataSubject.next(message); + } + } + public onReconnected() { this.reconnectSubject.next(); } diff --git a/ui-ngx/src/app/shared/models/widget.models.ts b/ui-ngx/src/app/shared/models/widget.models.ts index 0ed1f6da07..3b3c29f44d 100644 --- a/ui-ngx/src/app/shared/models/widget.models.ts +++ b/ui-ngx/src/app/shared/models/widget.models.ts @@ -23,6 +23,7 @@ import { AlarmSearchStatus } from '@shared/models/alarm.models'; import { DataKeyType } from './telemetry/telemetry.models'; import { EntityId } from '@shared/models/id/entity-id'; import * as moment_ from 'moment'; +import { EntityDataPageLink, EntityFilter, KeyFilter } from '@shared/models/query/query.models'; export enum widgetType { timeseries = 'timeseries', @@ -263,6 +264,10 @@ export interface Datasource { entityDescription?: string; generated?: boolean; isAdditional?: boolean; + pageLink?: EntityDataPageLink; + keyFilters?: Array; + entityFilter?: EntityFilter; + dataKeyStartIndex?: number; [key: string]: any; }