import { Injectable } from '@angular/core'; import { Router, ActivatedRoute } from '@angular/router'; import { TranslateService } from '@ngx-translate/core'; import { io } from "socket.io-client"; // RxJs import { Observable, Subject, fromEvent, of, Subscription } from 'rxjs'; import { mergeMap, map, catchError, take, first } from 'rxjs/operators'; import { ConfigService } from '@arrow/bom/config'; import { UserService } from '../user/user.service'; import { DialogService } from '../dialog/dialog.service'; import { AuthService } from '../auth.service'; import { AuthTokenService } from '../auth-token.service'; import { Utils } from '../utils.service'; import { ApiSocketListenAction } from './api.socket.listen.action.enum'; import { ApiSocketEmitAction } from './api.socket.emit.action.enum'; import { SocketEvent } from './socket.event.enum'; const REDIRECT_PARAMETER = 'url'; @Injectable() export class SocketsService { private _isConnecting = true; private _token: string | null = null; private _socket: any; private _clientId: string; private _socket$$: Subject = new Subject(); private _apiError = new Subject(); private _socketError = new Subject(); private _actionsMap = new Map(); private _listenActions: Array; private _emitActions: Array; private _actionSubscriptions: Subscription[] = []; constructor( private _config: ConfigService, private _userService: UserService, private _dialog: DialogService, private _router: Router, private _route: ActivatedRoute, private _authService: AuthService, private _authTokenService: AuthTokenService, private _translate: TranslateService, private _utilsService: Utils ) { console.log('Instantiating SocketsService...'); this._initialize(); } /* =================== = Private Methods = =================== */ private _initialize(): void { this._listenActions = this._utilsService.getAllEnumValuesAsList( ApiSocketListenAction ); this._emitActions = this._utilsService.getAllEnumValuesAsList( ApiSocketEmitAction ); this._mapActions(this._listenActions, this._emitActions); this._userService .isReady() .pipe(first()) // Connects after we have the session cookie .subscribe(() => this._getFreshToken().subscribe((token: string | null) => { this._connect(token); }) ); } private _connect(token?: string) { this._token = token; this._isConnecting = true; const query: any = token ? { jwt: token } : {}; const myArrowAppClientId = this._authTokenService.getClientId(); const refreshToken = this._authTokenService.getValidRefreshToken(); if (myArrowAppClientId) query.clientid = myArrowAppClientId; if (refreshToken) query.refreshtoken = refreshToken; this._socket = io(this._config.apiUrlWithoutVersion, { reconnection: true, upgrade: false, transports: ['websocket'], reconnectionAttempts: 10, query }); this._socket.on(SocketEvent.CONNECT_FAILED, this._handleConnectionError); this._socket.on(SocketEvent.RECONNECT_FAILED, this._handleConnectionError); this._socket.on(SocketEvent.CONNECT, () => { const existingClientId = localStorage.getItem( this._config.clientIdLocalStorageKey ); const registerPayload = existingClientId ? { clientId: existingClientId } : {}; this._socket.emit('register', registerPayload); }); this._socket.on(SocketEvent.CLIENT_ID, (clientId: string) => { localStorage.setItem(this._config.clientIdLocalStorageKey, clientId); this._clientId = clientId; this._actionSubscriptions = this._addActionListeners(this._listenActions); this._listenSocketErrors(); this._listenApiErrors(); this._socket$$.next(this._socket); this._isConnecting = false; }); } /** * disconect sockets * */ private _disconnect(): Observable { if (this._socket) { this._socket.removeAllListeners(); this._actionSubscriptions.forEach((sub: Subscription) => sub.unsubscribe() ); this._actionSubscriptions = []; if (!this._socket.connected) return of(null); this._socket.disconnect(); return fromEvent(this._socket, SocketEvent.DISCONNECT); } throw new Error('Error while trying to disconnect sockets'); } /** * map actions * @param listen * @param emmit * */ private _mapActions(listen: string[], emmit: string[]) { const allActions = [...listen, ...emmit]; allActions.forEach((action: string) => { this._actionsMap.set(action, new Subject()); }); } /** * add actions listener * @param actionsList * */ private _addActionListeners(actionsList: string[]): Subscription[] { const subs: Subscription[] = []; actionsList.forEach(action => { const sub = fromEvent(this._socket, action).subscribe(data => this._actionsMap.get(action).next(data) ); subs.push(sub); }); return subs; } /** * display modal alert if an error occurs * */ private _handleConnectionError() { this._dialog .alert( this._translate.instant('An-error-has-occured'), this._translate.instant('An-error-has-occured') ) .subscribe(() => { this._router.navigate(['../'], { relativeTo: this._route }); }); } /** * handle 401 erros * @param actionName * @param message */ private _handle401Errors(actionName: string, message: string) { // if (!this._router) this._router = this.injector.get(_router); const data: any = {}; let actionType = ''; data.message = message; // TODO: This needs improvements since calls are not standard, some delete calls are prefixed remove while others delete, same with update and alter if (actionName.match(/^add/)) { actionType = 'C'; // Create } else if (actionName.match(/^get/)) { actionType = 'R'; // Read } else if (actionName.match(/^update/)) { actionType = 'U'; // Update } else if (actionName.match(/^remove/)) { actionType = 'D'; // Delete } let url = encodeURIComponent(location.href); switch (actionType) { case 'R': // Read if (this._userService.isLoggedIn()) { this._dialog.alert(data.title, data.message).subscribe(() => { this._router.navigate(['../'], { relativeTo: this._route }); }); } else { this._dialog.alert(data.title, data.message).subscribe(() => { url = encodeURIComponent(location.href); location.replace(`${this._config.getRoute('login')}?url=${url}`); }); } break; default: // Create, Update, Delete data.title = 'Please login'; this._dialog.alert(data.title, data.message).subscribe(() => { url = encodeURIComponent(location.href); location.replace(`/${this._config.getRoute('login')}?url=${url}`); }); break; } } private _handle500Errors(actionName: string, message: string) { const data: any = {}; data.message = message; data.title = 'Server Error'; this._dialog.alert(data.title, data.message); } private _handleApiError(err: any) { const { status, action, statusMessage } = err; switch (status) { case 401: this._handle401Errors(action, statusMessage); break; case 500: this._handle500Errors(action, statusMessage); break; default: this._dialog .alert( this._translate.instant('An-error-has-occured'), this._translate.instant('An-error-has-occured') ) .subscribe(() => { this._router.navigate(['../'], { relativeTo: this._route }); }); break; } } private _listenApiErrors() { fromEvent(this._socket, SocketEvent.API_ERROR).subscribe((err: any) => { const { action } = err; const subject: Subject = this._actionsMap.get(action); if (subject) subject.error(err); this._handleApiError(err); this._apiError.next(err); }); } /** * socket errors */ private _listenSocketErrors() { fromEvent(this._socket, SocketEvent.SOCKET_ERROR).subscribe((err: any) => { const { action } = err; const subject: Subject = this._actionsMap.get(action); if (subject) subject.error(err); this._socketError.next(err); }); } private _reset(token?: string): void { this._disconnect(); this._connect(token); } /** * Return the current token if it is fresh, attempt to refresh the token and return result otherwise */ private _getFreshToken(): Observable { const needsToken = this._config.getTarget() === 'myarrow'; const freshToken$ = needsToken ? this._authService.checkAndRenewTokens().pipe( map((tokens: any) => tokens.access_token), catchError((error: Error) => { throw error; }) ) : of(null); return freshToken$; } /** * If token is fresh returns current socket, reconnect and return new socket otherwise */ private _getSocket(): Observable { if (this._isConnecting) return this._socket$$.asObservable(); this._isConnecting = true; return this._getFreshToken().pipe( mergeMap((freshToken: string) => { if (freshToken === this._token) { this._isConnecting = false; return of(this._socket); } this._reset(freshToken); return this._socket$$.asObservable(); }) ); } /* ================== = Public Methods = ================== */ /** * emit socket action to be listened * @param action * @param params * @return observable */ public emitAction(action: string, params: any = {}): Observable { this._getSocket() .pipe(take(1)) .subscribe((socket: any) => { if (!socket || socket.disconnected) throw new Error( `WS are not connected, impossible to emit action ${action}` ); socket.emit('action', { clientId: this._clientId, action, params }); }); return this._actionsMap.get(action).asObservable(); } /** * observe socket action * @param action */ public observeAction(action: string): Observable { const subject: Subject = this._actionsMap.get(action); if (!subject) throw new Error(`Action: ${action} is not recognized`); return subject.asObservable(); } public observeApiErrors(): Observable { return this._apiError.asObservable(); } }