Browse Source

UI: Change ws subscription to single session

pull/9717/head
Vladyslav_Prykhodko 3 years ago
parent
commit
31f95fbe62
  1. 108
      ui-ngx/src/app/core/ws/notification-websocket.service.ts
  2. 75
      ui-ngx/src/app/core/ws/telemetry-websocket.service.ts
  3. 32
      ui-ngx/src/app/core/ws/websocket.service.ts
  4. 3
      ui-ngx/src/app/modules/home/components/notification/notification-bell.component.ts
  5. 74
      ui-ngx/src/app/shared/models/telemetry/telemetry.models.ts
  6. 109
      ui-ngx/src/app/shared/models/websocket/notification-ws.models.ts

108
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<NotificationSubscriber> {
export class NotificationWebsocketService extends WebsocketService<TelemetrySubscriber> {
cmdWrapper: NotificationPluginCmdWrapper;
constructor(protected store: Store<AppState>,
constructor(private telemetryWebsocketService: TelemetryWebsocketService,
protected store: Store<AppState>,
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);
}
}

75
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<TelemetrySubscri
this.cmdWrapper.entityCountCmds.push(subscriptionCommand);
} else if (subscriptionCommand instanceof AlarmCountCmd) {
this.cmdWrapper.alarmCountCmds.push(subscriptionCommand);
} else if (subscriptionCommand instanceof UnreadCountSubCmd) {
this.cmdWrapper.unreadNotificationsCountSubCmds.push(subscriptionCommand);
} else if (subscriptionCommand instanceof UnreadSubCmd) {
this.cmdWrapper.unreadNotificationsSubCmds.push(subscriptionCommand);
} else if (subscriptionCommand instanceof MarkAsReadCmd) {
this.cmdWrapper.markNotificationAsReadCmds.push(subscriptionCommand);
this.subscribersMap.delete(cmdId);
} else if (subscriptionCommand instanceof MarkAllAsReadCmd) {
this.cmdWrapper.markAllNotificationsAsReadCmds.push(subscriptionCommand);
this.subscribersMap.delete(cmdId);
}
}
);
@ -97,6 +121,8 @@ export class TelemetryWebsocketService extends WebsocketService<TelemetrySubscri
(subscriptionCommand) => {
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<TelemetrySubscri
const alarmCountUnsubscribeCmd = new AlarmCountUnsubscribeCmd();
alarmCountUnsubscribeCmd.cmdId = subscriptionCommand.cmdId;
this.cmdWrapper.alarmCountUnsubscribeCmds.push(alarmCountUnsubscribeCmd);
} else if (subscriptionCommand instanceof UnreadCountSubCmd || subscriptionCommand instanceof UnreadSubCmd) {
const notificationsUnsubCmds = new UnsubscribeCmd();
notificationsUnsubCmds.cmdId = subscriptionCommand.cmdId;
this.cmdWrapper.notificationsUnsubCmds.push(notificationsUnsubCmds);
}
const cmdId = subscriptionCommand.cmdId;
if (cmdId) {
@ -145,29 +175,28 @@ export class TelemetryWebsocketService extends WebsocketService<TelemetrySubscri
}
processOnMessage(message: WebsocketDataMsg) {
let subscriber: TelemetrySubscriber;
if (isEntityDataUpdateMsg(message)) {
subscriber = this.subscribersMap.get(message.cmdId);
if (subscriber) {
subscriber.onEntityData(new EntityDataUpdate(message));
}
} else if (isAlarmDataUpdateMsg(message)) {
subscriber = this.subscribersMap.get(message.cmdId);
if (subscriber) {
subscriber.onAlarmData(new AlarmDataUpdate(message));
}
} else if (isEntityCountUpdateMsg(message)) {
let subscriber: TelemetrySubscriber | NotificationSubscriber;
if ('cmdId' in message && message.cmdId) {
subscriber = this.subscribersMap.get(message.cmdId);
if (subscriber) {
subscriber.onEntityCount(new EntityCountUpdate(message));
}
} else if (isAlarmCountUpdateMsg(message)) {
subscriber = this.subscribersMap.get(message.cmdId);
if (subscriber) {
subscriber.onAlarmCount(new AlarmCountUpdate(message));
if (subscriber instanceof NotificationSubscriber) {
if (isNotificationCountUpdateMsg(message) && subscriber) {
subscriber.onNotificationCountUpdate(new NotificationCountUpdate(message));
} else if (isNotificationsUpdateMsg(message) && subscriber) {
subscriber.onNotificationsUpdate(new NotificationsUpdate(message));
}
} else if (subscriber instanceof TelemetrySubscriber) {
if (isEntityDataUpdateMsg(message) && subscriber) {
subscriber.onEntityData(new EntityDataUpdate(message));
} else if (isAlarmDataUpdateMsg(message) && subscriber) {
subscriber.onAlarmData(new AlarmDataUpdate(message));
} else if (isEntityCountUpdateMsg(message) && subscriber) {
subscriber.onEntityCount(new EntityCountUpdate(message));
} else if (isAlarmCountUpdateMsg(message) && subscriber) {
subscriber.onAlarmCount(new AlarmCountUpdate(message));
}
}
} else if (message.subscriptionId) {
subscriber = this.subscribersMap.get(message.subscriptionId);
} else if ('subscriptionId' in message && message.subscriptionId) {
subscriber = this.subscribersMap.get(message.subscriptionId) as TelemetrySubscriber;
if (subscriber) {
subscriber.onData(new SubscriptionUpdate(message));
}

32
ui-ngx/src/app/core/ws/websocket.service.ts

@ -21,10 +21,10 @@ import { AuthService } from '@core/auth/auth.service';
import { NgZone } from '@angular/core';
import { selectIsAuthenticated } from '@core/auth/auth.selectors';
import { webSocket, WebSocketSubject } from 'rxjs/webSocket';
import { WebsocketNotificationMsg } from '@shared/models/websocket/notification-ws.models';
import { CmdUpdateMsg } from '@shared/models/telemetry/telemetry.models';
import { CmdUpdateMsg, TelemetrySubscriber, WebsocketDataMsg } from '@shared/models/telemetry/telemetry.models';
import { ActionNotificationShow } from '@core/notification/notification.actions';
import Timeout = NodeJS.Timeout;
import { NotificationSubscriber } from '@shared/models/websocket/notification-ws.models';
const RECONNECT_INTERVAL = 2000;
const WS_IDLE_TIMEOUT = 90000;
@ -42,11 +42,11 @@ export abstract class WebsocketService<T extends WsSubscriber> implements WsServ
lastCmdId = 0;
subscribersCount = 0;
subscribersMap = new Map<number, T>();
subscribersMap = new Map<number, TelemetrySubscriber | NotificationSubscriber>();
reconnectSubscribers = new Set<T>();
reconnectSubscribers = new Set<WsSubscriber>();
notificationUri: string;
wsUri: string;
dataStream: WebSocketSubject<CmdWrapper | CmdUpdateMsg>;
@ -69,23 +69,23 @@ export abstract class WebsocketService<T extends WsSubscriber> 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<T extends WsSubscriber> implements WsServ
}
private openSocket(token: string) {
const uri = `${this.notificationUri}?token=${token}`;
this.dataStream = webSocket(
const uri = `${this.wsUri}?token=${token}`;
this.dataStream = webSocket<CmdUpdateMsg>(
{
url: uri,
openObserver: {
@ -176,9 +176,9 @@ export abstract class WebsocketService<T extends WsSubscriber> 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<T extends WsSubscriber> 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();
}

3
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;
}
}

74
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<TimeseriesDeleteStra
[TimeseriesDeleteStrategy.DELETE_LATEST_VALUE, 'attribute.delete-timeseries.latest-value'],
[TimeseriesDeleteStrategy.DELETE_ALL_DATA_FOR_TIME_PERIOD, 'attribute.delete-timeseries.all-data-for-time-period']
]
)
);
export interface AttributeData {
lastUpdateTs?: number;
@ -278,6 +290,11 @@ export class TelemetryPluginCmdsWrapper implements CmdWrapper {
this.entityCountUnsubscribeCmds = [];
this.alarmCountCmds = [];
this.alarmCountUnsubscribeCmds = [];
this.unreadNotificationsCountSubCmds = [];
this.unreadNotificationsSubCmds = [];
this.notificationsUnsubCmds = [];
this.markNotificationAsReadCmds = [];
this.markAllNotificationsAsReadCmds = [];
}
attrSubCmds: Array<AttributesSubscriptionCmd>;
tsSubCmds: Array<TimeseriesSubscriptionCmd>;
@ -290,6 +307,11 @@ export class TelemetryPluginCmdsWrapper implements CmdWrapper {
entityCountUnsubscribeCmds: Array<EntityCountUnsubscribeCmd>;
alarmCountCmds: Array<AlarmCountCmd>;
alarmCountUnsubscribeCmds: Array<AlarmCountUnsubscribeCmd>;
unreadNotificationsCountSubCmds: Array<UnreadCountSubCmd>;
unreadNotificationsSubCmds: Array<UnreadSubCmd>;
notificationsUnsubCmds: Array<UnsubscribeCmd>;
markNotificationAsReadCmds: Array<MarkAsReadCmd>;
markAllNotificationsAsReadCmds: Array<MarkAllAsReadCmd>;
private static popCmds<T>(cmds: Array<T>, leftCount: number): Array<T> {
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<SubscriptionUpdate>(1);

109
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<NotificationCountUpdate>(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<WsSubscriber>,
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<WsSubscriber>,
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<WsSubscriber>,
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<WsSubscriber>): 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<any>, 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;
}
}

Loading…
Cancel
Save