socketio drop in for bsock
This commit is contained in:
+14
-24
@@ -1,17 +1,17 @@
|
||||
'use strict'
|
||||
|
||||
import http = require('http');
|
||||
import bsock = require('bsock');
|
||||
import { PromiseIO } from "./PromiseIO/Server";
|
||||
import * as T from './Types';
|
||||
import * as U from './Utils';
|
||||
import * as I from './Interfaces';
|
||||
|
||||
export class RPCServer<
|
||||
InterfaceT extends T.RPCInterface = T.RPCInterface,
|
||||
> implements I.Destroyable {
|
||||
> {
|
||||
|
||||
private ws = http.createServer()
|
||||
private io = bsock.createServer()
|
||||
private pio = PromiseIO.createServer()
|
||||
private visibility: T.Visibility
|
||||
private closeHandler: T.CloseHandler
|
||||
private errorHandler: T.ErrorHandler
|
||||
@@ -42,7 +42,7 @@ export class RPCServer<
|
||||
})
|
||||
|
||||
|
||||
this.errorHandler = (socket: I.Socket) => (error: any, rpcName: string, args: any[]) => {
|
||||
this.errorHandler = (socket: I.Socket | PromiseIO) => (error: any, rpcName: string, args: any[]) => {
|
||||
if (conf.errorHandler) conf.errorHandler(socket, error, rpcName, args)
|
||||
else throw error
|
||||
}
|
||||
@@ -72,20 +72,24 @@ export class RPCServer<
|
||||
|
||||
private startWebsocket() {
|
||||
try {
|
||||
this.io.attach(this.ws)
|
||||
this.io.on('socket', (socket: I.Socket) => {
|
||||
this.pio.attach(this.ws)
|
||||
|
||||
this.pio.on('socket', (socket: I.Socket) => {
|
||||
socket.on('error', (err) => this.errorHandler(socket, err, "system", []))
|
||||
socket.on('close', () => this.closeHandler(socket))
|
||||
this.connectionHandler(socket)
|
||||
this.initRPCs(socket)
|
||||
})
|
||||
this.ws = this.ws.listen(this.port, this.visibility)
|
||||
|
||||
this.pio.listen(this.port)
|
||||
|
||||
} catch (e) {
|
||||
this.errorHandler(this.io, e, 'system', [])
|
||||
this.errorHandler(this.pio, e, 'system', [])
|
||||
}
|
||||
}
|
||||
|
||||
protected initRPCs(socket: I.Socket) {
|
||||
|
||||
socket.hook('info', async (sesame?: string) => {
|
||||
const rpcs = await Promise.all(this.exporters.map(async exp => {
|
||||
const allowed = await this.accessFilter(sesame, exp)
|
||||
@@ -96,22 +100,8 @@ export class RPCServer<
|
||||
})
|
||||
}
|
||||
|
||||
/**
|
||||
* Publishes a new list of Exporters. This destroys and restarts the socket
|
||||
* @param exporters the exporters to publish
|
||||
public setExporters(exporters: T.ExporterArray<InterfaceT>): any {
|
||||
exporters.forEach(U.fixNames)
|
||||
this.destroy()
|
||||
this.ws = http.createServer()
|
||||
this.io = bsock.createServer()
|
||||
this.exporters = exporters
|
||||
this.startWebsocket()
|
||||
}
|
||||
*/
|
||||
|
||||
|
||||
destroy(): void {
|
||||
this.io.close()
|
||||
close(): void {
|
||||
this.pio.close()
|
||||
this.ws.close()
|
||||
}
|
||||
}
|
||||
+42
-23
@@ -1,7 +1,6 @@
|
||||
'use strict'
|
||||
|
||||
import bsock = require('bsock');
|
||||
|
||||
import { PromiseIOClient } from './PromiseIO/Client'
|
||||
import * as T from './Types';
|
||||
import * as I from './Interfaces';
|
||||
import { stripAfterEquals, appendComma } from './Utils';
|
||||
@@ -18,8 +17,12 @@ export class RPCSocket<Ifc extends T.RPCInterface = T.RPCInterface> implements I
|
||||
}
|
||||
|
||||
private socket: I.Socket
|
||||
private closeHandlers: T.FrontEndHandlerType['close'][] = []
|
||||
private errorHandlers: T.FrontEndHandlerType['error'][] = []
|
||||
private handlers : {
|
||||
[name in string]: T.AnyFunction[]
|
||||
} = {
|
||||
error: [],
|
||||
close: []
|
||||
}
|
||||
private hooks : {[name in string]: T.AnyFunction} = {}
|
||||
|
||||
/**
|
||||
@@ -45,6 +48,19 @@ export class RPCSocket<Ifc extends T.RPCInterface = T.RPCInterface> implements I
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Hooks a handler to a function name. Use {@link call} to trigger it.
|
||||
* @param name The function name to listen on
|
||||
* @param handler The handler to attach
|
||||
*/
|
||||
public bind(name: string, handler: (...args:any[]) => any | Promise<any>){
|
||||
if(!this.socket){
|
||||
this.hooks[name] = handler
|
||||
}else{
|
||||
this.socket.bind(name, handler)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Removes a {@link hook} listener by name.
|
||||
* @param name The function name
|
||||
@@ -62,13 +78,12 @@ export class RPCSocket<Ifc extends T.RPCInterface = T.RPCInterface> implements I
|
||||
* @param type 'error' or 'close'
|
||||
* @param f The listener to attach
|
||||
*/
|
||||
public on<T extends "error" | "close">(type: T, f: T.FrontEndHandlerType[T]){
|
||||
public on(type: string, f: T.AnyFunction){
|
||||
if(!this.socket){
|
||||
switch(type){
|
||||
case "error": this.errorHandlers.push(<T.FrontEndHandlerType['error']> f); break;
|
||||
case "close": this.closeHandlers.push(<T.FrontEndHandlerType['close']> f); break;
|
||||
default: throw new Error('socket.on only supports ´error´ and ´close´ as first parameter. Got: ´'+type+'´')
|
||||
}
|
||||
if(!this.handlers[type])
|
||||
this.handlers[type] = []
|
||||
|
||||
this.handlers[type].push(f)
|
||||
}else{
|
||||
this.socket.on(type, f)
|
||||
}
|
||||
@@ -84,14 +99,6 @@ export class RPCSocket<Ifc extends T.RPCInterface = T.RPCInterface> implements I
|
||||
this.socket.emit(eventName, data)
|
||||
}
|
||||
|
||||
/**
|
||||
* Destroys the socket
|
||||
*/
|
||||
public destroy(){
|
||||
if(!this.socket) return;
|
||||
this.socket.destroy()
|
||||
}
|
||||
|
||||
/**
|
||||
* Closes the socket. It may attempt to reconnect.
|
||||
*/
|
||||
@@ -108,7 +115,8 @@ export class RPCSocket<Ifc extends T.RPCInterface = T.RPCInterface> implements I
|
||||
public async call (rpcname: string, ...args: any[]) : Promise<any>{
|
||||
if(!this.socket) throw new Error("The socket is not connected! Use socket.connect() first")
|
||||
try{
|
||||
return await this.socket.call.apply(this.socket, [rpcname, ...args])
|
||||
const val = await this.socket.call.apply(this.socket, [rpcname, ...args])
|
||||
return val
|
||||
}catch(e){
|
||||
this.emit('error', e)
|
||||
throw e
|
||||
@@ -129,14 +137,23 @@ export class RPCSocket<Ifc extends T.RPCInterface = T.RPCInterface> implements I
|
||||
* Connects to the server and attaches available RPCs to this object
|
||||
*/
|
||||
public async connect( sesame?: string ) : Promise<T.ConnectedSocket<Ifc>> {
|
||||
this.socket = await bsock.connect(this.port, this.server, this.conf.tls?this.conf.tls:false)
|
||||
this.errorHandlers.forEach(h => this.socket.on('error', h))
|
||||
this.closeHandlers.forEach(h => this.socket.on('close', h))
|
||||
|
||||
try{
|
||||
this.socket = await PromiseIOClient.connect(this.port, this.server, /*this.conf.tls?this.conf.tls:false*/)
|
||||
}catch(e){
|
||||
this.handlers['error'].forEach(h => h(e))
|
||||
throw e
|
||||
}
|
||||
|
||||
Object.entries(this.handlers).forEach(([k,v])=>{
|
||||
v.forEach(h => this.socket.on(k, h))
|
||||
})
|
||||
|
||||
Object.entries(this.hooks).forEach((kv: [string, T.AnyFunction]) => {
|
||||
this.socket.hook(kv[0], kv[1])
|
||||
})
|
||||
|
||||
const info:T.ExtendedRpcInfo[] = await this.info(sesame)
|
||||
|
||||
info.forEach(i => {
|
||||
let f: any
|
||||
|
||||
@@ -153,6 +170,8 @@ export class RPCSocket<Ifc extends T.RPCInterface = T.RPCInterface> implements I
|
||||
this[i.owner][i.name] = f
|
||||
this[i.owner][i.name].bind(this)
|
||||
})
|
||||
|
||||
|
||||
return <T.ConnectedSocket<Ifc>> (this as any)
|
||||
}
|
||||
|
||||
|
||||
+12
-17
@@ -12,20 +12,15 @@ export type RPCExporter<
|
||||
exportRPCs() : T.RPCDefinitions<Ifc>[Name]
|
||||
}
|
||||
|
||||
/**
|
||||
* Generic socket interface that can apply to bsock as well as RPCSocket
|
||||
*/
|
||||
export interface Socket extends Destroyable {
|
||||
port: number
|
||||
hook: (rpcname: string, handler: T.AnyFunction) => void
|
||||
unhook: (rpcname:string) => void
|
||||
call: (rpcname:string, ...args: any[]) => Promise<any>
|
||||
fire: (rpcname:string, ...args: any[]) => Promise<any>
|
||||
on: T.OnFunction
|
||||
emit: (eventName: string, data:any) => void
|
||||
close() : void
|
||||
}
|
||||
|
||||
export interface Destroyable{
|
||||
destroy() : void
|
||||
}
|
||||
export interface Socket {
|
||||
id?: string
|
||||
bind: (name: string, listener: T.PioBindListener) => void
|
||||
|
||||
hook: (rpcname: string, handler: T.PioHookListener) => void
|
||||
unhook: (rpcname: string, listener?:T.AnyFunction) => void
|
||||
call: (rpcname: string, ...args: any[]) => Promise<any>
|
||||
fire: (rpcname: string, ...args: any[]) => Promise<any>
|
||||
on: (type: string, f: T.AnyFunction)=>any
|
||||
emit: (eventName: string, data: any) => void
|
||||
close(): void
|
||||
}
|
||||
@@ -0,0 +1,38 @@
|
||||
import { Socket } from "socket.io"
|
||||
import * as U from '../Utils'
|
||||
import * as I from '../Interfaces'
|
||||
import * as socketio from 'socket.io-client'
|
||||
|
||||
export class PromiseIOClient {
|
||||
|
||||
static connect = (port: number, host = "localhost"): Promise<I.Socket> => new Promise((res, rej) => {
|
||||
try {
|
||||
const socket = socketio(`http://${host}:${port}`, {
|
||||
reconnectionAttempts: 2,
|
||||
reconnectionDelay: 200,
|
||||
timeout: 450,
|
||||
reconnection: false
|
||||
})
|
||||
socket.on('connect_error', e => {
|
||||
sock.emit('error', e)
|
||||
rej(e)
|
||||
})
|
||||
|
||||
const sock = U.makePioSocket(socket)
|
||||
socket.on('connect', ()=>{ res(sock) })
|
||||
|
||||
|
||||
/*
|
||||
socket.on('connect_timeout', ()=>console.log('connect_timeout'))
|
||||
socket.on('disconnect', ()=>console.log('disconnect'))
|
||||
socket.on('reconnect', ()=>console.log('reconnect'))
|
||||
socket.on('reconnect_attempt', ()=>console.log('reconnect_attempt'))
|
||||
socket.on('reconnecting', ()=>console.log('reconnecting'));
|
||||
socket.on('reconnect_failed', ()=>console.log('reconnect_failed'));
|
||||
socket.on('reconnecting', ()=>console.log('reconnecting'));
|
||||
*/
|
||||
} catch (e) {
|
||||
rej(e)
|
||||
}
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,57 @@
|
||||
import { Server, Socket } from "socket.io"
|
||||
import { Server as httpServer } from "http"
|
||||
import * as U from '../Utils'
|
||||
import * as T from '../Types'
|
||||
const socketio = require('socket.io')
|
||||
|
||||
export class PromiseIO {
|
||||
io?: Server
|
||||
private listeners: { [eventName in string]: ((...args: any) => void)[] } = {
|
||||
socket: [],
|
||||
connect: []
|
||||
}
|
||||
|
||||
static createServer(): PromiseIO {
|
||||
return new PromiseIO();
|
||||
}
|
||||
|
||||
attach(httpServer: httpServer) {
|
||||
this.io = socketio(httpServer)
|
||||
}
|
||||
|
||||
listen(port: number) {
|
||||
this.io!.on('connection', (sock: Socket) => {
|
||||
const pioSock = U.makePioSocket(sock)
|
||||
this.listeners['socket'].forEach(listener => listener(pioSock))
|
||||
this.listeners['connect'].forEach(listener => listener(pioSock))
|
||||
/*
|
||||
pioSock.on('error', ()=>console.log('error'));
|
||||
|
||||
pioSock.on('connect_timeout', ()=>console.log('connect_timeout'))
|
||||
pioSock.on('disconnect', ()=>console.log('disconnect'))
|
||||
pioSock.on('reconnect', ()=>console.log('reconnect'))
|
||||
pioSock.on('reconnect_attempt', ()=>console.log('reconnect_attempt'))
|
||||
pioSock.on('reconnecting', ()=>console.log('reconnecting'));
|
||||
pioSock.on('reconnect_failed', ()=>console.log('reconnect_failed'));
|
||||
pioSock.on('reconnecting', ()=>console.log('reconnecting'));
|
||||
*/
|
||||
})
|
||||
this.io!.listen(port)
|
||||
}
|
||||
|
||||
on(eventName: string, listener: T.AnyFunction) {
|
||||
if (this.listeners[eventName] == null) {
|
||||
this.listeners[eventName] = []
|
||||
}
|
||||
this.listeners[eventName].push(listener)
|
||||
}
|
||||
|
||||
close = () => {
|
||||
if(this.io){
|
||||
this.io.engine.ws.close()
|
||||
this.io.close()
|
||||
this.io = undefined
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
+7
-1
@@ -1,12 +1,18 @@
|
||||
import * as I from "./Interfaces";
|
||||
import { RPCSocket } from "./Frontend";
|
||||
import { PromiseIO } from "./PromiseIO/Server";
|
||||
|
||||
export type PioBindListener = (...args: any) => void
|
||||
export type PioHookListener = AnyFunction
|
||||
|
||||
|
||||
|
||||
export type AnyFunction = (...args:any) => any
|
||||
export type HookFunction = AnyFunction
|
||||
export type AccessFilter<InterfaceT extends RPCInterface = RPCInterface> = (sesame:string|undefined, exporter: I.RPCExporter<InterfaceT, keyof InterfaceT>) => Promise<boolean> | boolean
|
||||
export type Visibility = "127.0.0.1" | "0.0.0.0"
|
||||
export type ConnectionHandler = (socket:I.Socket) => void
|
||||
export type ErrorHandler = (socket:I.Socket, error:any, rpcName: string, args: any[]) => void
|
||||
export type ErrorHandler = (socket:I.Socket | PromiseIO, error:any, rpcName: string, args: any[]) => void
|
||||
export type CloseHandler = (socket:I.Socket) => void
|
||||
export type SesameFunction = (sesame : string) => boolean
|
||||
export type SesameConf = {
|
||||
|
||||
+129
-44
@@ -2,6 +2,8 @@ import * as uuidv4 from "uuid/v4"
|
||||
|
||||
import * as T from "./Types";
|
||||
import * as I from "./Interfaces";
|
||||
import { Server as ioServer, Socket as ioSocket, Socket } from "socket.io"
|
||||
import { Socket as ioClientSocket } from "socket.io-client"
|
||||
|
||||
/**
|
||||
* Translate an RPC to RPCInfo for serialization.
|
||||
@@ -11,18 +13,18 @@ import * as I from "./Interfaces";
|
||||
* @param sesame optional sesame phrase to prepend before all RPC arguments
|
||||
* @throws Error on RPC without name property
|
||||
*/
|
||||
export const rpcToRpcinfo = (socket: I.Socket, rpc : T.RPC<any, any>, owner: string, errorHandler: T.ErrorHandler, sesame?:T.SesameFunction):T.RpcInfo => {
|
||||
switch (typeof rpc){
|
||||
case "object":
|
||||
if(rpc['call']){
|
||||
export const rpcToRpcinfo = (socket: I.Socket, rpc: T.RPC<any, any>, owner: string, errorHandler: T.ErrorHandler, sesame?: T.SesameFunction): T.RpcInfo => {
|
||||
switch (typeof rpc) {
|
||||
case "object":
|
||||
if (rpc['call']) {
|
||||
return {
|
||||
owner: owner,
|
||||
argNames: extractArgs(rpc['call']),
|
||||
type: "Call",
|
||||
name: rpc.name,
|
||||
call: sesame?async (_sesame, ...args) => {if(sesame(_sesame)) return await rpc['call'].apply({}, args); socket.destroy()}:rpc['call'], // check & remove sesame
|
||||
call: sesame ? async (_sesame, ...args) => { if (sesame(_sesame)) return await rpc['call'].apply({}, args); socket.close() } : rpc['call'], // check & remove sesame
|
||||
}
|
||||
}else{
|
||||
} else {
|
||||
const generator = hookGenerator(<T.HookRPC<any, any>>rpc, errorHandler, sesame)
|
||||
return {
|
||||
owner: owner,
|
||||
@@ -33,7 +35,7 @@ export const rpcToRpcinfo = (socket: I.Socket, rpc : T.RPC<any, any>, owner: str
|
||||
}
|
||||
}
|
||||
case "function":
|
||||
if(!rpc.name) throw new Error(`
|
||||
if (!rpc.name) throw new Error(`
|
||||
RPC did not provide a name.
|
||||
\nUse 'funtion name(..){ .. }' syntax instead.
|
||||
\n
|
||||
@@ -41,14 +43,14 @@ RPC did not provide a name.
|
||||
\n${rpc.toString()}
|
||||
\n>------------OFFENDING RPC`)
|
||||
return {
|
||||
owner : owner,
|
||||
owner: owner,
|
||||
argNames: extractArgs(rpc),
|
||||
type: "Call",
|
||||
name: rpc.name,
|
||||
call: sesame?async (_sesame, ...args) => {if(sesame(_sesame)) return await rpc.apply({}, args); throw makeError(rpc.name)}:rpc, // check & remove sesame
|
||||
call: sesame ? async (_sesame, ...args) => { if (sesame(_sesame)) return await rpc.apply({}, args); throw makeError(rpc.name) } : rpc, // check & remove sesame
|
||||
}
|
||||
}
|
||||
throw new Error("Bad socketIORPC type "+ typeof rpc)
|
||||
throw new Error("Bad socketIORPC type " + typeof rpc)
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -57,37 +59,37 @@ RPC did not provide a name.
|
||||
* @param exporter The exporter
|
||||
* @param makeUnique @default true Attach a suffix to RPC names
|
||||
*/
|
||||
export function rpcHooker(socket: I.Socket, exporter:I.RPCExporter<any, any>, errorHandler: T.ErrorHandler, sesame?:T.SesameFunction, makeUnique = true):T.ExtendedRpcInfo[]{
|
||||
export function rpcHooker(socket: I.Socket, exporter: I.RPCExporter<any, any>, errorHandler: T.ErrorHandler, sesame?: T.SesameFunction, makeUnique = true): T.ExtendedRpcInfo[] {
|
||||
const owner = exporter.name
|
||||
const RPCs = exporter.exportRPCs()
|
||||
|
||||
return RPCs.map(rpc => rpcToRpcinfo(socket, rpc, owner, errorHandler, sesame))
|
||||
.map(info => {
|
||||
const suffix = makeUnique?"-"+uuidv4().substr(0,4):""
|
||||
const ret:any = info
|
||||
ret.uniqueName = info.name+suffix
|
||||
let rpcFunction = info.type === 'Hook'? info.generator(socket)
|
||||
: info.call
|
||||
.map(info => {
|
||||
const suffix = makeUnique ? "-" + uuidv4().substr(0, 4) : ""
|
||||
const ret: any = info
|
||||
ret.uniqueName = info.name + suffix
|
||||
let rpcFunction = info.type === 'Hook' ? info.generator(socket)
|
||||
: info.call
|
||||
|
||||
socket.hook(ret.uniqueName, callGenerator(info.name, socket, rpcFunction, errorHandler))
|
||||
return ret
|
||||
})
|
||||
socket.hook(ret.uniqueName, callGenerator(info.name, socket, rpcFunction, errorHandler))
|
||||
return ret
|
||||
})
|
||||
}
|
||||
|
||||
/**
|
||||
* Decorate an RPC with the error handler
|
||||
* @param rpcFunction the function to decorate
|
||||
*/
|
||||
const callGenerator = (rpcName : string, socket: I.Socket, rpcFunction : T.AnyFunction, errorHandler: T.ErrorHandler) : T.AnyFunction => {
|
||||
const callGenerator = (rpcName: string, socket: I.Socket, rpcFunction: T.AnyFunction, errorHandler: T.ErrorHandler): T.AnyFunction => {
|
||||
const argsArr = extractArgs(rpcFunction)
|
||||
const args = argsArr.join(',')
|
||||
const argsStr = argsArr.map(stripAfterEquals).join(',')
|
||||
|
||||
return eval(`async (`+args+`) => {
|
||||
return eval(`async (` + args + `) => {
|
||||
try{
|
||||
return await rpcFunction(`+argsStr+`)
|
||||
return await rpcFunction(`+ argsStr + `)
|
||||
}catch(e){
|
||||
errorHandler(socket)(e, rpcName, [`+args+`])
|
||||
errorHandler(socket)(e, rpcName, [`+ args + `])
|
||||
}
|
||||
}`)
|
||||
}
|
||||
@@ -96,7 +98,7 @@ const callGenerator = (rpcName : string, socket: I.Socket, rpcFunction : T.AnyFu
|
||||
* Utility function to strip parameters like "a = 3" of their defaults
|
||||
* @param str The parameter to modify
|
||||
*/
|
||||
export function stripAfterEquals(str:string):string{
|
||||
export function stripAfterEquals(str: string): string {
|
||||
return str.split("=")[0]
|
||||
}
|
||||
|
||||
@@ -105,13 +107,14 @@ export function stripAfterEquals(str:string):string{
|
||||
* @param rpc The RPC to transform
|
||||
* @returns A {@link HookFunction}
|
||||
*/
|
||||
const hookGenerator = (rpc:T.HookRPC<any, any>, /*not unused!*/ errorHandler: T.ErrorHandler, sesameFn?: T.SesameFunction): T.HookInfo['generator'] => {
|
||||
const hookGenerator = (rpc: T.HookRPC<any, any>, /*not unused!*/ errorHandler: T.ErrorHandler, sesameFn?: T.SesameFunction): T.HookInfo['generator'] => {
|
||||
|
||||
let argsArr = extractArgs(rpc.hook)
|
||||
argsArr.pop() //remove 'callback' from the end
|
||||
let callArgs = argsArr.join(',')
|
||||
|
||||
const args = sesameFn?(['sesame', ...argsArr].join(','))
|
||||
:callArgs
|
||||
const args = sesameFn ? (['sesame', ...argsArr].join(','))
|
||||
: callArgs
|
||||
|
||||
callArgs = appendComma(callArgs, false)
|
||||
|
||||
@@ -139,35 +142,35 @@ const hookGenerator = (rpc:T.HookRPC<any, any>, /*not unused!*/ errorHandler: T.
|
||||
}
|
||||
|
||||
const makeError = (callName: string) => {
|
||||
return new Error("Call not found: "+callName+". ; Zone: <root> ; Task: Promise.then ; Value: Error: Call not found: "+callName)
|
||||
return new Error("Call not found: " + callName + ". ; Zone: <root> ; Task: Promise.then ; Value: Error: Call not found: " + callName)
|
||||
}
|
||||
|
||||
/**
|
||||
* Extract a string list of parameters from a function
|
||||
* @param f The source function
|
||||
*/
|
||||
const extractArgs = (f:Function):string[] => {
|
||||
let fn:string
|
||||
fn = (fn = String(f)).substr(0, fn.indexOf(")")).substr(fn.indexOf("(")+1)
|
||||
return fn!==""?fn.split(',') : []
|
||||
const extractArgs = (f: Function): string[] => {
|
||||
let fn: string
|
||||
fn = (fn = String(f)).substr(0, fn.indexOf(")")).substr(fn.indexOf("(") + 1)
|
||||
return fn !== "" ? fn.split(',') : []
|
||||
}
|
||||
|
||||
|
||||
export function makeSesameFunction (sesame : T.SesameFunction | string) : T.SesameFunction {
|
||||
if(typeof sesame === 'function'){
|
||||
export function makeSesameFunction(sesame: T.SesameFunction | string): T.SesameFunction {
|
||||
if (typeof sesame === 'function') {
|
||||
return sesame
|
||||
}
|
||||
|
||||
return (testSesame : string) => {
|
||||
return testSesame === sesame
|
||||
return (testSesame: string) => {
|
||||
return testSesame === sesame
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
export function appendComma(s?:string, turnToString = true):string{
|
||||
if(turnToString)
|
||||
return s?`'${s}',`:""
|
||||
return s?`${s},`:""
|
||||
export function appendComma(s?: string, turnToString = true): string {
|
||||
if (turnToString)
|
||||
return s ? `'${s}',` : ""
|
||||
return s ? `${s},` : ""
|
||||
}
|
||||
|
||||
|
||||
@@ -176,12 +179,94 @@ export function appendComma(s?:string, turnToString = true):string{
|
||||
* This was supposedly fixed (https://github.com/microsoft/TypeScript/issues/5611) but it still is the case.
|
||||
* This function sets the name value for all object members that are functions.
|
||||
*/
|
||||
export function fixNames(o:Object):void{
|
||||
export function fixNames(o: Object): void {
|
||||
Object.keys(o).forEach(key => {
|
||||
if(typeof o[key] === 'function' && !o[key].name){
|
||||
if (typeof o[key] === 'function' && !o[key].name) {
|
||||
Object.defineProperty(o[key], 'name', {
|
||||
value: key
|
||||
})
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
export const makePioSocket = (socket: any): I.Socket => {
|
||||
return {
|
||||
bind: (name: string, listener: T.PioBindListener) => socket.on(name, (...args: any) => {
|
||||
const ack = args.pop()
|
||||
listener.apply(null, args)
|
||||
ack()
|
||||
}),
|
||||
|
||||
hook: (name: string, listener: T.PioHookListener) => {
|
||||
const args = extractArgs(listener)
|
||||
let argNames
|
||||
let restParam = args.find(e => e.includes('...'))
|
||||
if(!restParam){
|
||||
argNames = [...args, '...__args__'].join(',')
|
||||
restParam = '__args__'
|
||||
}else{
|
||||
argNames = [...args].join(',')
|
||||
restParam = restParam.replace('...','')
|
||||
}
|
||||
|
||||
const decoratedListener = eval(`(() => async (${argNames}) => {
|
||||
const __ack__ = ${restParam}.pop()
|
||||
try{
|
||||
const response = await listener.apply(null, [${argNames}])
|
||||
__ack__(response)
|
||||
}catch(e){
|
||||
__ack__({
|
||||
...e,
|
||||
stack: e.stack,
|
||||
message: e.message,
|
||||
name: e.name,
|
||||
})
|
||||
}
|
||||
})()`)
|
||||
socket.on(name, decoratedListener)
|
||||
},
|
||||
|
||||
call: (name: string, ...args: any) => {
|
||||
return new Promise((res, rej) => {
|
||||
const params: any = [name, ...args, (resp) => {
|
||||
if(isError(resp)){
|
||||
const err = new Error()
|
||||
err.stack = resp.stack
|
||||
err.name = resp.name
|
||||
err.message = resp.message
|
||||
return rej(err)
|
||||
}
|
||||
res(resp)
|
||||
}]
|
||||
socket.emit.apply(socket, params)
|
||||
})
|
||||
},
|
||||
|
||||
fire: (name: string, ...args: any) => new Promise((res, rej) => {
|
||||
const params: any = [name, ...args]
|
||||
socket.emit.apply(socket, params)
|
||||
res()
|
||||
}),
|
||||
|
||||
unhook: (name: string, listener?: T.AnyFunction) => {
|
||||
if (listener) {
|
||||
socket.removeListener(name, listener)
|
||||
} else {
|
||||
socket.removeAllListeners(name)
|
||||
}
|
||||
},
|
||||
|
||||
id: socket.id,
|
||||
on: (...args) => socket.on.apply(socket, args),
|
||||
emit: (...args) => socket.emit.apply(socket, args),
|
||||
close: () => {
|
||||
socket
|
||||
socket.disconnect(true)
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
export const isError = function(e){
|
||||
return e && e.stack && e.message && typeof e.stack === 'string'
|
||||
&& typeof e.message === 'string';
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user