2.0 rework
This commit is contained in:
+15
-20
@@ -11,12 +11,12 @@ export class RPCServer<
|
||||
> {
|
||||
|
||||
private pio = PromiseIO.createServer()
|
||||
private visibility: T.Visibility
|
||||
private closeHandler: T.CloseHandler
|
||||
private errorHandler: T.ErrorHandler
|
||||
private connectionHandler: T.ConnectionHandler
|
||||
private sesame?: T.SesameFunction
|
||||
private accessFilter: T.AccessFilter<InterfaceT>
|
||||
private attached = false
|
||||
|
||||
/**
|
||||
* @throws On RPC with no name
|
||||
@@ -25,13 +25,9 @@ export class RPCServer<
|
||||
* @param conf A {@link SocketConf} object with optional settings
|
||||
*/
|
||||
constructor(
|
||||
private port: number,
|
||||
private exporters: T.ExporterArray<InterfaceT> = [],
|
||||
private conf: T.ServerConf<InterfaceT> = {},
|
||||
private ws = http.createServer()
|
||||
conf: T.ServerConf<InterfaceT> = {},
|
||||
) {
|
||||
if (!conf.visibility) this.visibility = "0.0.0.0"
|
||||
|
||||
if (conf.sesame) {
|
||||
this.sesame = U.makeSesameFunction(conf.sesame)
|
||||
}
|
||||
@@ -57,7 +53,7 @@ export class RPCServer<
|
||||
|
||||
exporters.forEach(U.fixNames) //TSC for some reason doesn't preserve name properties of methods
|
||||
|
||||
let badRPC = exporters.flatMap(ex => ex.exportRPCs()).find(rpc => !rpc.name)
|
||||
let badRPC = exporters.flatMap(ex => typeof ex.RPCs === "function"?ex.RPCs():(ex as any)).find(rpc => !rpc.name)
|
||||
if (badRPC) {
|
||||
throw new Error(`
|
||||
RPC did not provide a name.
|
||||
@@ -67,32 +63,32 @@ export class RPCServer<
|
||||
\n`+ badRPC.toString() + `
|
||||
\n>------------OFFENDING RPC`)
|
||||
}
|
||||
|
||||
|
||||
this.startWebsocket()
|
||||
}
|
||||
|
||||
private startWebsocket() {
|
||||
try {
|
||||
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)
|
||||
})
|
||||
|
||||
if(this.conf.selfStart == null || this.conf.selfStart == true)
|
||||
this.pio.listen(this.port)
|
||||
|
||||
} catch (e) {
|
||||
this.errorHandler(this.pio, e, 'system', [])
|
||||
}
|
||||
}
|
||||
|
||||
public attach = (httpServer = new http.Server()) : RPCServer<InterfaceT> => {
|
||||
this.pio.attach(httpServer)
|
||||
this.attached = true
|
||||
return this
|
||||
}
|
||||
|
||||
public listen(port:number) : RPCServer<InterfaceT>{
|
||||
if(!this.attached) this.attach()
|
||||
this.pio.listen(port)
|
||||
return this
|
||||
}
|
||||
|
||||
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)
|
||||
@@ -105,6 +101,5 @@ export class RPCServer<
|
||||
|
||||
close(): void {
|
||||
this.pio.close()
|
||||
this.ws.close()
|
||||
}
|
||||
}
|
||||
+1
-2
@@ -9,13 +9,12 @@ export type RPCExporter<
|
||||
Name extends keyof Ifc = keyof Ifc,
|
||||
> = {
|
||||
name: Name
|
||||
exportRPCs() : T.RPCDefinitions<Ifc>[Name]
|
||||
RPCs : T.RPCDefinitions<Ifc>[Name] | (() => T.RPCDefinitions<Ifc>[Name])
|
||||
}
|
||||
|
||||
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>
|
||||
|
||||
@@ -6,6 +6,7 @@ const socketio = require('socket.io')
|
||||
|
||||
export class PromiseIO {
|
||||
io?: Server
|
||||
httpServer: httpServer
|
||||
private listeners: { [eventName in string]: ((...args: any) => void)[] } = {
|
||||
socket: [],
|
||||
connect: []
|
||||
@@ -16,10 +17,8 @@ export class PromiseIO {
|
||||
}
|
||||
|
||||
attach(httpServer: httpServer) {
|
||||
this.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))
|
||||
@@ -36,7 +35,10 @@ export class PromiseIO {
|
||||
pioSock.on('reconnecting', ()=>console.log('reconnecting'));
|
||||
*/
|
||||
})
|
||||
this.io!.listen(port)
|
||||
}
|
||||
|
||||
listen(port: number) {
|
||||
this.httpServer!.listen(port)
|
||||
}
|
||||
|
||||
on(eventName: string, listener: T.AnyFunction) {
|
||||
|
||||
@@ -28,12 +28,10 @@ export type ExporterArray<InterfaceT extends RPCInterface = RPCInterface> = I.RP
|
||||
export type ConnectedSocket<T extends RPCInterface = RPCInterface> = RPCSocket & AsyncIfc<T>
|
||||
|
||||
export type ServerConf<InterfaceT extends RPCInterface> = {
|
||||
selfStart?: boolean
|
||||
accessFilter?: AccessFilter<InterfaceT>
|
||||
connectionHandler?: ConnectionHandler
|
||||
errorHandler?: ErrorHandler
|
||||
closeHandler?: CloseHandler
|
||||
visibility?: Visibility
|
||||
} & SesameConf
|
||||
|
||||
export type SocketConf = {
|
||||
|
||||
+2
-2
@@ -61,7 +61,7 @@ RPC did not provide a name.
|
||||
*/
|
||||
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()
|
||||
const RPCs = typeof exporter.RPCs === "function" ? exporter.RPCs() : exporter.RPCs
|
||||
|
||||
return RPCs.map(rpc => rpcToRpcinfo(socket, rpc, owner, errorHandler, sesame))
|
||||
.map(info => {
|
||||
@@ -262,7 +262,7 @@ export const makePioSocket = (socket: any): I.Socket => {
|
||||
close: () => {
|
||||
socket
|
||||
socket.disconnect(true)
|
||||
},
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user