diff --git a/@sammo/server_util/src/index.ts b/@sammo/server_util/src/index.ts index 8c59361..c16bb97 100644 --- a/@sammo/server_util/src/index.ts +++ b/@sammo/server_util/src/index.ts @@ -1,2 +1,3 @@ export { UniqueNumberAllocator, type StateIncrementer } from './UniqueNumberAllocator.js'; -export * from './StartSession.js'; \ No newline at end of file +export * from './StartSession.js'; +export * from './portRPC.js'; \ No newline at end of file diff --git a/@sammo/server_util/src/portRPC.ts b/@sammo/server_util/src/portRPC.ts new file mode 100644 index 0000000..346392e --- /dev/null +++ b/@sammo/server_util/src/portRPC.ts @@ -0,0 +1,111 @@ +import { ArrayBufferFromString, delay, type PlainJson } from "@sammo/util"; +import { randomUUID } from "node:crypto"; +import { type MessagePort } from "node:worker_threads"; + + +type RawRPCCall = [name: string, id: number, arg: unknown]; +type RawRPCResponse = [id: number, success: 1 | 0, result: unknown]; + +export type RPCItem = (arg: Arg) => RetVal | Promise; + +export type RPCLists = { + [name: string]: RPCItem; +} + +export class RPCServer{ + private listeners: Map unknown | Promise> = new Map(); + + constructor(private port: MessagePort, rpcLists: RPCL, public readonly portName: string = randomUUID()) { + port.on('message', async (buffer: ArrayBuffer) => { + if(!(buffer instanceof ArrayBuffer)){ + throw new Error(`RPC(${this.portName}): buffer is not ArrayBuffer`); + } + const rawText = Buffer.from(buffer).toString(); + const data = JSON.parse(rawText) as RawRPCCall; + const [name, id, arg] = data; + const func = this.listeners.get(name); + + if (!func) { + const result: RawRPCResponse = [id, 0, `Function ${name} not found`]; + const resultBuffer = ArrayBufferFromString(JSON.stringify(result)); + port.postMessage(resultBuffer, [resultBuffer]); + return; + } + + try { + const result: RawRPCResponse = [id, 1, await func(arg)]; + const resultBuffer = ArrayBufferFromString(JSON.stringify(result)); + port.postMessage(resultBuffer, [resultBuffer]); + } + catch (e) { + const errMsg = e instanceof Error ? e.message : e; + const result: RawRPCResponse = [id, 0, String(errMsg)]; + const resultBuffer = ArrayBufferFromString(JSON.stringify(result)); + port.postMessage(resultBuffer, [resultBuffer]); + } + }); + + for (const [name, func] of Object.entries(rpcLists)) { + this.listeners.set(name, func); + } + + console.info(`RPC(${this.portName}): RPCServer started`); + } +} + +export class RPCClient { + private prevID = 0; + + private waiters = new Mapvoid, reject: (reason: unknown)=>void]>(); + + constructor(private port: MessagePort, public readonly portName: string = randomUUID()) { + port.on('message', (buffer: ArrayBuffer) => { + const data = JSON.parse(Buffer.from(buffer).toString()) as RawRPCResponse; + const [id, success, result] = data; + const waiter = this.waiters.get(id); + if (!waiter) { + console.error(`RPC(${this.portName}): Response for unknown id: ${id}`); + return; + } + this.waiters.delete(id); + if (success) { + waiter[0](result); + } + else { + waiter[1](result); + } + }); + + console.info(`RPC(${this.portName}): RPCClient started`); + } + + public async callFunction(name: T, arg: Parameters[0]): Promise>> { + if(this.prevID >= Number.MAX_SAFE_INTEGER){ + this.prevID = 0; + } + const id = this.prevID++; + + if(this.waiters.has(id)){ + throw new Error(`RPC(${this.portName}): id(${id}) is already used`); + } + + let done = false; + + const waiter = new Promise((resolve, reject) => { + const call: RawRPCCall = [name, id, arg]; + const callBuffer = ArrayBufferFromString(JSON.stringify(call)); + + this.waiters.set(id, [resolve, reject]); + done = true; + this.port.postMessage(callBuffer, [callBuffer]); + }); + + await delay(0); + + if(!done){ + throw new Error(`RPC(${this.portName}): callFunction failed(${name})`); + } + + return waiter as Promise>>; + } +} \ No newline at end of file