|
|
|
@ -23,7 +23,7 @@ import { map } from 'rxjs/operators'; |
|
|
|
import { NgZone } from '@angular/core'; |
|
|
|
import { |
|
|
|
AlarmData, |
|
|
|
AlarmDataQuery, |
|
|
|
AlarmDataQuery, EntityCountQuery, |
|
|
|
EntityData, |
|
|
|
EntityDataQuery, |
|
|
|
EntityKey, |
|
|
|
@ -36,7 +36,8 @@ export enum DataKeyType { |
|
|
|
attribute = 'attribute', |
|
|
|
function = 'function', |
|
|
|
alarm = 'alarm', |
|
|
|
entityField = 'entityField' |
|
|
|
entityField = 'entityField', |
|
|
|
count = 'count' |
|
|
|
} |
|
|
|
|
|
|
|
export enum LatestTelemetry { |
|
|
|
@ -181,6 +182,11 @@ export class EntityDataCmd implements WebsocketCmd { |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
export class EntityCountCmd implements WebsocketCmd { |
|
|
|
cmdId: number; |
|
|
|
query?: EntityCountQuery; |
|
|
|
} |
|
|
|
|
|
|
|
export class AlarmDataCmd implements WebsocketCmd { |
|
|
|
cmdId: number; |
|
|
|
query?: AlarmDataQuery; |
|
|
|
@ -194,6 +200,10 @@ export class EntityDataUnsubscribeCmd implements WebsocketCmd { |
|
|
|
cmdId: number; |
|
|
|
} |
|
|
|
|
|
|
|
export class EntityCountUnsubscribeCmd implements WebsocketCmd { |
|
|
|
cmdId: number; |
|
|
|
} |
|
|
|
|
|
|
|
export class AlarmDataUnsubscribeCmd implements WebsocketCmd { |
|
|
|
cmdId: number; |
|
|
|
} |
|
|
|
@ -206,6 +216,8 @@ export class TelemetryPluginCmdsWrapper { |
|
|
|
entityDataUnsubscribeCmds: Array<EntityDataUnsubscribeCmd>; |
|
|
|
alarmDataCmds: Array<AlarmDataCmd>; |
|
|
|
alarmDataUnsubscribeCmds: Array<AlarmDataUnsubscribeCmd>; |
|
|
|
entityCountCmds: Array<EntityCountCmd>; |
|
|
|
entityCountUnsubscribeCmds: Array<EntityCountUnsubscribeCmd>; |
|
|
|
|
|
|
|
constructor() { |
|
|
|
this.attrSubCmds = []; |
|
|
|
@ -215,6 +227,8 @@ export class TelemetryPluginCmdsWrapper { |
|
|
|
this.entityDataUnsubscribeCmds = []; |
|
|
|
this.alarmDataCmds = []; |
|
|
|
this.alarmDataUnsubscribeCmds = []; |
|
|
|
this.entityCountCmds = []; |
|
|
|
this.entityCountUnsubscribeCmds = []; |
|
|
|
} |
|
|
|
|
|
|
|
public hasCommands(): boolean { |
|
|
|
@ -224,7 +238,9 @@ export class TelemetryPluginCmdsWrapper { |
|
|
|
this.entityDataCmds.length > 0 || |
|
|
|
this.entityDataUnsubscribeCmds.length > 0 || |
|
|
|
this.alarmDataCmds.length > 0 || |
|
|
|
this.alarmDataUnsubscribeCmds.length > 0; |
|
|
|
this.alarmDataUnsubscribeCmds.length > 0 || |
|
|
|
this.entityCountCmds.length > 0 || |
|
|
|
this.entityCountUnsubscribeCmds.length > 0; |
|
|
|
} |
|
|
|
|
|
|
|
public clear() { |
|
|
|
@ -235,6 +251,8 @@ export class TelemetryPluginCmdsWrapper { |
|
|
|
this.entityDataUnsubscribeCmds.length = 0; |
|
|
|
this.alarmDataCmds.length = 0; |
|
|
|
this.alarmDataUnsubscribeCmds.length = 0; |
|
|
|
this.entityCountCmds.length = 0; |
|
|
|
this.entityCountUnsubscribeCmds.length = 0; |
|
|
|
} |
|
|
|
|
|
|
|
public preparePublishCommands(maxCommands: number): TelemetryPluginCmdsWrapper { |
|
|
|
@ -253,6 +271,10 @@ export class TelemetryPluginCmdsWrapper { |
|
|
|
preparedWrapper.alarmDataCmds = this.popCmds(this.alarmDataCmds, leftCount); |
|
|
|
leftCount -= preparedWrapper.alarmDataCmds.length; |
|
|
|
preparedWrapper.alarmDataUnsubscribeCmds = this.popCmds(this.alarmDataUnsubscribeCmds, leftCount); |
|
|
|
leftCount -= preparedWrapper.alarmDataUnsubscribeCmds.length; |
|
|
|
preparedWrapper.entityCountCmds = this.popCmds(this.entityCountCmds, leftCount); |
|
|
|
leftCount -= preparedWrapper.entityCountCmds.length; |
|
|
|
preparedWrapper.entityCountUnsubscribeCmds = this.popCmds(this.entityCountUnsubscribeCmds, leftCount); |
|
|
|
return preparedWrapper; |
|
|
|
} |
|
|
|
|
|
|
|
@ -280,40 +302,54 @@ export interface SubscriptionUpdateMsg extends SubscriptionDataHolder { |
|
|
|
errorMsg: string; |
|
|
|
} |
|
|
|
|
|
|
|
export enum DataUpdateType { |
|
|
|
export enum CmdUpdateType { |
|
|
|
ENTITY_DATA = 'ENTITY_DATA', |
|
|
|
ALARM_DATA = 'ALARM_DATA' |
|
|
|
ALARM_DATA = 'ALARM_DATA', |
|
|
|
COUNT_DATA = 'COUNT_DATA' |
|
|
|
} |
|
|
|
|
|
|
|
export interface DataUpdateMsg<T> { |
|
|
|
export interface CmdUpdateMsg { |
|
|
|
cmdId: number; |
|
|
|
data?: PageData<T>; |
|
|
|
update?: Array<T>; |
|
|
|
errorCode: number; |
|
|
|
errorMsg: string; |
|
|
|
dataUpdateType: DataUpdateType; |
|
|
|
cmdUpdateType: CmdUpdateType; |
|
|
|
} |
|
|
|
|
|
|
|
export interface DataUpdateMsg<T> extends CmdUpdateMsg { |
|
|
|
data?: PageData<T>; |
|
|
|
update?: Array<T>; |
|
|
|
} |
|
|
|
|
|
|
|
export interface EntityDataUpdateMsg extends DataUpdateMsg<EntityData> { |
|
|
|
dataUpdateType: DataUpdateType.ENTITY_DATA; |
|
|
|
cmdUpdateType: CmdUpdateType.ENTITY_DATA; |
|
|
|
} |
|
|
|
|
|
|
|
export interface AlarmDataUpdateMsg extends DataUpdateMsg<AlarmData> { |
|
|
|
dataUpdateType: DataUpdateType.ALARM_DATA; |
|
|
|
cmdUpdateType: CmdUpdateType.ALARM_DATA; |
|
|
|
allowedEntities: number; |
|
|
|
totalEntities: number; |
|
|
|
} |
|
|
|
|
|
|
|
export type WebsocketDataMsg = AlarmDataUpdateMsg | EntityDataUpdateMsg | SubscriptionUpdateMsg; |
|
|
|
export interface EntityCountUpdateMsg extends CmdUpdateMsg { |
|
|
|
cmdUpdateType: CmdUpdateType.COUNT_DATA; |
|
|
|
count: number; |
|
|
|
} |
|
|
|
|
|
|
|
export type WebsocketDataMsg = AlarmDataUpdateMsg | EntityDataUpdateMsg | EntityCountUpdateMsg | SubscriptionUpdateMsg; |
|
|
|
|
|
|
|
export function isEntityDataUpdateMsg(message: WebsocketDataMsg): message is EntityDataUpdateMsg { |
|
|
|
const updateMsg = (message as DataUpdateMsg<any>); |
|
|
|
return updateMsg.cmdId !== undefined && updateMsg.dataUpdateType === DataUpdateType.ENTITY_DATA; |
|
|
|
const updateMsg = (message as CmdUpdateMsg); |
|
|
|
return updateMsg.cmdId !== undefined && updateMsg.cmdUpdateType === CmdUpdateType.ENTITY_DATA; |
|
|
|
} |
|
|
|
|
|
|
|
export function isAlarmDataUpdateMsg(message: WebsocketDataMsg): message is AlarmDataUpdateMsg { |
|
|
|
const updateMsg = (message as DataUpdateMsg<any>); |
|
|
|
return updateMsg.cmdId !== undefined && updateMsg.dataUpdateType === DataUpdateType.ALARM_DATA; |
|
|
|
const updateMsg = (message as CmdUpdateMsg); |
|
|
|
return updateMsg.cmdId !== undefined && updateMsg.cmdUpdateType === CmdUpdateType.ALARM_DATA; |
|
|
|
} |
|
|
|
|
|
|
|
export function isEntityCountUpdateMsg(message: WebsocketDataMsg): message is EntityCountUpdateMsg { |
|
|
|
const updateMsg = (message as CmdUpdateMsg); |
|
|
|
return updateMsg.cmdId !== undefined && updateMsg.cmdUpdateType === CmdUpdateType.COUNT_DATA; |
|
|
|
} |
|
|
|
|
|
|
|
export class SubscriptionUpdate implements SubscriptionUpdateMsg { |
|
|
|
@ -365,21 +401,28 @@ export class SubscriptionUpdate implements SubscriptionUpdateMsg { |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
export class DataUpdate<T> implements DataUpdateMsg<T> { |
|
|
|
export class CmdUpdate implements CmdUpdateMsg { |
|
|
|
cmdId: number; |
|
|
|
errorCode: number; |
|
|
|
errorMsg: string; |
|
|
|
data?: PageData<T>; |
|
|
|
update?: Array<T>; |
|
|
|
dataUpdateType: DataUpdateType; |
|
|
|
cmdUpdateType: CmdUpdateType; |
|
|
|
|
|
|
|
constructor(msg: DataUpdateMsg<T>) { |
|
|
|
constructor(msg: CmdUpdateMsg) { |
|
|
|
this.cmdId = msg.cmdId; |
|
|
|
this.errorCode = msg.errorCode; |
|
|
|
this.errorMsg = msg.errorMsg; |
|
|
|
this.cmdUpdateType = msg.cmdUpdateType; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
export class DataUpdate<T> extends CmdUpdate implements DataUpdateMsg<T> { |
|
|
|
data?: PageData<T>; |
|
|
|
update?: Array<T>; |
|
|
|
|
|
|
|
constructor(msg: DataUpdateMsg<T>) { |
|
|
|
super(msg); |
|
|
|
this.data = msg.data; |
|
|
|
this.update = msg.update; |
|
|
|
this.dataUpdateType = msg.dataUpdateType; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@ -400,6 +443,15 @@ export class AlarmDataUpdate extends DataUpdate<AlarmData> { |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
export class EntityCountUpdate extends CmdUpdate { |
|
|
|
count: number; |
|
|
|
|
|
|
|
constructor(msg: EntityCountUpdateMsg) { |
|
|
|
super(msg); |
|
|
|
this.count = msg.count; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
export interface TelemetryService { |
|
|
|
subscribe(subscriber: TelemetrySubscriber); |
|
|
|
update(subscriber: TelemetrySubscriber); |
|
|
|
@ -411,6 +463,7 @@ export class TelemetrySubscriber { |
|
|
|
private dataSubject = new ReplaySubject<SubscriptionUpdate>(1); |
|
|
|
private entityDataSubject = new ReplaySubject<EntityDataUpdate>(1); |
|
|
|
private alarmDataSubject = new ReplaySubject<AlarmDataUpdate>(1); |
|
|
|
private entityCountSubject = new ReplaySubject<EntityCountUpdate>(1); |
|
|
|
private reconnectSubject = new Subject(); |
|
|
|
|
|
|
|
private zone: NgZone; |
|
|
|
@ -420,6 +473,7 @@ export class TelemetrySubscriber { |
|
|
|
public data$ = this.dataSubject.asObservable(); |
|
|
|
public entityData$ = this.entityDataSubject.asObservable(); |
|
|
|
public alarmData$ = this.alarmDataSubject.asObservable(); |
|
|
|
public entityCount$ = this.entityCountSubject.asObservable(); |
|
|
|
public reconnect$ = this.reconnectSubject.asObservable(); |
|
|
|
|
|
|
|
public static createEntityAttributesSubscription(telemetryService: TelemetryService, |
|
|
|
@ -464,6 +518,7 @@ export class TelemetrySubscriber { |
|
|
|
this.dataSubject.complete(); |
|
|
|
this.entityDataSubject.complete(); |
|
|
|
this.alarmDataSubject.complete(); |
|
|
|
this.entityCountSubject.complete(); |
|
|
|
this.reconnectSubject.complete(); |
|
|
|
} |
|
|
|
|
|
|
|
@ -513,6 +568,18 @@ export class TelemetrySubscriber { |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
public onEntityCount(message: EntityCountUpdate) { |
|
|
|
if (this.zone) { |
|
|
|
this.zone.run( |
|
|
|
() => { |
|
|
|
this.entityCountSubject.next(message); |
|
|
|
} |
|
|
|
); |
|
|
|
} else { |
|
|
|
this.entityCountSubject.next(message); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
public onReconnected() { |
|
|
|
this.reconnectSubject.next(); |
|
|
|
} |
|
|
|
|