175 lines
4.3 KiB
TypeScript
175 lines
4.3 KiB
TypeScript
/**
|
|
* @author Martin Karkowski
|
|
* @email m.karkowski@zema.de
|
|
* @create date 2020-11-06 08:52:36
|
|
* @modify date 2020-12-03 16:34:27
|
|
* @desc [description]
|
|
*/
|
|
|
|
import { EventEmitter } from "events";
|
|
import { ILogger } from "js-logger";
|
|
import { NopeObservable } from "../observables/nopeObservable";
|
|
import {
|
|
IAvailableInstanceGeneratorsMsg,
|
|
IAvailableInstancesMsg,
|
|
IAvailableServicesMsg,
|
|
IAvailableTopicsMsg,
|
|
ICommunicationInterface,
|
|
IExternalEventMsg,
|
|
IRequestTaskMsg,
|
|
IResponseTaskMsg,
|
|
ITaskCancelationMsg
|
|
} from "../types/nope/nopeCommunication.interface";
|
|
import { INopeObservable } from "../types/nope/nopeObservable.interface";
|
|
|
|
/**
|
|
* A Communication Layer for the Dispatchers.
|
|
* Here, only a Events are used.
|
|
*
|
|
* @export
|
|
* @class EventLayer
|
|
* @implements {ICommunicationInterface}
|
|
*/
|
|
export class EventLayer implements ICommunicationInterface {
|
|
connected: INopeObservable<boolean>;
|
|
|
|
protected _emitter = new EventEmitter();
|
|
|
|
constructor(
|
|
public readonly subscriptionMode: "individual" | "generic" = "individual",
|
|
public readonly resultSharing: "individual" | "generic" = "individual",
|
|
protected _logger?: ILogger
|
|
) {
|
|
this.connected = new NopeObservable();
|
|
this.connected.setContent(true);
|
|
}
|
|
|
|
async onNewInstancesAvailable(
|
|
cb: (instances: IAvailableInstancesMsg) => void
|
|
): Promise<void> {
|
|
this._emitter.on("newInstancesAvailable", cb);
|
|
}
|
|
|
|
async emitNewInstancesAvailable(
|
|
instances: IAvailableInstancesMsg
|
|
): Promise<void> {
|
|
this._emitter.emit("newInstancesAvailable", instances);
|
|
}
|
|
|
|
async onTaskCancelation(
|
|
cb: (msg: ITaskCancelationMsg) => void
|
|
): Promise<void> {
|
|
this._emitter.on("cancel", cb);
|
|
}
|
|
|
|
async emitTaskCancelation(msg: ITaskCancelationMsg): Promise<void> {
|
|
this._emitter.emit("cancel", msg);
|
|
}
|
|
|
|
async onAurevoir(cb: (dispatcher: string) => void): Promise<void> {
|
|
this._emitter.on("aurevoir", cb);
|
|
}
|
|
|
|
async emitAurevoir(dispatcher: string): Promise<void> {
|
|
this._emitter.emit("aurevoir", dispatcher);
|
|
}
|
|
|
|
async emitNewInstanceGeneratorsAvailable(
|
|
generators: IAvailableInstanceGeneratorsMsg
|
|
): Promise<void> {
|
|
this._emitter.emit("generators", generators);
|
|
}
|
|
|
|
async onNewInstanceGeneratorsAvailable(
|
|
cb: (generators: IAvailableInstanceGeneratorsMsg) => void
|
|
): Promise<void> {
|
|
this._emitter.on("generators", cb);
|
|
}
|
|
|
|
async emitRpcRequest(name: string, request: IRequestTaskMsg): Promise<void> {
|
|
this._emitter.emit(name, request);
|
|
}
|
|
|
|
async emitRpcResult(name: string, result: IResponseTaskMsg): Promise<void> {
|
|
this._emitter.emit(name, result);
|
|
}
|
|
|
|
async onRpcResult(
|
|
name: string,
|
|
cb: (result: IResponseTaskMsg) => void
|
|
): Promise<void> {
|
|
this._emitter.on(name, cb);
|
|
}
|
|
|
|
async offRpcResponse(
|
|
name: string,
|
|
cb: (result: IResponseTaskMsg) => void
|
|
): Promise<void> {
|
|
this._emitter.off(name, cb);
|
|
}
|
|
|
|
async onRpcRequest(
|
|
name: string,
|
|
cb: (data: IRequestTaskMsg) => void
|
|
): Promise<void> {
|
|
this._emitter.on(name, cb);
|
|
}
|
|
|
|
async offRpcRequest(
|
|
name: string,
|
|
cb: (data: IRequestTaskMsg) => void
|
|
): Promise<void> {
|
|
this._emitter.off(name, cb);
|
|
}
|
|
|
|
async emitNewServicesAvailable(
|
|
services: IAvailableServicesMsg
|
|
): Promise<void> {
|
|
this._emitter.emit("services", services);
|
|
}
|
|
|
|
async onNewServicesAvailable(
|
|
cb: (services: IAvailableServicesMsg) => void
|
|
): Promise<void> {
|
|
this._emitter.on("services", cb);
|
|
}
|
|
|
|
async onBonjour(cb): Promise<void> {
|
|
this._emitter.on("bonjour", cb);
|
|
}
|
|
|
|
async emitBonjour(msg): Promise<void> {
|
|
this._emitter.emit("bonjour", msg);
|
|
}
|
|
|
|
async emitNewObersvablesAvailable(
|
|
topics: IAvailableTopicsMsg
|
|
): Promise<void> {
|
|
this._emitter.emit("topics", topics);
|
|
}
|
|
|
|
async onNewObservablesAvailable(
|
|
cb: (topics: IAvailableTopicsMsg) => void
|
|
): Promise<void> {
|
|
this._emitter.on("topics", cb);
|
|
}
|
|
|
|
async onEvent(
|
|
event: string,
|
|
cb: (data: IExternalEventMsg) => void
|
|
): Promise<void> {
|
|
this._emitter.on("event_" + event, cb);
|
|
}
|
|
|
|
async emitEvent(event: string, data: IExternalEventMsg): Promise<void> {
|
|
this._emitter.emit("event_" + event, data);
|
|
}
|
|
|
|
async offEvent(
|
|
event: string,
|
|
cb: (data: IExternalEventMsg) => void
|
|
): Promise<void> {
|
|
this._emitter.off("event_" + event, cb);
|
|
}
|
|
}
|