|
|
|
@ -19,14 +19,16 @@ import { |
|
|
|
DataEntry, |
|
|
|
DataSet, |
|
|
|
DataSetHolder, |
|
|
|
DatasourceType, IndexedData, |
|
|
|
DatasourceType, |
|
|
|
IndexedData, |
|
|
|
widgetType |
|
|
|
} from '@shared/models/widget.models'; |
|
|
|
import { |
|
|
|
AggregationType, |
|
|
|
ComparisonDuration, |
|
|
|
createTimewindowForComparison, |
|
|
|
getCurrentTime, IntervalMath, |
|
|
|
getCurrentTime, |
|
|
|
IntervalMath, |
|
|
|
SubscriptionTimewindow |
|
|
|
} from '@shared/models/time/time.models'; |
|
|
|
import { |
|
|
|
@ -51,7 +53,8 @@ import { |
|
|
|
IndexedSubscriptionData, |
|
|
|
IntervalType, |
|
|
|
NOT_SUPPORTED, |
|
|
|
SubscriptionData, SubscriptionDataEntry, |
|
|
|
SubscriptionData, |
|
|
|
SubscriptionDataEntry, |
|
|
|
TelemetrySubscriber |
|
|
|
} from '@shared/models/telemetry/telemetry.models'; |
|
|
|
import { UtilsService } from '@core/services/utils.service'; |
|
|
|
@ -61,9 +64,16 @@ import { PageData } from '@shared/models/page/page-data'; |
|
|
|
import { DataAggregator, onAggregatedData } from '@core/api/data-aggregator'; |
|
|
|
import { NULL_UUID } from '@shared/models/id/has-uuid'; |
|
|
|
import { EntityType } from '@shared/models/entity-type.models'; |
|
|
|
import { Observable, of, ReplaySubject, Subject } from 'rxjs'; |
|
|
|
import { firstValueFrom, from, Observable, of, ReplaySubject, Subject, switchMap } from 'rxjs'; |
|
|
|
import { EntityId } from '@shared/models/id/entity-id'; |
|
|
|
import { TelemetryWebsocketService } from '@core/ws/telemetry-websocket.service'; |
|
|
|
import { |
|
|
|
CompiledTbFunction, |
|
|
|
compileTbFunction, |
|
|
|
isNotEmptyTbFunction, |
|
|
|
TbFunction |
|
|
|
} from '@shared/models/js-function.models'; |
|
|
|
import { HttpClient } from '@angular/common/http'; |
|
|
|
import Timeout = NodeJS.Timeout; |
|
|
|
|
|
|
|
declare type DataKeyFunction = (time: number, prevValue: any) => any; |
|
|
|
@ -81,8 +91,8 @@ export interface SubscriptionDataKey { |
|
|
|
comparisonResultType?: ComparisonResultType; |
|
|
|
funcBody: string; |
|
|
|
func?: DataKeyFunction; |
|
|
|
postFuncBody: string; |
|
|
|
postFunc?: DataKeyPostFunction; |
|
|
|
postFuncBody: TbFunction; |
|
|
|
postFunc?: CompiledTbFunction<DataKeyPostFunction>; |
|
|
|
index?: number; |
|
|
|
listIndex?: number; |
|
|
|
key?: string; |
|
|
|
@ -109,8 +119,8 @@ export class EntityDataSubscription { |
|
|
|
|
|
|
|
constructor(private listener: EntityDataListener, |
|
|
|
private telemetryService: TelemetryWebsocketService, |
|
|
|
private utils: UtilsService) { |
|
|
|
this.initializeSubscription(); |
|
|
|
private utils: UtilsService, |
|
|
|
private http: HttpClient) { |
|
|
|
} |
|
|
|
|
|
|
|
private entityDataSubscriptionOptions = this.listener.subscriptionOptions; |
|
|
|
@ -193,7 +203,7 @@ export class EntityDataSubscription { |
|
|
|
return [[timestamp, value]]; |
|
|
|
} |
|
|
|
|
|
|
|
private initializeSubscription() { |
|
|
|
private async initializeSubscription() { |
|
|
|
for (let i = 0; i < this.entityDataSubscriptionOptions.dataKeys.length; i++) { |
|
|
|
const dataKey = deepClone(this.entityDataSubscriptionOptions.dataKeys[i]); |
|
|
|
this.dataKeysList.push(dataKey); |
|
|
|
@ -203,9 +213,10 @@ export class EntityDataSubscription { |
|
|
|
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; |
|
|
|
if (isNotEmptyTbFunction(dataKey.postFuncBody) && !dataKey.postFunc) { |
|
|
|
try { |
|
|
|
dataKey.postFunc = await firstValueFrom(compileTbFunction(this.http, dataKey.postFuncBody, 'time', 'value', 'prevValue', 'timePrev', 'prevOrigValue')); |
|
|
|
} catch (e) {} |
|
|
|
} |
|
|
|
} |
|
|
|
let key: string; |
|
|
|
@ -265,354 +276,358 @@ export class EntityDataSubscription { |
|
|
|
} |
|
|
|
|
|
|
|
public subscribe(): Observable<EntityDataLoadResult> { |
|
|
|
this.entityDataResolveSubject = new ReplaySubject(1); |
|
|
|
if (this.entityDataSubscriptionOptions.isPaginatedDataSubscription) { |
|
|
|
this.started = true; |
|
|
|
this.dataResolved = true; |
|
|
|
this.prepareSubscriptionTimewindow(); |
|
|
|
} |
|
|
|
if (this.datasourceType === DatasourceType.entity) { |
|
|
|
const entityFields: Array<EntityKey> = |
|
|
|
this.dataKeysList.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' |
|
|
|
}); |
|
|
|
} |
|
|
|
if (!entityFields.find(key => key.key === 'label')) { |
|
|
|
entityFields.push({ |
|
|
|
type: EntityKeyType.ENTITY_FIELD, |
|
|
|
key: 'label' |
|
|
|
}); |
|
|
|
} |
|
|
|
if (!entityFields.find(key => key.key === 'additionalInfo')) { |
|
|
|
entityFields.push({ |
|
|
|
type: EntityKeyType.ENTITY_FIELD, |
|
|
|
key: 'additionalInfo' |
|
|
|
}); |
|
|
|
} |
|
|
|
return from(this.initializeSubscription()).pipe( |
|
|
|
switchMap(() => { |
|
|
|
this.entityDataResolveSubject = new ReplaySubject(1); |
|
|
|
if (this.entityDataSubscriptionOptions.isPaginatedDataSubscription) { |
|
|
|
this.started = true; |
|
|
|
this.dataResolved = true; |
|
|
|
this.prepareSubscriptionTimewindow(); |
|
|
|
} |
|
|
|
if (this.datasourceType === DatasourceType.entity) { |
|
|
|
const entityFields: Array<EntityKey> = |
|
|
|
this.dataKeysList.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' |
|
|
|
}); |
|
|
|
} |
|
|
|
if (!entityFields.find(key => key.key === 'label')) { |
|
|
|
entityFields.push({ |
|
|
|
type: EntityKeyType.ENTITY_FIELD, |
|
|
|
key: 'label' |
|
|
|
}); |
|
|
|
} |
|
|
|
if (!entityFields.find(key => key.key === 'additionalInfo')) { |
|
|
|
entityFields.push({ |
|
|
|
type: EntityKeyType.ENTITY_FIELD, |
|
|
|
key: 'additionalInfo' |
|
|
|
}); |
|
|
|
} |
|
|
|
|
|
|
|
this.attrFields = this.dataKeysList.filter(dataKey => dataKey.type === DataKeyType.attribute).map( |
|
|
|
dataKey => ({ type: EntityKeyType.ATTRIBUTE, key: dataKey.name }) |
|
|
|
); |
|
|
|
this.attrFields = this.dataKeysList.filter(dataKey => dataKey.type === DataKeyType.attribute).map( |
|
|
|
dataKey => ({ type: EntityKeyType.ATTRIBUTE, key: dataKey.name }) |
|
|
|
); |
|
|
|
|
|
|
|
this.tsFields = this.dataKeysList. |
|
|
|
filter(dataKey => dataKey.type === DataKeyType.timeseries && |
|
|
|
this.tsFields = this.dataKeysList. |
|
|
|
filter(dataKey => dataKey.type === DataKeyType.timeseries && |
|
|
|
(!dataKey.aggregationType || dataKey.aggregationType === AggregationType.NONE) && !dataKey.latest).map( |
|
|
|
dataKey => ({ type: EntityKeyType.TIME_SERIES, key: dataKey.name }) |
|
|
|
); |
|
|
|
dataKey => ({ type: EntityKeyType.TIME_SERIES, key: dataKey.name }) |
|
|
|
); |
|
|
|
|
|
|
|
if (this.entityDataSubscriptionOptions.type === widgetType.timeseries) { |
|
|
|
const latestTsFields = this.dataKeysList. |
|
|
|
filter(dataKey => dataKey.type === DataKeyType.timeseries && dataKey.latest && |
|
|
|
(!dataKey.aggregationType || dataKey.aggregationType === AggregationType.NONE)).map( |
|
|
|
dataKey => ({ type: EntityKeyType.TIME_SERIES, key: dataKey.name }) |
|
|
|
); |
|
|
|
this.latestValues = this.attrFields.concat(latestTsFields); |
|
|
|
} else { |
|
|
|
this.latestValues = this.attrFields.concat(this.tsFields); |
|
|
|
} |
|
|
|
if (this.entityDataSubscriptionOptions.type === widgetType.timeseries) { |
|
|
|
const latestTsFields = this.dataKeysList. |
|
|
|
filter(dataKey => dataKey.type === DataKeyType.timeseries && dataKey.latest && |
|
|
|
(!dataKey.aggregationType || dataKey.aggregationType === AggregationType.NONE)).map( |
|
|
|
dataKey => ({ type: EntityKeyType.TIME_SERIES, key: dataKey.name }) |
|
|
|
); |
|
|
|
this.latestValues = this.attrFields.concat(latestTsFields); |
|
|
|
} else { |
|
|
|
this.latestValues = this.attrFields.concat(this.tsFields); |
|
|
|
} |
|
|
|
|
|
|
|
this.aggTsValues = this.dataKeysList. |
|
|
|
this.aggTsValues = this.dataKeysList. |
|
|
|
filter(dataKey => dataKey.type === DataKeyType.timeseries && |
|
|
|
dataKey.aggregationType && dataKey.aggregationType !== AggregationType.NONE && !dataKey.comparisonEnabled).map( |
|
|
|
dataKey => ({ id: dataKey.index, key: dataKey.name, agg: dataKey.aggregationType }) |
|
|
|
); |
|
|
|
|
|
|
|
this.aggTsComparisonValues = this.dataKeysList. |
|
|
|
filter(dataKey => dataKey.type === DataKeyType.timeseries && |
|
|
|
dataKey.aggregationType && dataKey.aggregationType !== AggregationType.NONE && dataKey.comparisonEnabled).map( |
|
|
|
dataKey => ({ id: dataKey.index, key: dataKey.name, agg: dataKey.aggregationType, |
|
|
|
previousValueOnly: dataKey.comparisonResultType === ComparisonResultType.PREVIOUS_VALUE }) |
|
|
|
); |
|
|
|
this.aggTsComparisonValues = this.dataKeysList. |
|
|
|
filter(dataKey => dataKey.type === DataKeyType.timeseries && |
|
|
|
dataKey.aggregationType && dataKey.aggregationType !== AggregationType.NONE && dataKey.comparisonEnabled).map( |
|
|
|
dataKey => ({ id: dataKey.index, key: dataKey.name, agg: dataKey.aggregationType, |
|
|
|
previousValueOnly: dataKey.comparisonResultType === ComparisonResultType.PREVIOUS_VALUE }) |
|
|
|
); |
|
|
|
|
|
|
|
this.subscriber = new TelemetrySubscriber(this.telemetryService); |
|
|
|
this.dataCommand = new EntityDataCmd(); |
|
|
|
this.subscriber = new TelemetrySubscriber(this.telemetryService); |
|
|
|
this.dataCommand = new EntityDataCmd(); |
|
|
|
|
|
|
|
let keyFilters = this.entityDataSubscriptionOptions.keyFilters; |
|
|
|
if (this.entityDataSubscriptionOptions.additionalKeyFilters) { |
|
|
|
if (keyFilters) { |
|
|
|
keyFilters = keyFilters.concat(this.entityDataSubscriptionOptions.additionalKeyFilters); |
|
|
|
} else { |
|
|
|
keyFilters = this.entityDataSubscriptionOptions.additionalKeyFilters; |
|
|
|
} |
|
|
|
} |
|
|
|
let keyFilters = this.entityDataSubscriptionOptions.keyFilters; |
|
|
|
if (this.entityDataSubscriptionOptions.additionalKeyFilters) { |
|
|
|
if (keyFilters) { |
|
|
|
keyFilters = keyFilters.concat(this.entityDataSubscriptionOptions.additionalKeyFilters); |
|
|
|
} else { |
|
|
|
keyFilters = this.entityDataSubscriptionOptions.additionalKeyFilters; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
this.dataCommand.query = { |
|
|
|
entityFilter: this.entityDataSubscriptionOptions.entityFilter, |
|
|
|
pageLink: this.entityDataSubscriptionOptions.pageLink, |
|
|
|
keyFilters, |
|
|
|
entityFields, |
|
|
|
latestValues: this.latestValues |
|
|
|
}; |
|
|
|
this.dataCommand.query = { |
|
|
|
entityFilter: this.entityDataSubscriptionOptions.entityFilter, |
|
|
|
pageLink: this.entityDataSubscriptionOptions.pageLink, |
|
|
|
keyFilters, |
|
|
|
entityFields, |
|
|
|
latestValues: this.latestValues |
|
|
|
}; |
|
|
|
|
|
|
|
if (this.entityDataSubscriptionOptions.isPaginatedDataSubscription) { |
|
|
|
this.prepareSubscriptionCommands(this.dataCommand); |
|
|
|
if (this.entityDataSubscriptionOptions.type === widgetType.timeseries) { |
|
|
|
this.subscriber.setTsOffset(this.subsTw.tsOffset); |
|
|
|
} else { |
|
|
|
this.subscriber.setTsOffset(this.latestTsOffset); |
|
|
|
} |
|
|
|
} |
|
|
|
if (this.entityDataSubscriptionOptions.isPaginatedDataSubscription) { |
|
|
|
this.prepareSubscriptionCommands(this.dataCommand); |
|
|
|
if (this.entityDataSubscriptionOptions.type === widgetType.timeseries) { |
|
|
|
this.subscriber.setTsOffset(this.subsTw.tsOffset); |
|
|
|
} else { |
|
|
|
this.subscriber.setTsOffset(this.latestTsOffset); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
this.subscriber.subscriptionCommands.push(this.dataCommand); |
|
|
|
this.subscriber.subscriptionCommands.push(this.dataCommand); |
|
|
|
|
|
|
|
this.subscriber.entityData$.subscribe( |
|
|
|
(entityDataUpdate) => { |
|
|
|
if (entityDataUpdate.data) { |
|
|
|
this.onPageData(entityDataUpdate.data); |
|
|
|
if (this.prematureUpdates) { |
|
|
|
for (const update of this.prematureUpdates) { |
|
|
|
this.onDataUpdate(update); |
|
|
|
this.subscriber.entityData$.subscribe( |
|
|
|
(entityDataUpdate) => { |
|
|
|
if (entityDataUpdate.data) { |
|
|
|
this.onPageData(entityDataUpdate.data); |
|
|
|
if (this.prematureUpdates) { |
|
|
|
for (const update of this.prematureUpdates) { |
|
|
|
this.onDataUpdate(update); |
|
|
|
} |
|
|
|
this.prematureUpdates = null; |
|
|
|
} |
|
|
|
} else if (entityDataUpdate.update) { |
|
|
|
if (!this.pageData) { |
|
|
|
if (!this.prematureUpdates) { |
|
|
|
this.prematureUpdates = []; |
|
|
|
} |
|
|
|
this.prematureUpdates.push(entityDataUpdate.update); |
|
|
|
} else { |
|
|
|
this.onDataUpdate(entityDataUpdate.update); |
|
|
|
} |
|
|
|
} |
|
|
|
this.prematureUpdates = null; |
|
|
|
} |
|
|
|
} else if (entityDataUpdate.update) { |
|
|
|
if (!this.pageData) { |
|
|
|
if (!this.prematureUpdates) { |
|
|
|
this.prematureUpdates = []; |
|
|
|
); |
|
|
|
|
|
|
|
this.subscriber.reconnect$.subscribe(() => { |
|
|
|
if (this.started) { |
|
|
|
const targetCommand = this.entityDataSubscriptionOptions.isPaginatedDataSubscription ? this.dataCommand : this.subsCommand; |
|
|
|
if (!this.history && (this.entityDataSubscriptionOptions.type === widgetType.timeseries && this.tsFields.length || |
|
|
|
this.aggTsValues.length > 0 && !this.isFloatingTimewindow)) { |
|
|
|
const newSubsTw = this.listener.updateRealtimeSubscription(); |
|
|
|
this.subsTw = newSubsTw; |
|
|
|
if (this.entityDataSubscriptionOptions.type === widgetType.timeseries && this.tsFields.length) { |
|
|
|
targetCommand.tsCmd.startTs = this.subsTw.startTs; |
|
|
|
targetCommand.tsCmd.timeWindow = this.subsTw.aggregation.timeWindow; |
|
|
|
if (typeof this.subsTw.aggregation.interval === 'number') { |
|
|
|
targetCommand.tsCmd.interval = this.subsTw.aggregation.interval; |
|
|
|
targetCommand.tsCmd.intervalType = IntervalType.MILLISECONDS; |
|
|
|
} else { |
|
|
|
targetCommand.tsCmd.intervalType = this.subsTw.aggregation.interval; |
|
|
|
} |
|
|
|
targetCommand.tsCmd.timeZoneId = this.subsTw.timezone; |
|
|
|
targetCommand.tsCmd.limit = this.subsTw.aggregation.limit; |
|
|
|
targetCommand.tsCmd.agg = this.subsTw.aggregation.type; |
|
|
|
targetCommand.tsCmd.fetchLatestPreviousPoint = this.subsTw.aggregation.stateData; |
|
|
|
this.dataAggregators.forEach((dataAggregator) => { |
|
|
|
dataAggregator.reset(newSubsTw); |
|
|
|
}); |
|
|
|
} |
|
|
|
if (this.aggTsValues.length > 0 && !this.isFloatingTimewindow) { |
|
|
|
targetCommand.aggTsCmd.startTs = this.subsTw.startTs; |
|
|
|
targetCommand.aggTsCmd.timeWindow = this.subsTw.aggregation.timeWindow; |
|
|
|
this.tsLatestDataAggregators.forEach((dataAggregator) => { |
|
|
|
dataAggregator.reset(newSubsTw); |
|
|
|
}); |
|
|
|
} |
|
|
|
} |
|
|
|
this.prematureUpdates.push(entityDataUpdate.update); |
|
|
|
if (this.entityDataSubscriptionOptions.type === widgetType.timeseries) { |
|
|
|
this.subscriber.setTsOffset(this.subsTw.tsOffset); |
|
|
|
} else { |
|
|
|
this.subscriber.setTsOffset(this.latestTsOffset); |
|
|
|
} |
|
|
|
targetCommand.query = this.dataCommand.query; |
|
|
|
this.subscriber.subscriptionCommands = [targetCommand]; |
|
|
|
} else { |
|
|
|
this.onDataUpdate(entityDataUpdate.update); |
|
|
|
this.subscriber.subscriptionCommands = [this.dataCommand]; |
|
|
|
} |
|
|
|
}); |
|
|
|
this.subscriber.subscribe(); |
|
|
|
} else if (this.datasourceType === DatasourceType.function) { |
|
|
|
let tsOffset = 0; |
|
|
|
if (this.entityDataSubscriptionOptions.type === widgetType.latest) { |
|
|
|
tsOffset = this.entityDataSubscriptionOptions.latestTsOffset; |
|
|
|
} else if (this.entityDataSubscriptionOptions.subscriptionTimewindow) { |
|
|
|
tsOffset = this.entityDataSubscriptionOptions.subscriptionTimewindow.tsOffset; |
|
|
|
} |
|
|
|
} |
|
|
|
); |
|
|
|
|
|
|
|
this.subscriber.reconnect$.subscribe(() => { |
|
|
|
if (this.started) { |
|
|
|
const targetCommand = this.entityDataSubscriptionOptions.isPaginatedDataSubscription ? this.dataCommand : this.subsCommand; |
|
|
|
if (!this.history && (this.entityDataSubscriptionOptions.type === widgetType.timeseries && this.tsFields.length || |
|
|
|
this.aggTsValues.length > 0 && !this.isFloatingTimewindow)) { |
|
|
|
const newSubsTw = this.listener.updateRealtimeSubscription(); |
|
|
|
this.subsTw = newSubsTw; |
|
|
|
if (this.entityDataSubscriptionOptions.type === widgetType.timeseries && this.tsFields.length) { |
|
|
|
targetCommand.tsCmd.startTs = this.subsTw.startTs; |
|
|
|
targetCommand.tsCmd.timeWindow = this.subsTw.aggregation.timeWindow; |
|
|
|
if (typeof this.subsTw.aggregation.interval === 'number') { |
|
|
|
targetCommand.tsCmd.interval = this.subsTw.aggregation.interval; |
|
|
|
targetCommand.tsCmd.intervalType = IntervalType.MILLISECONDS; |
|
|
|
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() + tsOffset, value: name} |
|
|
|
}; |
|
|
|
const pageData: PageData<EntityData> = { |
|
|
|
data: [entityData], |
|
|
|
hasNext: false, |
|
|
|
totalElements: 1, |
|
|
|
totalPages: 1 |
|
|
|
}; |
|
|
|
this.onPageData(pageData); |
|
|
|
} else if (this.datasourceType === DatasourceType.entityCount) { |
|
|
|
this.latestTsOffset = this.entityDataSubscriptionOptions.latestTsOffset; |
|
|
|
this.subscriber = new TelemetrySubscriber(this.telemetryService); |
|
|
|
this.subscriber.setTsOffset(this.latestTsOffset); |
|
|
|
this.countCommand = new EntityCountCmd(); |
|
|
|
let keyFilters = this.entityDataSubscriptionOptions.keyFilters; |
|
|
|
if (this.entityDataSubscriptionOptions.additionalKeyFilters) { |
|
|
|
if (keyFilters) { |
|
|
|
keyFilters = keyFilters.concat(this.entityDataSubscriptionOptions.additionalKeyFilters); |
|
|
|
} else { |
|
|
|
keyFilters = this.entityDataSubscriptionOptions.additionalKeyFilters; |
|
|
|
} |
|
|
|
} |
|
|
|
this.countCommand.query = { |
|
|
|
entityFilter: this.entityDataSubscriptionOptions.entityFilter, |
|
|
|
keyFilters |
|
|
|
}; |
|
|
|
this.subscriber.subscriptionCommands.push(this.countCommand); |
|
|
|
|
|
|
|
const entityId: EntityId = { |
|
|
|
id: NULL_UUID, |
|
|
|
entityType: null |
|
|
|
}; |
|
|
|
|
|
|
|
const countKey = this.dataKeysList[0]; |
|
|
|
|
|
|
|
let dataReceived = false; |
|
|
|
|
|
|
|
this.subscriber.entityCount$.subscribe( |
|
|
|
(entityCountUpdate) => { |
|
|
|
if (!dataReceived) { |
|
|
|
const entityData: EntityData = { |
|
|
|
entityId, |
|
|
|
latest: { |
|
|
|
[EntityKeyType.ENTITY_FIELD]: { |
|
|
|
name: { |
|
|
|
ts: Date.now() + this.latestTsOffset, |
|
|
|
value: DatasourceType.entityCount |
|
|
|
} |
|
|
|
}, |
|
|
|
[EntityKeyType.COUNT]: { |
|
|
|
[countKey.name]: { |
|
|
|
ts: Date.now() + this.latestTsOffset, |
|
|
|
value: entityCountUpdate.count + '' |
|
|
|
} |
|
|
|
} |
|
|
|
}, |
|
|
|
timeseries: {} |
|
|
|
}; |
|
|
|
const pageData: PageData<EntityData> = { |
|
|
|
data: [entityData], |
|
|
|
hasNext: false, |
|
|
|
totalElements: 1, |
|
|
|
totalPages: 1 |
|
|
|
}; |
|
|
|
this.onPageData(pageData); |
|
|
|
dataReceived = true; |
|
|
|
} else { |
|
|
|
targetCommand.tsCmd.intervalType = this.subsTw.aggregation.interval; |
|
|
|
const update: EntityData[] = [{ |
|
|
|
entityId, |
|
|
|
latest: { |
|
|
|
[EntityKeyType.COUNT]: { |
|
|
|
[countKey.name]: { |
|
|
|
ts: Date.now() + this.latestTsOffset, |
|
|
|
value: entityCountUpdate.count + '' |
|
|
|
} |
|
|
|
} |
|
|
|
}, |
|
|
|
timeseries: {} |
|
|
|
}]; |
|
|
|
this.onDataUpdate(update); |
|
|
|
} |
|
|
|
targetCommand.tsCmd.timeZoneId = this.subsTw.timezone; |
|
|
|
targetCommand.tsCmd.limit = this.subsTw.aggregation.limit; |
|
|
|
targetCommand.tsCmd.agg = this.subsTw.aggregation.type; |
|
|
|
targetCommand.tsCmd.fetchLatestPreviousPoint = this.subsTw.aggregation.stateData; |
|
|
|
this.dataAggregators.forEach((dataAggregator) => { |
|
|
|
dataAggregator.reset(newSubsTw); |
|
|
|
}); |
|
|
|
} |
|
|
|
if (this.aggTsValues.length > 0 && !this.isFloatingTimewindow) { |
|
|
|
targetCommand.aggTsCmd.startTs = this.subsTw.startTs; |
|
|
|
targetCommand.aggTsCmd.timeWindow = this.subsTw.aggregation.timeWindow; |
|
|
|
this.tsLatestDataAggregators.forEach((dataAggregator) => { |
|
|
|
dataAggregator.reset(newSubsTw); |
|
|
|
}); |
|
|
|
); |
|
|
|
this.subscriber.subscribe(); |
|
|
|
} else if (this.datasourceType === DatasourceType.alarmCount) { |
|
|
|
this.latestTsOffset = this.entityDataSubscriptionOptions.latestTsOffset; |
|
|
|
this.subscriber = new TelemetrySubscriber(this.telemetryService); |
|
|
|
this.subscriber.setTsOffset(this.latestTsOffset); |
|
|
|
this.alarmCountCommand = new AlarmCountCmd(); |
|
|
|
let keyFilters = this.entityDataSubscriptionOptions.keyFilters; |
|
|
|
if (this.entityDataSubscriptionOptions.additionalKeyFilters) { |
|
|
|
if (keyFilters) { |
|
|
|
keyFilters = keyFilters.concat(this.entityDataSubscriptionOptions.additionalKeyFilters); |
|
|
|
} else { |
|
|
|
keyFilters = this.entityDataSubscriptionOptions.additionalKeyFilters; |
|
|
|
} |
|
|
|
} |
|
|
|
if (this.entityDataSubscriptionOptions.type === widgetType.timeseries) { |
|
|
|
this.subscriber.setTsOffset(this.subsTw.tsOffset); |
|
|
|
} else { |
|
|
|
this.subscriber.setTsOffset(this.latestTsOffset); |
|
|
|
this.alarmCountCommand.query = { |
|
|
|
entityFilter: this.entityDataSubscriptionOptions.entityFilter, |
|
|
|
keyFilters |
|
|
|
}; |
|
|
|
if (this.entityDataSubscriptionOptions.alarmFilter) { |
|
|
|
this.alarmCountCommand.query = {...this.alarmCountCommand.query, ...this.entityDataSubscriptionOptions.alarmFilter}; |
|
|
|
} |
|
|
|
targetCommand.query = this.dataCommand.query; |
|
|
|
this.subscriber.subscriptionCommands = [targetCommand]; |
|
|
|
} else { |
|
|
|
this.subscriber.subscriptionCommands = [this.dataCommand]; |
|
|
|
} |
|
|
|
}); |
|
|
|
this.subscriber.subscribe(); |
|
|
|
} else if (this.datasourceType === DatasourceType.function) { |
|
|
|
let tsOffset = 0; |
|
|
|
if (this.entityDataSubscriptionOptions.type === widgetType.latest) { |
|
|
|
tsOffset = this.entityDataSubscriptionOptions.latestTsOffset; |
|
|
|
} else if (this.entityDataSubscriptionOptions.subscriptionTimewindow) { |
|
|
|
tsOffset = this.entityDataSubscriptionOptions.subscriptionTimewindow.tsOffset; |
|
|
|
} |
|
|
|
|
|
|
|
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() + tsOffset, value: name} |
|
|
|
}; |
|
|
|
const pageData: PageData<EntityData> = { |
|
|
|
data: [entityData], |
|
|
|
hasNext: false, |
|
|
|
totalElements: 1, |
|
|
|
totalPages: 1 |
|
|
|
}; |
|
|
|
this.onPageData(pageData); |
|
|
|
} else if (this.datasourceType === DatasourceType.entityCount) { |
|
|
|
this.latestTsOffset = this.entityDataSubscriptionOptions.latestTsOffset; |
|
|
|
this.subscriber = new TelemetrySubscriber(this.telemetryService); |
|
|
|
this.subscriber.setTsOffset(this.latestTsOffset); |
|
|
|
this.countCommand = new EntityCountCmd(); |
|
|
|
let keyFilters = this.entityDataSubscriptionOptions.keyFilters; |
|
|
|
if (this.entityDataSubscriptionOptions.additionalKeyFilters) { |
|
|
|
if (keyFilters) { |
|
|
|
keyFilters = keyFilters.concat(this.entityDataSubscriptionOptions.additionalKeyFilters); |
|
|
|
} else { |
|
|
|
keyFilters = this.entityDataSubscriptionOptions.additionalKeyFilters; |
|
|
|
} |
|
|
|
} |
|
|
|
this.countCommand.query = { |
|
|
|
entityFilter: this.entityDataSubscriptionOptions.entityFilter, |
|
|
|
keyFilters |
|
|
|
}; |
|
|
|
this.subscriber.subscriptionCommands.push(this.countCommand); |
|
|
|
this.subscriber.subscriptionCommands.push(this.alarmCountCommand); |
|
|
|
|
|
|
|
const entityId: EntityId = { |
|
|
|
id: NULL_UUID, |
|
|
|
entityType: null |
|
|
|
}; |
|
|
|
|
|
|
|
const countKey = this.dataKeysList[0]; |
|
|
|
|
|
|
|
let dataReceived = false; |
|
|
|
const entityId: EntityId = { |
|
|
|
id: NULL_UUID, |
|
|
|
entityType: null |
|
|
|
}; |
|
|
|
|
|
|
|
this.subscriber.entityCount$.subscribe( |
|
|
|
(entityCountUpdate) => { |
|
|
|
if (!dataReceived) { |
|
|
|
const entityData: EntityData = { |
|
|
|
entityId, |
|
|
|
latest: { |
|
|
|
[EntityKeyType.ENTITY_FIELD]: { |
|
|
|
name: { |
|
|
|
ts: Date.now() + this.latestTsOffset, |
|
|
|
value: DatasourceType.entityCount |
|
|
|
} |
|
|
|
}, |
|
|
|
[EntityKeyType.COUNT]: { |
|
|
|
[countKey.name]: { |
|
|
|
ts: Date.now() + this.latestTsOffset, |
|
|
|
value: entityCountUpdate.count + '' |
|
|
|
} |
|
|
|
} |
|
|
|
}, |
|
|
|
timeseries: {} |
|
|
|
}; |
|
|
|
const pageData: PageData<EntityData> = { |
|
|
|
data: [entityData], |
|
|
|
hasNext: false, |
|
|
|
totalElements: 1, |
|
|
|
totalPages: 1 |
|
|
|
}; |
|
|
|
this.onPageData(pageData); |
|
|
|
dataReceived = true; |
|
|
|
} else { |
|
|
|
const update: EntityData[] = [{ |
|
|
|
entityId, |
|
|
|
latest: { |
|
|
|
[EntityKeyType.COUNT]: { |
|
|
|
[countKey.name]: { |
|
|
|
ts: Date.now() + this.latestTsOffset, |
|
|
|
value: entityCountUpdate.count + '' |
|
|
|
} |
|
|
|
} |
|
|
|
}, |
|
|
|
timeseries: {} |
|
|
|
}]; |
|
|
|
this.onDataUpdate(update); |
|
|
|
} |
|
|
|
const countKey = this.dataKeysList[0]; |
|
|
|
|
|
|
|
let dataReceived = false; |
|
|
|
|
|
|
|
this.subscriber.alarmCount$.subscribe( |
|
|
|
(alarmCountUpdate) => { |
|
|
|
if (!dataReceived) { |
|
|
|
const entityData: EntityData = { |
|
|
|
entityId, |
|
|
|
latest: { |
|
|
|
[EntityKeyType.ENTITY_FIELD]: { |
|
|
|
name: { |
|
|
|
ts: Date.now() + this.latestTsOffset, |
|
|
|
value: DatasourceType.alarmCount |
|
|
|
} |
|
|
|
}, |
|
|
|
[EntityKeyType.COUNT]: { |
|
|
|
[countKey.name]: { |
|
|
|
ts: Date.now() + this.latestTsOffset, |
|
|
|
value: alarmCountUpdate.count + '' |
|
|
|
} |
|
|
|
} |
|
|
|
}, |
|
|
|
timeseries: {} |
|
|
|
}; |
|
|
|
const pageData: PageData<EntityData> = { |
|
|
|
data: [entityData], |
|
|
|
hasNext: false, |
|
|
|
totalElements: 1, |
|
|
|
totalPages: 1 |
|
|
|
}; |
|
|
|
this.onPageData(pageData); |
|
|
|
dataReceived = true; |
|
|
|
} else { |
|
|
|
const update: EntityData[] = [{ |
|
|
|
entityId, |
|
|
|
latest: { |
|
|
|
[EntityKeyType.COUNT]: { |
|
|
|
[countKey.name]: { |
|
|
|
ts: Date.now() + this.latestTsOffset, |
|
|
|
value: alarmCountUpdate.count + '' |
|
|
|
} |
|
|
|
} |
|
|
|
}, |
|
|
|
timeseries: {} |
|
|
|
}]; |
|
|
|
this.onDataUpdate(update); |
|
|
|
} |
|
|
|
} |
|
|
|
); |
|
|
|
this.subscriber.subscribe(); |
|
|
|
} |
|
|
|
); |
|
|
|
this.subscriber.subscribe(); |
|
|
|
} else if (this.datasourceType === DatasourceType.alarmCount) { |
|
|
|
this.latestTsOffset = this.entityDataSubscriptionOptions.latestTsOffset; |
|
|
|
this.subscriber = new TelemetrySubscriber(this.telemetryService); |
|
|
|
this.subscriber.setTsOffset(this.latestTsOffset); |
|
|
|
this.alarmCountCommand = new AlarmCountCmd(); |
|
|
|
let keyFilters = this.entityDataSubscriptionOptions.keyFilters; |
|
|
|
if (this.entityDataSubscriptionOptions.additionalKeyFilters) { |
|
|
|
if (keyFilters) { |
|
|
|
keyFilters = keyFilters.concat(this.entityDataSubscriptionOptions.additionalKeyFilters); |
|
|
|
if (this.entityDataSubscriptionOptions.isPaginatedDataSubscription) { |
|
|
|
return of(null); |
|
|
|
} else { |
|
|
|
keyFilters = this.entityDataSubscriptionOptions.additionalKeyFilters; |
|
|
|
return this.entityDataResolveSubject.asObservable(); |
|
|
|
} |
|
|
|
} |
|
|
|
this.alarmCountCommand.query = { |
|
|
|
entityFilter: this.entityDataSubscriptionOptions.entityFilter, |
|
|
|
keyFilters |
|
|
|
}; |
|
|
|
if (this.entityDataSubscriptionOptions.alarmFilter) { |
|
|
|
this.alarmCountCommand.query = {...this.alarmCountCommand.query, ...this.entityDataSubscriptionOptions.alarmFilter}; |
|
|
|
} |
|
|
|
this.subscriber.subscriptionCommands.push(this.alarmCountCommand); |
|
|
|
|
|
|
|
const entityId: EntityId = { |
|
|
|
id: NULL_UUID, |
|
|
|
entityType: null |
|
|
|
}; |
|
|
|
|
|
|
|
const countKey = this.dataKeysList[0]; |
|
|
|
|
|
|
|
let dataReceived = false; |
|
|
|
|
|
|
|
this.subscriber.alarmCount$.subscribe( |
|
|
|
(alarmCountUpdate) => { |
|
|
|
if (!dataReceived) { |
|
|
|
const entityData: EntityData = { |
|
|
|
entityId, |
|
|
|
latest: { |
|
|
|
[EntityKeyType.ENTITY_FIELD]: { |
|
|
|
name: { |
|
|
|
ts: Date.now() + this.latestTsOffset, |
|
|
|
value: DatasourceType.alarmCount |
|
|
|
} |
|
|
|
}, |
|
|
|
[EntityKeyType.COUNT]: { |
|
|
|
[countKey.name]: { |
|
|
|
ts: Date.now() + this.latestTsOffset, |
|
|
|
value: alarmCountUpdate.count + '' |
|
|
|
} |
|
|
|
} |
|
|
|
}, |
|
|
|
timeseries: {} |
|
|
|
}; |
|
|
|
const pageData: PageData<EntityData> = { |
|
|
|
data: [entityData], |
|
|
|
hasNext: false, |
|
|
|
totalElements: 1, |
|
|
|
totalPages: 1 |
|
|
|
}; |
|
|
|
this.onPageData(pageData); |
|
|
|
dataReceived = true; |
|
|
|
} else { |
|
|
|
const update: EntityData[] = [{ |
|
|
|
entityId, |
|
|
|
latest: { |
|
|
|
[EntityKeyType.COUNT]: { |
|
|
|
[countKey.name]: { |
|
|
|
ts: Date.now() + this.latestTsOffset, |
|
|
|
value: alarmCountUpdate.count + '' |
|
|
|
} |
|
|
|
} |
|
|
|
}, |
|
|
|
timeseries: {} |
|
|
|
}]; |
|
|
|
this.onDataUpdate(update); |
|
|
|
} |
|
|
|
} |
|
|
|
); |
|
|
|
this.subscriber.subscribe(); |
|
|
|
} |
|
|
|
if (this.entityDataSubscriptionOptions.isPaginatedDataSubscription) { |
|
|
|
return of(null); |
|
|
|
} else { |
|
|
|
return this.entityDataResolveSubject.asObservable(); |
|
|
|
} |
|
|
|
}) |
|
|
|
); |
|
|
|
} |
|
|
|
|
|
|
|
public start() { |
|
|
|
@ -1104,7 +1119,7 @@ export class EntityDataSubscription { |
|
|
|
this.datasourceOrigData[dataIndex][datasourceKey].data.push([series[0], series[1], series[2]]); |
|
|
|
let value = EntityDataSubscription.convertValue(series[1]); |
|
|
|
if (dataKey.postFunc) { |
|
|
|
value = dataKey.postFunc(time, value, prevSeries[1], prevOrigSeries[0], prevOrigSeries[1]); |
|
|
|
value = dataKey.postFunc.execute(time, value, prevSeries[1], prevOrigSeries[0], prevOrigSeries[1]); |
|
|
|
} |
|
|
|
prevOrigSeries = [series[0], series[1], series[2]]; |
|
|
|
series = [series[0], value, series[2]]; |
|
|
|
@ -1119,7 +1134,7 @@ export class EntityDataSubscription { |
|
|
|
this.datasourceOrigData[dataIndex][datasourceKey].data.push([series[0], series[1], series[2]]); |
|
|
|
let value = EntityDataSubscription.convertValue(series[1]); |
|
|
|
if (dataKey.postFunc) { |
|
|
|
value = dataKey.postFunc(time, value, prevSeries[1], prevOrigSeries[0], prevOrigSeries[1]); |
|
|
|
value = dataKey.postFunc.execute(time, value, prevSeries[1], prevOrigSeries[0], prevOrigSeries[1]); |
|
|
|
} |
|
|
|
series = [time, value, series[2]]; |
|
|
|
data.push([series[0], series[1], series[2]]); |
|
|
|
|