feat: portRPC
- port를 이용해서 worker_thread간 RPC 호출 가능 - 호출, 호출반환은 ArrayBuffer를 이용 - 호출당 100us
This commit is contained in:
@@ -1,2 +1,3 @@
|
||||
export { UniqueNumberAllocator, type StateIncrementer } from './UniqueNumberAllocator.js';
|
||||
export * from './StartSession.js';
|
||||
export * from './StartSession.js';
|
||||
export * from './portRPC.js';
|
||||
@@ -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 extends PlainJson, RetVal extends PlainJson> = (arg: Arg) => RetVal | Promise<RetVal>;
|
||||
|
||||
export type RPCLists = {
|
||||
[name: string]: RPCItem<any, any>;
|
||||
}
|
||||
|
||||
export class RPCServer<RPCL extends RPCLists>{
|
||||
private listeners: Map<string, (arg: unknown) => unknown | Promise<unknown>> = 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<RPCL extends RPCLists> {
|
||||
private prevID = 0;
|
||||
|
||||
private waiters = new Map<number, [resolve: (result: unknown)=>void, 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<T extends keyof RPCL & string>(name: T, arg: Parameters<RPCL[T]>[0]): Promise<Awaited<ReturnType<RPCL[T]>>> {
|
||||
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<Awaited<ReturnType<RPCL[T]>>>;
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user