import { Injectable } from '@angular/core'; import { Router, ActivatedRoute } from '@angular/router'; import { TranslateService } from '@ngx-translate/core'; // 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'; // Next two lines are for .Net Core - Server // import * as signalR from '@microsoft/signalr'; // import { IHttpConnectionOptions } from '@microsoft/signalr'; declare var $: any; @Injectable() export class SignalRSocketsService { private isConnecting = true; private token: string | null = null; //private _socket: any; private hubConnection: any; private clientId: string; private hubProxy: any; private hubProxy$$: 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 SignalRSocketsService...'); this._initialize(); } 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); }) ); // this.testLatestSignalR() } private _connect(token?: string) { this.token = token; this.isConnecting = true; const myArrowAppClientId = this._authTokenService.getClientId(); const refreshToken = this._authTokenService.getValidRefreshToken(); const queryString: any = token ? { jwt: token } : {}; if (myArrowAppClientId) queryString.clientid = myArrowAppClientId; if (refreshToken) queryString.refreshtoken = refreshToken; const connOptions = { logging: false, // TODO: control logging verbosity in consol useDefaultPath: false } console.log(`SignalR URL: ${this._config.signalRUrl}`); this.hubConnection = $.hubConnection(this._config.signalRUrl, connOptions); this.hubConnection.qs = queryString; this.registerConnectionLifecycleHooks(); this.hubProxy = this.hubConnection.createHubProxy('bomHub'); this.registerSocketHubHooks(); this.hubConnection.start({ transport: 'webSockets', withCredentials: true, // skipNegotiation: true }) .then((connection) => { console.log(`SignalR Connection started: ${connection.id}`); localStorage.setItem(this._config.clientIdLocalStorageKey, connection.id) }) .then(() => this.registerClientId()) .catch((err) => console.log('Error while starting SignalR connection: ' + err)) } private registerClientId() { const connectionId = localStorage.getItem(this._config.clientIdLocalStorageKey) || ""; this.hubProxy.invoke('register', connectionId) .done((clientId) => { localStorage.setItem(this._config.clientIdLocalStorageKey, clientId); this.clientId = clientId; this.actionSubscriptions = this._addActionListeners(this.listenActions); this._listenSocketErrors(); this._listenApiErrors(); this.hubProxy$$.next(this.hubProxy); this.isConnecting = false; }) .fail((error) => { console.log(`Invocation of 'Register' server function failed. Error: ${error}`); }); } private registerSocketHubHooks() { this.hubProxy.on(SocketEvent.CONNECT_FAILED, this._handleConnectionError); this.hubProxy.on(SocketEvent.RECONNECT_FAILED, this._handleConnectionError); // TODO: delete as redundant. Now 'Register' sends connectionId back // this.hubProxy.on(SocketEvent.CLIENT_ID, (clientId: string) => { // console.log(`SignalR: onClientId handler received ${clientId}`) // localStorage.setItem(this._config.clientIdLocalStorageKey, clientId); // this.clientId = clientId; // this.actionSubscriptions = this._addActionListeners(this.listenActions); // this._listenSocketErrors(); // this._listenApiErrors(); // this.hubProxy$$.next(this.hubProxy); // this.isConnecting = false; // }); this.hubProxy.on('starting', () => { // Raised before any data is sent over the connection. console.log(`SignalR (Hub) connection starting...`) }); this.hubProxy.on('received', (data) => { // Raised when any data is received on the connection. Provides the received data. console.log(`SignalR (Hub) connection received: ${JSON.stringify(data)}`) }); this.hubProxy.on('connectionSlow', () => { // Raised when the client detects a slow or frequently dropping connection. console.log(`SignalR (Hub): We are currently experiencing difficulties with the connection.`) }); this.hubProxy.on('reconnecting', (data) => { // Raised when the underlying transport begins reconnecting. console.log(`SignalR (Hub) is reconnecting: ${JSON.stringify(data)}`) }); this.hubProxy.on('reconnected', (data) => { // Raised when the underlying transport has reconnected. console.log(`SignalR (Hub) is reconnected: ${JSON.stringify(data)}`) }); this.hubProxy.on('stateChanged', (data) => { // Raised when the connection state changes. Provides the old state and the new state (Connecting, Connected, Reconnecting, or Disconnected). console.log(`SignalR (Hub) state changed: ${JSON.stringify(data)}`) }); this.hubProxy.on('disconnected', (data) => { // Raised when the connection has disconnected. console.log(`SignalR (Hub): disconnected: ${data}`); }); this.hubProxy.on('error', (e) => { console.log('SignalR (Hub) Client-Socket Error: ' + e) }); this.hubProxy.on('onclose', (e) => { console.log(`SignalR (Hub): disconnected: ${e}`); }); this.hubProxy.on('close', (e) => { console.log(`SignalR (Hub): close: ${e}`); }); } private registerConnectionLifecycleHooks() { this.hubConnection.starting(() => { // TODO: tested + // Raised before any data is sent over the connection. console.log(`SignalR connection starting...`) }); this.hubConnection.received(data => { // TODO: tested+ | TRACKS EVERY RECEIVED MESSAGE // Raised when any data is received on the connection. Provides the received data. //Could be a JSON or plane text message // TODO: enable as necessary for debugging - too verbose! // console.log(`SignalR connection message received. ${JSON.stringify(data)}}`) }); this.hubConnection.connectionSlow(() => { // TODO: tested+ // Raised when the client detects a slow or frequently dropping connection. console.log(`SignalR: We are currently experiencing difficulties with the connection.`) }); this.hubConnection.reconnecting(() => { // TODO: tested+ // Raised when the underlying transport begins reconnecting. console.log(`SignalR is reconnecting...`) }); this.hubConnection.reconnected(() => { // TODO: tested+ // Raised when the underlying transport has reconnected. console.log(`SignalR is reconnected.`) }); this.hubConnection.stateChanged(data => { // Raised when the connection state changes. Provides the old state and the new state (Connecting, Connected, Reconnecting, or Disconnected). // TODO: tested+ if (data.newState === 1) console.log(`SignalR state changed: connecting...`) else if (data.newState === 4) console.log(`SignalR state changed: disconnecting...`) // TODO: enable as necessary - too verbose // else // console.log(`SignalR state changed: ${JSON.stringify(data)}`) }); this.hubConnection.disconnected(() => { // TODO: tested + // Raised when the connection has disconnected. console.log(`SignalR: is disconnected`); }); // TODO: tested + this.hubConnection.error = (error) => { console.log('SignalR connection error occurred: ' + error) }; } /** * disconnect sockets * */ private _disconnect(): Observable { if (this.hubProxy && this.hubConnection) { //this._hubProxy.removeAllListeners(); this.actionSubscriptions.forEach((sub: Subscription) => sub.unsubscribe() ); this.actionSubscriptions = []; //if (!this._hubProxy.connected) return of(null); //this._hubProxy.disconnect(); // TODO: tested+ this.hubConnection.stop(); return fromEvent(this.hubProxy, 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.hubProxy, 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; } } // TODO: tested + private _listenApiErrors() { fromEvent(this.hubProxy, SocketEvent.API_ERROR).subscribe((err: any) => { console.log('SignalR: API_ERROR error handler in action') const { action } = err; const subject: Subject = this.actionsMap.get(action); if (subject) subject.error(err); this._handleApiError(err); this.apiError.next(err); }); } /** * socket errors */ // TODO: test private _listenSocketErrors() { fromEvent(this.hubProxy, SocketEvent.SOCKET_ERROR).subscribe((err: any) => { console.log('SignalR: SOCKET_ERROR error handler in action') 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.hubProxy$$.asObservable(); this.isConnecting = true; return this._getFreshToken().pipe( mergeMap((freshToken: string) => { if (freshToken === this.token) { this.isConnecting = false; return of(this.hubProxy); } this._reset(freshToken); return this.hubProxy$$.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((hubProxy: any) => { if (!hubProxy || hubProxy.disconnected) throw new Error( `WS are not connected, impossible to emit action ${action}` ); hubProxy.invoke('action', { clientId: this.clientId, action, params }) .done(function () { }) .fail(function (error) { console.log('Invocation of server function - failed. Error: ' + error); }); }); 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(); } // private testLatestSignalR() { // // TODO: @microsoft/signalr after migration to .Net Core+ // const connection = new signalR.HubConnectionBuilder() // .configureLogging(signalR.LogLevel.Debug) // .withUrl(this.config.signalRUrl, { // withCredentials: true, // skipNegotiation: true, // transport: signalR.HttpTransportType.WebSockets, // //accessTokenFactory: this.getToken, // //headers: ... // } as IHttpConnectionOptions) // .build(); // connection.start().then(function () { // console.log(`SignalR connected: ${connection.state}`); // }).catch(function (err) { // console.log('SignalR error'); // return console.error(err.toString()); // }); // connection.on("messageToClients", (msg) => { // console.log(`SignalR - messageToClients sent: ${msg}`); // }); // // connection.invoke("sendMessageToServer", "ping").catch(function (err) { // // return console.error(err.toString()); // // }); // //connection.send("sendMessageToServer", "ping"); // // connection.stop().then(function () { // // console.log('SignalR Disconnected!'); // // }).catch(function (err) { // // return console.error(err.toString()); // // }); // } // // for reference, from node.js socket-serve project // private getToken() { // const xhr = new XMLHttpRequest(); // return new Promise ((resolve, reject) => { // xhr.onreadystatechange = function() { // if (this.readyState !== 4) return; // if (this.status == 200) { // resolve(this.responseText); // } else { // reject(this.statusText); // } // }; // xhr.open("GET", "/api/token"); // xhr.send(); // }); // } }