### FILE: C:\Users\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xFrameworkPushServiceHolder\projects\x-framework-push-service\src\public-api.ts /* * Public API Surface of x-framework-push-service */ // export * from './lib/tokens/x-injectable-tokens'; export * from './lib/models/x-hub-connection.dto'; export * from './lib/typings/x-push-notification.typings'; export * from './lib/policies/x-push-service-token-fetcher'; export * from './lib/config/x-framework-push-service.config'; export * from './lib/policies/x-push-service-connection-retry-policy'; // // Base ... export * from './lib/base/x-rtc-peer.service'; export * from './lib/base/x-base-push.service'; export * from './lib/base/x-base-web-rtc.service'; export * from './lib/base/x-base-entity.push.service'; // // Interfaces ... // // Models ... // // Tools ... // // Module ... export * from './lib/x-framework-push-service.module'; ### FILE: C:\Users\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xFrameworkPushServiceHolder\projects\x-framework-push-service\src\lib\x-framework-push-service.module.ts import { NgModule } from '@angular/core'; @NgModule({ declarations: [ ], imports: [ ], exports: [ ] }) export class XFrameworkPushServiceModule { } ### FILE: C:\Users\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xFrameworkPushServiceHolder\projects\x-framework-push-service\src\lib\base\x-base-entity.push.service.ts import { Subject } from 'rxjs'; import { Injectable } from '@angular/core'; import { XBasePushService } from './x-base-push.service'; import { hasChild, parseJson, isNullOrUndefined } from 'x-framework-core'; import { XBaseEntityPushAction } from '../typings/x-push-notification.typings'; @Injectable() export abstract class XBaseEntityPushService extends XBasePushService { // //#region Props ... onAdded$ = new Subject<{ model: TEntity; actorId: string }>(); onUpdated$ = new Subject<{ model: TEntity; actorId: string }>(); onDeleted$ = new Subject<{ model: TEntity; actorId: string }>(); onAddOrUpdated$ = new Subject<{ model: TEntity; actorId: string }>(); onManyAdded$ = new Subject<{ model: Array; actorId: string }>(); onManyUpdated$ = new Subject<{ model: Array; actorId: string }>(); onManyDeleted$ = new Subject<{ model: Array; actorId: string }>(); //#endregion // //#region Actions ... /** * Fire Add Action ... * * @param model * @returns */ async add(model: TEntity) { // // Validate ... let isValid = !isNullOrUndefined(model); if (!isValid) { return; } // await this.invoke(XBaseEntityPushAction.Add, JSON.stringify(model)); } /** * Fire Update Action ... * * @param model * @returns */ async update(model: TEntity) { // // Validate ... let isValid = !isNullOrUndefined(model); if (!isValid) { return; } // await this.invoke(XBaseEntityPushAction.Update, JSON.stringify(model)); } /** * Fire Delete Action ... * * @param model * @returns */ async delete(model: TEntity) { // // Validate ... let isValid = !isNullOrUndefined(model); if (!isValid) { return; } // await this.invoke(XBaseEntityPushAction.Delete, JSON.stringify(model)); } /** * Fire Add Or Update Action ... * * @param model * @returns */ async addOrUpdate(model: TEntity) { // // Validate ... let isValid = !isNullOrUndefined(model); if (!isValid) { return; } // await this.invoke(XBaseEntityPushAction.AddOrUpdate, JSON.stringify(model)); } /** * Fire Add Many Action ... * * @param model * @returns */ async addMany(models: Array) { // // Validate ... let isValid = !isNullOrUndefined(models) && hasChild(models); if (!isValid) { return; } // await this.invoke(XBaseEntityPushAction.AddMany, JSON.stringify(models)); } /** * Fire Update Many Action ... * * @param model * @returns */ async updateMany(models: Array) { // // Validate ... let isValid = !isNullOrUndefined(models) && hasChild(models); if (!isValid) { return; } // await this.invoke(XBaseEntityPushAction.UpdateMany, JSON.stringify(models)); } /** * Fire Delete Many Action ... * * @param model * @returns */ async deleteMany(models: Array) { // // Validate ... let isValid = !isNullOrUndefined(models) && hasChild(models); if (!isValid) { return; } // await this.invoke(XBaseEntityPushAction.DeleteMany, JSON.stringify(models)); } //#endregion // registerSubjects(): void { super.registerSubjects(); // this.onAdded((model: TEntity, actorId: string) => { this.onAdded$.next({ model, actorId }); }); // this.onUpdated((model: TEntity, actorId: string) => { this.onUpdated$.next({ model, actorId }); }); // this.onDeleted((model: TEntity, actorId: string) => { this.onDeleted$.next({ model, actorId }); }); // this.onAddOrUpdated((model: TEntity, actorId: string) => { this.onAddOrUpdated$.next({ model, actorId }); }); // this.onManyAdded((model: Array, actorId: string) => { this.onManyAdded$.next({ model, actorId }); }); // this.onManyUpdated((model: Array, actorId: string) => { this.onManyUpdated$.next({ model, actorId }); }); // this.onManyDeleted((model: Array, actorId: string) => { this.onManyDeleted$.next({ model, actorId }); }); } // //#region Event Oners ... /** * Register Add Event Listener ... * * @param callback */ public onAdded(callback: (model: TEntity, actorId: string) => any) { // this.on(XBaseEntityPushAction.Add, (modelJson: string, actorId: string) => { // const model = parseJson(modelJson); callback(model, actorId); }); } /** * Register Update Event Listener ... * * @param callback */ public onUpdated(callback: (model: TEntity, actorId: string) => any) { // this.on( XBaseEntityPushAction.Update, (modelJson: string, actorId: string) => { // const model = parseJson(modelJson); callback(model, actorId); } ); } /** * Register Delete Event Listener ... * * @param callback */ public onDeleted(callback: (model: TEntity, actorId: string) => any) { // this.on( XBaseEntityPushAction.Delete, (modelJson: string, actorId: string) => { // const model = parseJson(modelJson); callback(model, actorId); } ); } /** * Register Add Or Update Event Listener ... * * @param callback */ public onAddOrUpdated(callback: (model: TEntity, actorId: string) => any) { // this.on( XBaseEntityPushAction.AddOrUpdate, (modelJson: string, actorId: string) => { // const model = parseJson(modelJson); callback(model, actorId); } ); } /** * Register Many Added Event Listener ... * * @param callback */ public onManyAdded( callback: (model: Array, actorId: string) => any ) { // this.on( XBaseEntityPushAction.AddMany, (modelJson: string, actorId: string) => { // const model = parseJson>(modelJson); callback(model, actorId); } ); } /** * Register Many Updated Event Listener ... * * @param callback */ public onManyUpdated( callback: (model: Array, actorId: string) => any ) { // this.on( XBaseEntityPushAction.UpdateMany, (modelJson: string, actorId: string) => { // const model = parseJson>(modelJson); callback(model, actorId); } ); } /** * Register Many Deleted Event Listener ... * * @param callback */ public onManyDeleted( callback: (model: Array, actorId: string) => any ) { // this.on( XBaseEntityPushAction.DeleteMany, (modelJson: string, actorId: string) => { // const model = parseJson>(modelJson); callback(model, actorId); } ); } //#endregion // //#region Event Offers ... /** * UnRegister Add Event Listener ... * * @param callback */ public offAdded(callback?: (model: TEntity, actorId: string) => any) { this.off(XBaseEntityPushAction.Add, callback); } /** * UnRegister Update Event Listener ... * * @param callback */ public offUpdated(callback?: (model: TEntity, actorId: string) => any) { this.off(XBaseEntityPushAction.Update, callback); } /** * UnRegister Delete Event Listener ... * * @param callback */ public offDeleted(callback?: (model: TEntity, actorId: string) => any) { this.off(XBaseEntityPushAction.Delete, callback); } /** * UnRegister Add Or Update Event Listener ... * * @param callback */ public offAddOrUpdated(callback?: (model: TEntity, actorId: string) => any) { this.off(XBaseEntityPushAction.AddOrUpdate, callback); } /** * UnRegister Many Added Event Listener ... * * @param callback */ public offManyAdded( callback?: (model: Array, actorId: string) => any ) { this.off(XBaseEntityPushAction.AddMany, callback); } /** * UnRegister Many Updated Event Listener ... * * @param callback */ public offManyUpdated( callback?: (model: Array, actorId: string) => any ) { this.off(XBaseEntityPushAction.UpdateMany, callback); } /** * UnRegister Many Deleted Event Listener ... * * @param callback */ public offManyDeleted( callback?: (model: Array, actorId: string) => any ) { this.off(XBaseEntityPushAction.DeleteMany, callback); } //#endregion } ### FILE: C:\Users\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xFrameworkPushServiceHolder\projects\x-framework-push-service\src\lib\base\x-base-push.service.ts import { isFunction, XBaseService, XExceptionIDs, throwException, isNullOrUndefined, isNullOrEmptyString, } from 'x-framework-core'; import { ITransport, HubConnection, MessageHeaders, HttpTransportType, HubConnectionState, HubConnectionBuilder, IHttpConnectionOptions, } from '@microsoft/signalr'; import { Subject } from 'rxjs'; import { X_FRAMEWORK_PUSH_SERVICE_CONFIG, X_FRAMEWORK_PUSH_SERVICE_TOKEN_FETCHER, } from '../tokens/x-injectable-tokens'; import { Inject, Injectable, Optional } from '@angular/core'; import { XBasePushAction } from '../typings/x-push-notification.typings'; import { MessagePackHubProtocol } from '@microsoft/signalr-protocol-msgpack'; import { IXPushServiceTokenFetcher } from '../policies/x-push-service-token-fetcher'; import { XFrameworkPushServiceConfig } from '../config/x-framework-push-service.config'; import { XPushServiceConnectionRetryPolicy } from '../policies/x-push-service-connection-retry-policy'; @Injectable() export abstract class XBasePushService extends XBaseService { // //#region Props ... // private CONNECTION_ID: string; private HUB_CONNECTION: HubConnection; private HUB_CONNECTION_POLICY: XPushServiceConnectionRetryPolicy; // public readonly connect$ = new Subject(); public readonly error$ = new Subject(); public readonly newConnection$ = new Subject(); public readonly connectionClosed$ = new Subject(); /** * readonly connection id ... */ get connectionId() { // let result = ''; // if (!isNullOrUndefined(this.HUB_CONNECTION)) { result = this.HUB_CONNECTION.connectionId; } // return result; } //#endregion // //#region Abstract ... /** * retrieve specific hub's name and route ... */ abstract hubName(): string; /** * specify the service required access token validator or not ... */ abstract requiredAccessToken(): boolean; /** * specify which protocols can used as transport protocol ... */ abstract transportType(): HttpTransportType | ITransport; /** * Specified Skip Negotiation ... */ public skipNegotiation: boolean = false; //#endregion // //#region Constructor ... constructor( // @Inject(X_FRAMEWORK_PUSH_SERVICE_CONFIG) public config: XFrameworkPushServiceConfig, // @Optional() @Inject(X_FRAMEWORK_PUSH_SERVICE_TOKEN_FETCHER) public tokenFetcher?: IXPushServiceTokenFetcher ) { super(config); } //#endregion // //#region Actions ... /** * determines connection is stablished ot not ... * @returns boolean value ... */ public isConnected() { // const result = !isNullOrUndefined(this.HUB_CONNECTION) && this.HUB_CONNECTION.state === HubConnectionState.Connected; return result; } /** * retrieve hub route ... * @returns string represent hub url ... */ public getHubRoute() { // let result = ''; // // check config exists and // check server route ... if ( isNullOrUndefined(this.config) || isNullOrEmptyString(this.hubName()) || isNullOrEmptyString(this.config.hubServerRoute) ) { return result; } // // Configure route ... result = this.config.hubServerRoute.endsWith('/') ? `${this.config.hubServerRoute}${ this.config.hubsBaseRoute }/${this.hubName()}` : `${this.config.hubServerRoute}/${ this.config.hubsBaseRoute }/${this.hubName()}`; // return result; } /** * connect to server ... */ public async connect() { // // Instantiate Hub Connection ... if (!this.HUB_CONNECTION) { // // Prepare Reconnection Policy ... this.HUB_CONNECTION_POLICY = new XPushServiceConnectionRetryPolicy( this.config ); // // retrieve and validate ... let hubRoute = this.getHubRoute(); if (isNullOrEmptyString(hubRoute)) { throwException(XExceptionIDs.InvalidConfiguration); } // // Normalize Transport Type ... let transport: HttpTransportType | ITransport = HttpTransportType.LongPolling; if (!isNullOrUndefined(this.transportType())) { transport = this.transportType(); } // // Add Custom Headers if Required ... let headers: MessageHeaders = {}; // // Create Connection Options ... const connectionOptions: IHttpConnectionOptions = { headers, transport, skipNegotiation: this.skipNegotiation, }; // // Attach Token Fetcher ... if ( !!this.requiredAccessToken() && !isNullOrUndefined(this.tokenFetcher) ) { connectionOptions.accessTokenFactory = () => this.tokenFetcher.getToken(); } // // Create Connection builder ... const connectionBuilder = new HubConnectionBuilder() .withUrl(hubRoute, connectionOptions) .configureLogging(this.config.connectionLogLevel) .withAutomaticReconnect(this.HUB_CONNECTION_POLICY); // // Handle Message Protocol ... if (!!this.config.addSupportMessageProtocol) { connectionBuilder.withHubProtocol(new MessagePackHubProtocol()); } // // Add LogLevel ... if (!isNullOrUndefined(this.config.connectionLogLevel)) { connectionBuilder.configureLogging(this.config.connectionLogLevel); } // // Create Connection ... this.HUB_CONNECTION = connectionBuilder.build(); // // Add Default Base Actions Event Handlers ... this.on(XBasePushAction.NewConnection, (connectionId) => this.newConnection$.next(connectionId) ); this.on(XBasePushAction.ConnectionClosed, (connectionId) => this.connectionClosed$.next(connectionId) ); } // await this.HUB_CONNECTION.start(); if (!isNullOrEmptyString(this.HUB_CONNECTION.connectionId)) { // this.CONNECTION_ID = this.HUB_CONNECTION.connectionId; this.connect$.next(); this.registerSubjects(); } else { this.CONNECTION_ID = ''; } } /** * Disconnect a Connection if Exists ... */ public async disconnect() { // if ( this.HUB_CONNECTION && this.HUB_CONNECTION.state !== HubConnectionState.Disconnected ) { // await this.HUB_CONNECTION.stop(); this.HUB_CONNECTION = undefined; } } /** * Attach an Event Listener to Hub ... * * @param event * @param callback */ on(event: string, callback: (...args: any[]) => any) { // // Validate ... if ( !isFunction(callback) || isNullOrEmptyString(event) || isNullOrUndefined(callback) || isNullOrUndefined(this.HUB_CONNECTION) ) { return; } // this.HUB_CONNECTION.on(event, callback); } /** * Detach an Event Listener to Hub ... * * @param event * @param callback */ off(event: string, callback?: (...args: any[]) => any): void { // // Validate ... if (isNullOrEmptyString(event) || isNullOrUndefined(this.HUB_CONNECTION)) { return; } // if (isNullOrUndefined(callback)) { this.HUB_CONNECTION.off(event); } else if (!isFunction(callback) || isNullOrUndefined(callback)) { this.HUB_CONNECTION.off(event, callback); } } /** * Invoke Specified Hub Actions ... * * @param method * @param args * @returns */ async invoke(method: string, ...args: any[]) { // // Validate Hub Connected ... if (!this.isConnected() || isNullOrUndefined(this.HUB_CONNECTION)) { return await Promise.resolve(undefined); } // return await this.HUB_CONNECTION.invoke(method, ...args); } //#endregion // //#region Abstractions ... /** * here we have to register all events subjects ... */ registerSubjects() {} //#endregion // //#region Private ... //#endregion } ### FILE: C:\Users\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xFrameworkPushServiceHolder\projects\x-framework-push-service\src\lib\base\x-base-web-rtc.service.ts import { Subject } from 'rxjs'; import { X_FRAMEWORK_PUSH_SERVICE_CONFIG, X_FRAMEWORK_PUSH_SERVICE_TOKEN_FETCHER, } from '../tokens/x-injectable-tokens'; import { isString } from 'x-framework-core'; import { XBasePushService } from './x-base-push.service'; import { Inject, Injectable, Optional } from '@angular/core'; import { XWebRTCAction } from '../typings/x-push-notification.typings'; import { IXPushServiceTokenFetcher } from '../policies/x-push-service-token-fetcher'; import { XFrameworkPushServiceConfig } from '../config/x-framework-push-service.config'; @Injectable() export abstract class XBaseWebRTCService extends XBasePushService { // //#region Props ... public onOfferRecieved$ = new Subject(); public onAnswerRecieved$ = new Subject<{ answer: any; requesterId: string; }>(); public onCandidateRecieved$ = new Subject<{ candidate: any; senderId: string; }>(); //#endregion // //#region Constructor ... constructor( // @Inject(X_FRAMEWORK_PUSH_SERVICE_CONFIG) public config: XFrameworkPushServiceConfig, // @Optional() @Inject(X_FRAMEWORK_PUSH_SERVICE_TOKEN_FETCHER) public tokenFetcher?: IXPushServiceTokenFetcher ) { super(config, tokenFetcher); } //#endregion // //#region Actions ... public async offer(offer: any, to: string) { // if (!isString(offer)) { offer = JSON.stringify(offer); } // await this.invoke(XWebRTCAction.Offer, offer, to); } public async answer(answer: any) { // if (!isString(answer)) { answer = JSON.stringify(answer); } // await this.invoke(XWebRTCAction.Answer, answer); } public async candidate(candidate: any, to: string) { // if (!isString(candidate)) { candidate = JSON.stringify(candidate); } // await this.invoke(XWebRTCAction.Candidate, candidate, to); } //#endregion /** * Register All Available */ registerSubjects(): void { super.registerSubjects(); // this.onOfferRecieved((offer: any) => { this.onOfferRecieved$.next(offer); }); // this.onAnswerRecieved((answer: any, requesterId: string) => { this.onAnswerRecieved$.next({ answer, requesterId }); }); // this.onCandidateRecieved((candidate: any, senderId: string) => { this.onCandidateRecieved$.next({ candidate, senderId }); }); } // //#region Event Handler Attaches ... public onOfferRecieved(callback: (offer: any) => any) { this.on(XWebRTCAction.Offer, callback); } public onAnswerRecieved(callback: (answer: any, requesterId: string) => any) { this.on(XWebRTCAction.Answer, callback); } public onCandidateRecieved( callback: (candidate: any, senderId: string) => any ) { this.on(XWebRTCAction.Candidate, callback); } //#endregion // //#region Event Handler Detachers ... public offOfferRecieved(callback?: (offer: any) => any) { this.off(XWebRTCAction.Offer, callback); } public offAnswerRecieved( callback?: (answer: any, requesterId: string) => any ) { this.off(XWebRTCAction.Answer, callback); } public offCandidateRecieved( callback?: (candidate: any, senderId: string) => any ) { this.off(XWebRTCAction.Candidate, callback); } //#endregion } ### FILE: C:\Users\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xFrameworkPushServiceHolder\projects\x-framework-push-service\src\lib\base\x-rtc-peer.service.ts import { Inject, Injectable } from '@angular/core'; import { X_FRAMEWORK_PUSH_SERVICE_CONFIG } from '../tokens/x-injectable-tokens'; import { XFrameworkPushServiceConfig } from '../config/x-framework-push-service.config'; @Injectable({ providedIn: 'root' }) export class XRTCPeerService { // //#region Constructor ... constructor( @Inject(X_FRAMEWORK_PUSH_SERVICE_CONFIG) public config: XFrameworkPushServiceConfig ) {} //#endregion // //#region Actions ... public createConnection() { // const result = new RTCPeerConnection({ iceServers: [...this.config.iceServers], }); // return result; } //#endregion } ### FILE: C:\Users\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xFrameworkPushServiceHolder\projects\x-framework-push-service\src\lib\config\x-framework-push-service.config.ts import { LogLevel } from '@microsoft/signalr'; import { XFrameworkCoreConfig } from 'x-framework-core'; import { XFrameworkServicesConfig } from 'x-framework-services'; export type XSharedConfig = XFrameworkCoreConfig & XFrameworkServicesConfig; export interface XFrameworkPushServiceConfig extends XSharedConfig { // // Push Service Configuration ... hubsBaseRoute: string; hubServerRoute: string; connectionLogLevel?: LogLevel; pushConnectionMaxRetry: number; addSupportMessageProtocol?: boolean; pushConnectionReconnectDelay: number; // iceServers: Array; } ### FILE: C:\Users\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xFrameworkPushServiceHolder\projects\x-framework-push-service\src\lib\models\x-hub-connection.dto.ts import { XBaseDto } from 'x-framework-core'; export interface XHubConnectionDto extends XBaseDto { userId: string; username?: string; connectionId: string; } ### FILE: C:\Users\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xFrameworkPushServiceHolder\projects\x-framework-push-service\src\lib\policies\x-push-service-connection-retry-policy.ts import { Inject } from '@angular/core'; import { IRetryPolicy, RetryContext } from '@microsoft/signalr'; import { XFrameworkPushServiceConfig } from '../config/x-framework-push-service.config'; import { X_FRAMEWORK_PUSH_SERVICE_CONFIG } from '../tokens/x-injectable-tokens'; export class XPushServiceConnectionRetryPolicy implements IRetryPolicy { // //#region Props ... retryReason: Error = null; previousRetryCount: number = 0; elapsedMilliseconds: number = 0; //#endregion // //#region Constructor ... constructor( @Inject(X_FRAMEWORK_PUSH_SERVICE_CONFIG) public config: XFrameworkPushServiceConfig ) {} //#endregion // nextRetryDelayInMilliseconds(retryContext: RetryContext): number { // this.retryReason = retryContext.retryReason; this.previousRetryCount = retryContext.previousRetryCount; this.elapsedMilliseconds = retryContext.elapsedMilliseconds; // if (retryContext.previousRetryCount < this.config.pushConnectionMaxRetry) { return this.config.pushConnectionReconnectDelay; } else { return null; } } } ### FILE: C:\Users\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xFrameworkPushServiceHolder\projects\x-framework-push-service\src\lib\policies\x-push-service-token-fetcher.ts import { Observable } from 'rxjs'; /** * an interface for accessing tokens and user related infos ... */ export interface IXPushServiceTokenFetcher { token$: Observable; getToken(): string; } ### FILE: C:\Users\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xFrameworkPushServiceHolder\projects\x-framework-push-service\src\lib\tokens\x-injectable-tokens.ts import { InjectionToken } from '@angular/core'; import { IXPushServiceTokenFetcher } from '../policies/x-push-service-token-fetcher'; import { XFrameworkPushServiceConfig } from '../config/x-framework-push-service.config'; export const X_FRAMEWORK_PUSH_SERVICE_CONFIG = new InjectionToken( 'x_config' ); export const X_FRAMEWORK_PUSH_SERVICE_TOKEN_FETCHER = new InjectionToken( 'x_push_token_fetcher' ); ### FILE: C:\Users\SaherElm\Documents\Projects\xSaherelmWorkspace\Modules\xFrameworkPushServiceHolder\projects\x-framework-push-service\src\lib\typings\x-push-notification.typings.ts /** * Base Hub Service Actions ... */ export enum XBasePushAction { /// /// new Connection Join ... /// NewConnection = 'NewConnection', /// /// an Exists Connection Closed ... /// ConnectionClosed = 'ConnectionClosed', } /** * Base Web RTC Actions ... */ export enum XWebRTCAction { Offer = 'Offer', Answer = 'Answer', Candidate = 'Candidate', } /** * Base Entity Push Service Actions ... */ export enum XBaseEntityPushAction { Add = 'Add', Update = 'Update', Delete = 'Delete', AddMany = 'AddMany', UpdateMany = 'UpdateMany', DeleteMany = 'DeleteMany', AddOrUpdate = 'AddOrUpdate', }