Initial Commit on git.saherelm.ir ...
This commit is contained in:
@@ -0,0 +1,11 @@
|
||||
# XFrameworkPushService
|
||||
|
||||
this is xFramework's PushService module and conatins all features of XFramework Library for providing Push functionalities.
|
||||
|
||||
## Maintainer
|
||||
|
||||
Hadi Khazaee asl
|
||||
|
||||
<https://www.saherelm.ir>
|
||||
|
||||
hadi_khazaee_asl@yahoo.com
|
||||
@@ -0,0 +1,7 @@
|
||||
{
|
||||
"$schema": "../../node_modules/ng-packagr/ng-package.schema.json",
|
||||
"dest": "../../dist/x-framework-push-service",
|
||||
"lib": {
|
||||
"entryFile": "src/public-api.ts"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,22 @@
|
||||
{
|
||||
"name": "x-framework-push-service",
|
||||
"version": "0.0.1",
|
||||
"description": "access xFrameworkPushServices features ...",
|
||||
"keywords": [
|
||||
"x-framework",
|
||||
"x-framework-services"
|
||||
],
|
||||
"homepage": "https://saherelm.ir",
|
||||
"author": {
|
||||
"name": "Hadi Khazaee Asl",
|
||||
"email": "hadi_khazaee_asl@yahoo.com",
|
||||
"url": "https://saherelm.ir"
|
||||
},
|
||||
"peerDependencies": {
|
||||
"@angular/common": "^12.0.4",
|
||||
"@angular/core": "^12.0.4"
|
||||
},
|
||||
"dependencies": {
|
||||
"tslib": "^2.1.0"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,385 @@
|
||||
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<TEntity> 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<TEntity>; actorId: string }>();
|
||||
onManyUpdated$ = new Subject<{ model: Array<TEntity>; actorId: string }>();
|
||||
onManyDeleted$ = new Subject<{ model: Array<TEntity>; 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<TEntity>) {
|
||||
//
|
||||
// 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<TEntity>) {
|
||||
//
|
||||
// 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<TEntity>) {
|
||||
//
|
||||
// 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<TEntity>, actorId: string) => {
|
||||
this.onManyAdded$.next({ model, actorId });
|
||||
});
|
||||
|
||||
//
|
||||
this.onManyUpdated((model: Array<TEntity>, actorId: string) => {
|
||||
this.onManyUpdated$.next({ model, actorId });
|
||||
});
|
||||
|
||||
//
|
||||
this.onManyDeleted((model: Array<TEntity>, 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<TEntity>(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<TEntity>(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<TEntity>(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<TEntity>(modelJson);
|
||||
callback(model, actorId);
|
||||
}
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Register Many Added Event Listener ...
|
||||
*
|
||||
* @param callback
|
||||
*/
|
||||
public onManyAdded(
|
||||
callback: (model: Array<TEntity>, actorId: string) => any
|
||||
) {
|
||||
//
|
||||
this.on(
|
||||
XBaseEntityPushAction.AddMany,
|
||||
(modelJson: string, actorId: string) => {
|
||||
//
|
||||
const model = parseJson<Array<TEntity>>(modelJson);
|
||||
callback(model, actorId);
|
||||
}
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Register Many Updated Event Listener ...
|
||||
*
|
||||
* @param callback
|
||||
*/
|
||||
public onManyUpdated(
|
||||
callback: (model: Array<TEntity>, actorId: string) => any
|
||||
) {
|
||||
//
|
||||
this.on(
|
||||
XBaseEntityPushAction.UpdateMany,
|
||||
(modelJson: string, actorId: string) => {
|
||||
//
|
||||
const model = parseJson<Array<TEntity>>(modelJson);
|
||||
callback(model, actorId);
|
||||
}
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Register Many Deleted Event Listener ...
|
||||
*
|
||||
* @param callback
|
||||
*/
|
||||
public onManyDeleted(
|
||||
callback: (model: Array<TEntity>, actorId: string) => any
|
||||
) {
|
||||
//
|
||||
this.on(
|
||||
XBaseEntityPushAction.DeleteMany,
|
||||
(modelJson: string, actorId: string) => {
|
||||
//
|
||||
const model = parseJson<Array<TEntity>>(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<TEntity>, actorId: string) => any
|
||||
) {
|
||||
this.off(XBaseEntityPushAction.AddMany, callback);
|
||||
}
|
||||
|
||||
/**
|
||||
* UnRegister Many Updated Event Listener ...
|
||||
*
|
||||
* @param callback
|
||||
*/
|
||||
public offManyUpdated(
|
||||
callback?: (model: Array<TEntity>, actorId: string) => any
|
||||
) {
|
||||
this.off(XBaseEntityPushAction.UpdateMany, callback);
|
||||
}
|
||||
|
||||
/**
|
||||
* UnRegister Many Deleted Event Listener ...
|
||||
*
|
||||
* @param callback
|
||||
*/
|
||||
public offManyDeleted(
|
||||
callback?: (model: Array<TEntity>, actorId: string) => any
|
||||
) {
|
||||
this.off(XBaseEntityPushAction.DeleteMany, callback);
|
||||
}
|
||||
//#endregion
|
||||
}
|
||||
@@ -0,0 +1,330 @@
|
||||
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<any>();
|
||||
public readonly newConnection$ = new Subject<string>();
|
||||
public readonly connectionClosed$ = new Subject<string>();
|
||||
|
||||
/**
|
||||
* 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<T = any>(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<T>(method, ...args);
|
||||
}
|
||||
//#endregion
|
||||
|
||||
//
|
||||
//#region Abstractions ...
|
||||
/**
|
||||
* here we have to register all events subjects ...
|
||||
*/
|
||||
registerSubjects() {}
|
||||
//#endregion
|
||||
|
||||
//
|
||||
//#region Private ...
|
||||
//#endregion
|
||||
}
|
||||
@@ -0,0 +1,134 @@
|
||||
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<any>();
|
||||
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
|
||||
}
|
||||
@@ -0,0 +1,30 @@
|
||||
|
||||
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
|
||||
}
|
||||
@@ -0,0 +1,19 @@
|
||||
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<RTCIceServer>;
|
||||
}
|
||||
@@ -0,0 +1,7 @@
|
||||
import { XBaseDto } from 'x-framework-core';
|
||||
|
||||
export interface XHubConnectionDto extends XBaseDto {
|
||||
userId: string;
|
||||
username?: string;
|
||||
connectionId: string;
|
||||
}
|
||||
+36
@@ -0,0 +1,36 @@
|
||||
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;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,9 @@
|
||||
import { Observable } from 'rxjs';
|
||||
|
||||
/**
|
||||
* an interface for accessing tokens and user related infos ...
|
||||
*/
|
||||
export interface IXPushServiceTokenFetcher {
|
||||
token$: Observable<string>;
|
||||
getToken(): string;
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
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<XFrameworkPushServiceConfig>(
|
||||
'x_config'
|
||||
);
|
||||
|
||||
export const X_FRAMEWORK_PUSH_SERVICE_TOKEN_FETCHER = new InjectionToken<IXPushServiceTokenFetcher>(
|
||||
'x_push_token_fetcher'
|
||||
);
|
||||
@@ -0,0 +1,35 @@
|
||||
/**
|
||||
* Base Hub Service Actions ...
|
||||
*/
|
||||
export enum XBasePushAction {
|
||||
/// <summary>
|
||||
/// new Connection Join ...
|
||||
/// </summary>
|
||||
NewConnection = 'NewConnection',
|
||||
/// <summary>
|
||||
/// an Exists Connection Closed ...
|
||||
/// </summary>
|
||||
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',
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
import { NgModule } from '@angular/core';
|
||||
|
||||
@NgModule({
|
||||
declarations: [
|
||||
],
|
||||
imports: [
|
||||
],
|
||||
exports: [
|
||||
]
|
||||
})
|
||||
export class XFrameworkPushServiceModule { }
|
||||
@@ -0,0 +1,31 @@
|
||||
/*
|
||||
* 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';
|
||||
@@ -0,0 +1,20 @@
|
||||
/* To learn more about this file see: https://angular.io/config/tsconfig. */
|
||||
{
|
||||
"extends": "../../tsconfig.json",
|
||||
"compilerOptions": {
|
||||
"outDir": "../../out-tsc/lib",
|
||||
"target": "es2015",
|
||||
"declaration": true,
|
||||
"declarationMap": true,
|
||||
"inlineSources": true,
|
||||
"types": [],
|
||||
"lib": [
|
||||
"dom",
|
||||
"es2018"
|
||||
]
|
||||
},
|
||||
"exclude": [
|
||||
"src/test.ts",
|
||||
"**/*.spec.ts"
|
||||
]
|
||||
}
|
||||
@@ -0,0 +1,10 @@
|
||||
/* To learn more about this file see: https://angular.io/config/tsconfig. */
|
||||
{
|
||||
"extends": "./tsconfig.lib.json",
|
||||
"compilerOptions": {
|
||||
"declarationMap": false
|
||||
},
|
||||
"angularCompilerOptions": {
|
||||
"compilationMode": "partial"
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user