diff --git a/ui-ngx/src/app/core/ws/notification-websocket.service.ts b/ui-ngx/src/app/core/ws/notification-websocket.service.ts index 8f219c029d..cb33c23635 100644 --- a/ui-ngx/src/app/core/ws/notification-websocket.service.ts +++ b/ui-ngx/src/app/core/ws/notification-websocket.service.ts @@ -15,117 +15,45 @@ /// import { Inject, Injectable, NgZone } from '@angular/core'; +import { + TelemetryPluginCmdsWrapper, + TelemetrySubscriber, + WebsocketDataMsg +} from '@shared/models/telemetry/telemetry.models'; +import { TelemetryWebsocketService } from '@core/ws/telemetry-websocket.service'; import { Store } from '@ngrx/store'; import { AppState } from '@core/core.state'; import { AuthService } from '@core/auth/auth.service'; import { WINDOW } from '@core/services/window.service'; -import { - isNotificationCountUpdateMsg, - isNotificationsUpdateMsg, - MarkAllAsReadCmd, - MarkAsReadCmd, - NotificationCountUpdate, - NotificationPluginCmdWrapper, - NotificationSubscriber, - NotificationsUpdate, - UnreadCountSubCmd, - UnreadSubCmd, - UnsubscribeCmd, - WebsocketNotificationMsg -} from '@shared/models/websocket/notification-ws.models'; import { WebsocketService } from '@core/ws/websocket.service'; - // @dynamic @Injectable({ providedIn: 'root' }) -export class NotificationWebsocketService extends WebsocketService { +export class NotificationWebsocketService extends WebsocketService { - cmdWrapper: NotificationPluginCmdWrapper; - - constructor(protected store: Store, + constructor(private telemetryWebsocketService: TelemetryWebsocketService, + protected store: Store, protected authService: AuthService, protected ngZone: NgZone, @Inject(WINDOW) protected window: Window) { - super(store, authService, ngZone, 'api/ws/plugins/notifications', new NotificationPluginCmdWrapper(), window); - this.errorName = 'WebSocket Notification Error'; + super(store, authService, ngZone, 'api/ws/plugins/telemetry', new TelemetryPluginCmdsWrapper(), window); } - public subscribe(subscriber: NotificationSubscriber) { - this.isActive = true; - subscriber.subscriptionCommands.forEach( - (subscriptionCommand) => { - const cmdId = this.nextCmdId(); - this.subscribersMap.set(cmdId, subscriber); - subscriptionCommand.cmdId = cmdId; - if (subscriptionCommand instanceof UnreadCountSubCmd) { - this.cmdWrapper.unreadCountSubCmd = subscriptionCommand; - } else if (subscriptionCommand instanceof UnreadSubCmd) { - this.cmdWrapper.unreadSubCmd = subscriptionCommand; - } else if (subscriptionCommand instanceof MarkAsReadCmd) { - this.cmdWrapper.markAsReadCmd = subscriptionCommand; - this.subscribersMap.delete(cmdId); - } else if (subscriptionCommand instanceof MarkAllAsReadCmd) { - this.cmdWrapper.markAllAsReadCmd = subscriptionCommand; - this.subscribersMap.delete(cmdId); - } - } - ); - if (this.cmdWrapper.unreadCountSubCmd || this.cmdWrapper.unreadSubCmd) { - this.subscribersCount++; - } - this.publishCommands(); + public subscribe(subscriber: TelemetrySubscriber) { + this.telemetryWebsocketService.subscribe(subscriber); } - public update(subscriber: NotificationSubscriber) { - if (!this.isReconnect) { - subscriber.subscriptionCommands.forEach( - (subscriptionCommand) => { - if (subscriptionCommand.cmdId && subscriptionCommand instanceof UnreadSubCmd) { - this.cmdWrapper.unreadSubCmd = subscriptionCommand; - } - } - ); - this.publishCommands(); - } + public update(subscriber: TelemetrySubscriber) { + this.telemetryWebsocketService.update(subscriber); } - public unsubscribe(subscriber: NotificationSubscriber) { - if (this.isActive) { - subscriber.subscriptionCommands.forEach( - (subscriptionCommand) => { - if (subscriptionCommand instanceof UnreadCountSubCmd - || subscriptionCommand instanceof UnreadSubCmd) { - const unreadCountUnsubscribeCmd = new UnsubscribeCmd(); - unreadCountUnsubscribeCmd.cmdId = subscriptionCommand.cmdId; - this.cmdWrapper.unsubCmd = unreadCountUnsubscribeCmd; - } - const cmdId = subscriptionCommand.cmdId; - if (cmdId) { - this.subscribersMap.delete(cmdId); - } - } - ); - this.reconnectSubscribers.delete(subscriber); - this.subscribersCount--; - this.publishCommands(); - } + public unsubscribe(subscriber: TelemetrySubscriber) { + this.telemetryWebsocketService.unsubscribe(subscriber); } - processOnMessage(message: WebsocketNotificationMsg) { - let subscriber: NotificationSubscriber; - if (isNotificationCountUpdateMsg(message)) { - subscriber = this.subscribersMap.get(message.cmdId); - if (subscriber) { - subscriber.onNotificationCountUpdate(new NotificationCountUpdate(message)); - } - } else if (isNotificationsUpdateMsg(message)) { - subscriber = this.subscribersMap.get(message.cmdId); - if (subscriber) { - subscriber.onNotificationsUpdate(new NotificationsUpdate(message)); - } - } + processOnMessage(message: WebsocketDataMsg) { + this.telemetryWebsocketService.processOnMessage(message); } - } 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 56c35328fb..45d8510687 100644 --- a/ui-ngx/src/app/core/ws/telemetry-websocket.service.ts +++ b/ui-ngx/src/app/core/ws/telemetry-websocket.service.ts @@ -16,7 +16,8 @@ import { Inject, Injectable, NgZone } from '@angular/core'; import { - AlarmCountCmd, AlarmCountUnsubscribeCmd, + AlarmCountCmd, + AlarmCountUnsubscribeCmd, AlarmCountUpdate, AlarmDataCmd, AlarmDataUnsubscribeCmd, @@ -28,10 +29,13 @@ import { EntityDataCmd, EntityDataUnsubscribeCmd, EntityDataUpdate, - GetHistoryCmd, isAlarmCountUpdateMsg, + GetHistoryCmd, + isAlarmCountUpdateMsg, isAlarmDataUpdateMsg, isEntityCountUpdateMsg, isEntityDataUpdateMsg, + NotificationCountUpdate, + NotificationsUpdate, SubscriptionCmd, SubscriptionUpdate, TelemetryFeature, @@ -45,6 +49,16 @@ import { AppState } from '@core/core.state'; import { AuthService } from '@core/auth/auth.service'; import { WINDOW } from '@core/services/window.service'; import { WebsocketService } from '@core/ws/websocket.service'; +import { + isNotificationCountUpdateMsg, + isNotificationsUpdateMsg, + MarkAllAsReadCmd, + MarkAsReadCmd, + NotificationSubscriber, + UnreadCountSubCmd, + UnreadSubCmd, + UnsubscribeCmd +} from '@shared/models/websocket/notification-ws.models'; // @dynamic @Injectable({ @@ -84,6 +98,16 @@ export class TelemetryWebsocketService extends WebsocketService { if (subscriptionCommand.cmdId && subscriptionCommand instanceof EntityDataCmd) { this.cmdWrapper.entityDataCmds.push(subscriptionCommand); + } else if (subscriptionCommand.cmdId && subscriptionCommand instanceof UnreadSubCmd) { + this.cmdWrapper.unreadNotificationsSubCmds.push(subscriptionCommand); } } ); @@ -131,6 +157,10 @@ export class TelemetryWebsocketService extends WebsocketService implements WsServ lastCmdId = 0; subscribersCount = 0; - subscribersMap = new Map(); + subscribersMap = new Map(); - reconnectSubscribers = new Set(); + reconnectSubscribers = new Set(); - notificationUri: string; + wsUri: string; dataStream: WebSocketSubject; @@ -69,23 +69,23 @@ export abstract class WebsocketService implements WsServ if (!port) { port = '443'; } - this.notificationUri = 'wss:'; + this.wsUri = 'wss:'; } else { if (!port) { port = '80'; } - this.notificationUri = 'ws:'; + this.wsUri = 'ws:'; } - this.notificationUri += `//${this.window.location.hostname}:${port}/${apiEndpoint}`; + this.wsUri += `//${this.window.location.hostname}:${port}/${apiEndpoint}`; } - abstract subscribe(subscriber: T); + abstract subscribe(subscriber: WsSubscriber); abstract update(subscriber: T); abstract unsubscribe(subscriber: T); - abstract processOnMessage(message: any); + abstract processOnMessage(message: WebsocketDataMsg); protected nextCmdId(): number { this.lastCmdId++; @@ -158,8 +158,8 @@ export abstract class WebsocketService implements WsServ } private openSocket(token: string) { - const uri = `${this.notificationUri}?token=${token}`; - this.dataStream = webSocket( + const uri = `${this.wsUri}?token=${token}`; + this.dataStream = webSocket( { url: uri, openObserver: { @@ -176,9 +176,9 @@ export abstract class WebsocketService implements WsServ ); this.dataStream.subscribe({ - next: (message) => { + next: (message: CmdUpdateMsg) => { this.ngZone.runOutsideAngular(() => { - this.onMessage(message as WebsocketNotificationMsg); + this.onMessage(message); }); }, error: (error) => { @@ -208,11 +208,11 @@ export abstract class WebsocketService implements WsServ } } - private onMessage(message: WebsocketNotificationMsg) { + private onMessage(message: CmdUpdateMsg) { if (message.errorCode) { this.showWsError(message.errorCode, message.errorMsg); } else { - this.processOnMessage(message); + this.processOnMessage(message as WebsocketDataMsg); } this.checkToClose(); } diff --git a/ui-ngx/src/app/modules/home/components/notification/notification-bell.component.ts b/ui-ngx/src/app/modules/home/components/notification/notification-bell.component.ts index 8ee5771fe7..7c68b2df65 100644 --- a/ui-ngx/src/app/modules/home/components/notification/notification-bell.component.ts +++ b/ui-ngx/src/app/modules/home/components/notification/notification-bell.component.ts @@ -77,11 +77,11 @@ export class NotificationBellComponent implements OnDestroy { if ($event) { $event.stopPropagation(); } - this.unsubscribeSubscription(); const trigger = createVersionButton._elementRef.nativeElement; if (this.popoverService.hasPopover(trigger)) { this.popoverService.hidePopover(trigger); } else { + this.unsubscribeSubscription(); const showNotificationPopover = this.popoverService.displayPopover(trigger, this.renderer, this.viewContainerRef, ShowNotificationPopoverComponent, 'bottom', true, null, { @@ -108,5 +108,6 @@ export class NotificationBellComponent implements OnDestroy { private unsubscribeSubscription() { this.notificationCountSubscriber.unsubscribe(); this.notificationSubscriber.unsubscribe(); + this.notificationSubscriber = null; } } 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 66f39d080a..57e94baef2 100644 --- a/ui-ngx/src/app/shared/models/telemetry/telemetry.models.ts +++ b/ui-ngx/src/app/shared/models/telemetry/telemetry.models.ts @@ -24,9 +24,11 @@ import { NgZone } from '@angular/core'; import { AlarmCountQuery, AlarmData, - AlarmDataQuery, EntityCountQuery, + AlarmDataQuery, + EntityCountQuery, EntityData, - EntityDataQuery, EntityFilter, + EntityDataQuery, + EntityFilter, EntityKey, TsValue } from '@shared/models/query/query.models'; @@ -36,6 +38,16 @@ import { entityFields } from '@shared/models/entity.models'; import { isUndefined } from '@core/utils'; import { CmdWrapper, WsSubscriber } from '@shared/models/websocket/websocket.models'; import { TelemetryWebsocketService } from '@core/ws/telemetry-websocket.service'; +import { + MarkAllAsReadCmd, + MarkAsReadCmd, + NotificationCountUpdateMsg, + NotificationsUpdateMsg, + UnreadCountSubCmd, + UnreadSubCmd, + UnsubscribeCmd +} from '@shared/models/websocket/notification-ws.models'; +import { Notification } from '@shared/models/notification.models'; export const NOT_SUPPORTED = 'Not supported!'; @@ -105,7 +117,7 @@ export const timeseriesDeleteStrategyTranslations = new Map; tsSubCmds: Array; @@ -290,6 +307,11 @@ export class TelemetryPluginCmdsWrapper implements CmdWrapper { entityCountUnsubscribeCmds: Array; alarmCountCmds: Array; alarmCountUnsubscribeCmds: Array; + unreadNotificationsCountSubCmds: Array; + unreadNotificationsSubCmds: Array; + notificationsUnsubCmds: Array; + markNotificationAsReadCmds: Array; + markAllNotificationsAsReadCmds: Array; private static popCmds(cmds: Array, leftCount: number): Array { const toPublish = Math.min(cmds.length, leftCount); @@ -311,7 +333,12 @@ export class TelemetryPluginCmdsWrapper implements CmdWrapper { this.entityCountCmds.length > 0 || this.entityCountUnsubscribeCmds.length > 0 || this.alarmCountCmds.length > 0 || - this.alarmCountUnsubscribeCmds.length > 0; + this.alarmCountUnsubscribeCmds.length > 0 || + this.unreadNotificationsCountSubCmds.length > 0 || + this.unreadNotificationsSubCmds.length > 0 || + this.notificationsUnsubCmds.length > 0 || + this.markNotificationAsReadCmds.length > 0 || + this.markAllNotificationsAsReadCmds.length > 0; } public clear() { @@ -326,6 +353,11 @@ export class TelemetryPluginCmdsWrapper implements CmdWrapper { this.entityCountUnsubscribeCmds.length = 0; this.alarmCountCmds.length = 0; this.alarmCountUnsubscribeCmds.length = 0; + this.unreadNotificationsSubCmds.length = 0; + this.unreadNotificationsCountSubCmds.length = 0; + this.notificationsUnsubCmds.length = 0; + this.markNotificationAsReadCmds.length = 0; + this.markAllNotificationsAsReadCmds.length = 0; } public preparePublishCommands(maxCommands: number): TelemetryPluginCmdsWrapper { @@ -352,6 +384,16 @@ export class TelemetryPluginCmdsWrapper implements CmdWrapper { preparedWrapper.alarmCountCmds = TelemetryPluginCmdsWrapper.popCmds(this.alarmCountCmds, leftCount); leftCount -= preparedWrapper.alarmCountCmds.length; preparedWrapper.alarmCountUnsubscribeCmds = TelemetryPluginCmdsWrapper.popCmds(this.alarmCountUnsubscribeCmds, leftCount); + leftCount -= preparedWrapper.unreadNotificationsSubCmds.length; + preparedWrapper.unreadNotificationsSubCmds = TelemetryPluginCmdsWrapper.popCmds(this.unreadNotificationsSubCmds, leftCount); + leftCount -= preparedWrapper.unreadNotificationsCountSubCmds.length; + preparedWrapper.unreadNotificationsCountSubCmds = TelemetryPluginCmdsWrapper.popCmds(this.unreadNotificationsCountSubCmds, leftCount); + leftCount -= preparedWrapper.notificationsUnsubCmds.length; + preparedWrapper.notificationsUnsubCmds = TelemetryPluginCmdsWrapper.popCmds(this.notificationsUnsubCmds, leftCount); + leftCount -= preparedWrapper.markNotificationAsReadCmds.length; + preparedWrapper.markNotificationAsReadCmds = TelemetryPluginCmdsWrapper.popCmds(this.markNotificationAsReadCmds, leftCount); + leftCount -= preparedWrapper.markAllNotificationsAsReadCmds.length; + preparedWrapper.markAllNotificationsAsReadCmds = TelemetryPluginCmdsWrapper.popCmds(this.markAllNotificationsAsReadCmds, leftCount); return preparedWrapper; } } @@ -416,7 +458,7 @@ export interface AlarmCountUpdateMsg extends CmdUpdateMsg { } export type WebsocketDataMsg = AlarmDataUpdateMsg | AlarmCountUpdateMsg | - EntityDataUpdateMsg | EntityCountUpdateMsg | SubscriptionUpdateMsg; + EntityDataUpdateMsg | EntityCountUpdateMsg | SubscriptionUpdateMsg | NotificationCountUpdateMsg | NotificationsUpdateMsg; export const isEntityDataUpdateMsg = (message: WebsocketDataMsg): message is EntityDataUpdateMsg => { const updateMsg = (message as CmdUpdateMsg); @@ -627,6 +669,28 @@ export class AlarmCountUpdate extends CmdUpdate { } } +export class NotificationCountUpdate extends CmdUpdate { + totalUnreadCount: number; + + constructor(msg: NotificationCountUpdateMsg) { + super(msg); + this.totalUnreadCount = msg.totalUnreadCount; + } +} + +export class NotificationsUpdate extends CmdUpdate { + totalUnreadCount: number; + update?: Notification; + notifications?: Notification[]; + + constructor(msg: NotificationsUpdateMsg) { + super(msg); + this.totalUnreadCount = msg.totalUnreadCount; + this.update = msg.update; + this.notifications = msg.notifications; + } +} + export class TelemetrySubscriber extends WsSubscriber { private dataSubject = new ReplaySubject(1); diff --git a/ui-ngx/src/app/shared/models/websocket/notification-ws.models.ts b/ui-ngx/src/app/shared/models/websocket/notification-ws.models.ts index dffa828ae1..16280fe81a 100644 --- a/ui-ngx/src/app/shared/models/websocket/notification-ws.models.ts +++ b/ui-ngx/src/app/shared/models/websocket/notification-ws.models.ts @@ -14,36 +14,21 @@ /// limitations under the License. /// -import { BehaviorSubject, ReplaySubject } from 'rxjs'; -import { CmdUpdate, CmdUpdateMsg, CmdUpdateType, WebsocketCmd } from '@shared/models/telemetry/telemetry.models'; -import { map } from 'rxjs/operators'; +import { + CmdUpdateMsg, + CmdUpdateType, + NotificationCountUpdate, + NotificationsUpdate, + WebsocketCmd, + WebsocketDataMsg +} from '@shared/models/telemetry/telemetry.models'; import { NgZone } from '@angular/core'; import { isDefinedAndNotNull } from '@core/utils'; import { Notification } from '@shared/models/notification.models'; -import { CmdWrapper, WsSubscriber } from '@shared/models/websocket/websocket.models'; -import { NotificationWebsocketService } from '@core/ws/notification-websocket.service'; - -export class NotificationCountUpdate extends CmdUpdate { - totalUnreadCount: number; - - constructor(msg: NotificationCountUpdateMsg) { - super(msg); - this.totalUnreadCount = msg.totalUnreadCount; - } -} - -export class NotificationsUpdate extends CmdUpdate { - totalUnreadCount: number; - update?: Notification; - notifications?: Notification[]; - - constructor(msg: NotificationsUpdateMsg) { - super(msg); - this.totalUnreadCount = msg.totalUnreadCount; - this.update = msg.update; - this.notifications = msg.notifications; - } -} +import { WsService, WsSubscriber } from '@shared/models/websocket/websocket.models'; +import { BehaviorSubject, ReplaySubject } from 'rxjs'; +import { map } from 'rxjs/operators'; +import { WebsocketService } from '@core/ws/websocket.service'; export class NotificationSubscriber extends WsSubscriber { private notificationCountSubject = new ReplaySubject(1); @@ -61,40 +46,40 @@ export class NotificationSubscriber extends WsSubscriber { public notificationCount$ = this.notificationCountSubject.asObservable().pipe(map(msg => msg.totalUnreadCount)); public notifications$ = this.notificationsSubject.asObservable().pipe(map(msg => msg.notifications )); - public static createNotificationCountSubscription(notificationWsService: NotificationWebsocketService, + public static createNotificationCountSubscription(websocketService: WebsocketService, zone: NgZone): NotificationSubscriber { const subscriptionCommand = new UnreadCountSubCmd(); - const subscriber = new NotificationSubscriber(notificationWsService, zone); + const subscriber = new NotificationSubscriber(websocketService, zone); subscriber.subscriptionCommands.push(subscriptionCommand); return subscriber; } - public static createNotificationsSubscription(notificationWsService: NotificationWebsocketService, + public static createNotificationsSubscription(websocketService: WebsocketService, zone: NgZone, limit = 10): NotificationSubscriber { const subscriptionCommand = new UnreadSubCmd(limit); - const subscriber = new NotificationSubscriber(notificationWsService, zone); + const subscriber = new NotificationSubscriber(websocketService, zone); subscriber.messageLimit = limit; subscriber.subscriptionCommands.push(subscriptionCommand); return subscriber; } - public static createMarkAsReadCommand(notificationWsService: NotificationWebsocketService, + public static createMarkAsReadCommand(websocketService: WebsocketService, ids: string[]): NotificationSubscriber { const subscriptionCommand = new MarkAsReadCmd(ids); - const subscriber = new NotificationSubscriber(notificationWsService); + const subscriber = new NotificationSubscriber(websocketService); subscriber.subscriptionCommands.push(subscriptionCommand); return subscriber; } - public static createMarkAllAsReadCommand(notificationWsService: NotificationWebsocketService): NotificationSubscriber { + public static createMarkAllAsReadCommand(websocketService: WebsocketService): NotificationSubscriber { const subscriptionCommand = new MarkAllAsReadCmd(); - const subscriber = new NotificationSubscriber(notificationWsService); + const subscriber = new NotificationSubscriber(websocketService); subscriber.subscriptionCommands.push(subscriptionCommand); return subscriber; } - constructor(private notificationWsService: NotificationWebsocketService, protected zone?: NgZone) { - super(notificationWsService, zone); + constructor(private websocketService: WsService, protected zone?: NgZone) { + super(websocketService, zone); } onNotificationCountUpdate(message: NotificationCountUpdate) { @@ -183,58 +168,12 @@ export interface NotificationsUpdateMsg extends CmdUpdateMsg { notifications?: Notification[]; } -export type WebsocketNotificationMsg = NotificationCountUpdateMsg | NotificationsUpdateMsg; - -export const isNotificationCountUpdateMsg = (message: WebsocketNotificationMsg): message is NotificationCountUpdateMsg => { +export const isNotificationCountUpdateMsg = (message: WebsocketDataMsg): message is NotificationCountUpdateMsg => { const updateMsg = (message as CmdUpdateMsg); return updateMsg.cmdId !== undefined && updateMsg.cmdUpdateType === CmdUpdateType.NOTIFICATIONS_COUNT; }; -export const isNotificationsUpdateMsg = (message: WebsocketNotificationMsg): message is NotificationsUpdateMsg => { +export const isNotificationsUpdateMsg = (message: WebsocketDataMsg): message is NotificationsUpdateMsg => { const updateMsg = (message as CmdUpdateMsg); return updateMsg.cmdId !== undefined && updateMsg.cmdUpdateType === CmdUpdateType.NOTIFICATIONS; }; - -export class NotificationPluginCmdWrapper implements CmdWrapper { - - constructor() { - this.unreadCountSubCmd = null; - this.unreadSubCmd = null; - this.unsubCmd = null; - this.markAsReadCmd = null; - this.markAllAsReadCmd = null; - } - - unreadCountSubCmd: UnreadCountSubCmd; - unreadSubCmd: UnreadSubCmd; - unsubCmd: UnsubscribeCmd; - markAsReadCmd: MarkAsReadCmd; - markAllAsReadCmd: MarkAllAsReadCmd; - - public hasCommands(): boolean { - return isDefinedAndNotNull(this.unreadCountSubCmd) || - isDefinedAndNotNull(this.unreadSubCmd) || - isDefinedAndNotNull(this.unsubCmd) || - isDefinedAndNotNull(this.markAsReadCmd) || - isDefinedAndNotNull(this.markAllAsReadCmd); - } - - public clear() { - this.unreadCountSubCmd = null; - this.unreadSubCmd = null; - this.unsubCmd = null; - this.markAsReadCmd = null; - this.markAllAsReadCmd = null; - } - - public preparePublishCommands(): NotificationPluginCmdWrapper { - const preparedWrapper = new NotificationPluginCmdWrapper(); - preparedWrapper.unreadCountSubCmd = this.unreadCountSubCmd || undefined; - preparedWrapper.unreadSubCmd = this.unreadSubCmd || undefined; - preparedWrapper.unsubCmd = this.unsubCmd || undefined; - preparedWrapper.markAsReadCmd = this.markAsReadCmd || undefined; - preparedWrapper.markAllAsReadCmd = this.markAllAsReadCmd || undefined; - this.clear(); - return preparedWrapper; - } -}