From 7885c8c8d77d039861d183461fc4c174ab18dfa8 Mon Sep 17 00:00:00 2001 From: Hide_D Date: Thu, 7 Mar 2024 16:58:03 +0000 Subject: [PATCH] test --- @sammo/server/src/worker_a.ts | 49 +++++++++++ @sammo/server/src/worker_b.ts | 40 +++++++++ @sammo/server/src/worker_common.ts | 5 ++ @sammo/server/src/worker_test.ts | 35 +++++++- @sammo/server/src/worker_util.ts | 136 +++++++++++++++++++++++++++++ 5 files changed, 264 insertions(+), 1 deletion(-) create mode 100644 @sammo/server/src/worker_a.ts create mode 100644 @sammo/server/src/worker_b.ts create mode 100644 @sammo/server/src/worker_common.ts create mode 100644 @sammo/server/src/worker_util.ts diff --git a/@sammo/server/src/worker_a.ts b/@sammo/server/src/worker_a.ts new file mode 100644 index 0000000..294d27f --- /dev/null +++ b/@sammo/server/src/worker_a.ts @@ -0,0 +1,49 @@ +import { delay, must } from "@sammo/util"; +import { parentPort as _parentPort } from "node:worker_threads"; +import { type MessagePort } from "node:worker_threads"; +import { generateMessage, PortRPCClient, type InitPort, type RandomMsg } from "./worker_util.js"; +import type { RPCSampleList } from "./worker_common.js"; + + +if (!_parentPort) { + throw new Error("parentPort is not defined"); +} + +const parentPort = must(_parentPort); +let run = true; + +async function worker_a(port: MessagePort) { + const rpcClient = new PortRPCClient(port); + + console.log("Hello from worker_a.ts!"); + + let cnt = 0; + + const now = Date.now(); + + while (run && cnt < 10000) { + const message = generateMessage(); + + await rpcClient.callFunction('do', message); + cnt++; + if(cnt % 100 === 0){ + console.log(`a_cnt: ${cnt}`); + } + } + + const elapsed = Date.now() - now; + console.log(`Elapsed: ${elapsed}, cnt: ${cnt}`); + + parentPort.postMessage(elapsed); +} + +await new Promise((resolve) => { + parentPort.once('message', async (ports: InitPort) => { + await worker_a(ports.oppose); + resolve(); + }); + + parentPort.once('close', () => { + run = false; + }); +}); diff --git a/@sammo/server/src/worker_b.ts b/@sammo/server/src/worker_b.ts new file mode 100644 index 0000000..77cbbac --- /dev/null +++ b/@sammo/server/src/worker_b.ts @@ -0,0 +1,40 @@ +import { delay, must } from "@sammo/util"; +import { parentPort as _parentPort } from "node:worker_threads"; +import { type MessagePort } from "node:worker_threads"; +import { PortRPCServer, type InitPort, type RandomMsg } from "./worker_util.js"; +import { unknown } from "zod"; + +if (!_parentPort) { + throw new Error("parentPort is not defined"); +} + +const parentPort = must(_parentPort); +let run = true; + +function doFunc(msg: RandomMsg): string { + return msg.key; +} + +let rpcServer: PortRPCServer<{}>|null = null; + +async function worker_b(port: MessagePort) { + + console.log("Hello from worker_b.ts!"); + + + rpcServer = new PortRPCServer(port, { + do: doFunc + }); + parentPort.postMessage("ready"); +} + +await new Promise((resolve) => { + parentPort.once('message', async (ports: InitPort) => { + await worker_b(ports.oppose); + resolve(); + }); + + parentPort.once('close', () => { + run = false; + }); +}); \ No newline at end of file diff --git a/@sammo/server/src/worker_common.ts b/@sammo/server/src/worker_common.ts new file mode 100644 index 0000000..dceddf0 --- /dev/null +++ b/@sammo/server/src/worker_common.ts @@ -0,0 +1,5 @@ +import type { RandomMsg } from "./worker_util.js"; + +export type RPCSampleList = { + do: (msg: RandomMsg) => string; +} \ No newline at end of file diff --git a/@sammo/server/src/worker_test.ts b/@sammo/server/src/worker_test.ts index cff9a41..1589dcd 100644 --- a/@sammo/server/src/worker_test.ts +++ b/@sammo/server/src/worker_test.ts @@ -1 +1,34 @@ -console.log("Hello from worker_test.ts!"); \ No newline at end of file +import { Worker, parentPort, MessageChannel } from "node:worker_threads"; +import * as nodePath from "node:path"; + +console.log("Hello from worker_test.ts!"); + +const __filename = new URL(import.meta.url).pathname; +const __dirname = nodePath.dirname(__filename); +const worker_b_ts = nodePath.join(__dirname, "worker_b.ts"); +const worker_a_ts = nodePath.join(__dirname, "worker_a.ts"); + + + + +const worker_b = new Worker(worker_b_ts); +const worker_a = new Worker(worker_a_ts); + +{ + const betweenPort = new MessageChannel(); + + worker_b.postMessage({ oppose: betweenPort.port2 }, [betweenPort.port2]); + worker_b.once('message', ()=>{ + worker_a.postMessage({ oppose: betweenPort.port1 }, [betweenPort.port1]); + }); + + +} + + +await new Promise((resolve) => { + worker_a.once('message', async (cnt: number) => { + await worker_b.terminate(); + resolve(); + }); +}); diff --git a/@sammo/server/src/worker_util.ts b/@sammo/server/src/worker_util.ts new file mode 100644 index 0000000..4011cee --- /dev/null +++ b/@sammo/server/src/worker_util.ts @@ -0,0 +1,136 @@ + +import { delay, must } from "@sammo/util"; +import { type MessagePort } from "node:worker_threads"; + +export function randomString(len: number): string { + const charSet = "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789"; + const charSetLen = charSet.length; + const str: string[] = []; + for (let i = 0; i < len; i++) { + str.push(charSet.charAt(Math.floor(Math.random() * charSetLen))); + } + return str.join(""); +} + +export type RandomMsg = { + key: string; + cnt: number; + arr: string[]; +}; + +export function ArrayBufferFromString(str: string): ArrayBuffer { + const buf = Buffer.from(str); + return buf.buffer.slice(buf.byteOffset, buf.byteOffset + buf.byteLength); +} + +export function generateMessage(): RandomMsg { + const randomObjCnt = Math.floor(Math.random() * 10); + const key = randomString(60); + const obj = { + key, + cnt: randomObjCnt, + arr: [] as string[], + }; + + for (let i = 0; i < randomObjCnt; i++) { + const randomStrLen = Math.floor(Math.random() * 20) + 4; + obj.arr.push(randomString(randomStrLen)); + } + + return obj; +} + +export type RawRPCCall = [name: string, id: number, arg: unknown]; + +export type RawRPCResponse = [id: number, success: 1 | 0, result: unknown]; + +export type RPCLists = { + [name: string]: (arg: any) => any | Promise; +} + +export class PortRPCServer{ + private listeners: Map unknown | Promise> = new Map(); + + constructor(private port: MessagePort, rpcLists: RPCL) { + port.on('message', async (buffer: ArrayBuffer) => { + if(!(buffer instanceof ArrayBuffer)){ + throw new Error("RPCServer: 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]); + } + const func = must(_func); + + 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: PortRPCServer started"); + } +} + +export class PortRPCClient { + private prevID = 0; + + private waiters = new Mapvoid, reject: (reason: unknown)=>void]>(); + + constructor(private port: MessagePort) { + 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: Response for unknown id: ${id}`); + return; + } + this.waiters.delete(id); + if (success) { + waiter[0](result); + } + else { + waiter[1](result); + } + }); + console.info("RPC: PortRPCClient started"); + } + + public async callFunction(name: T, arg: Parameters[0]): Promise> { + const id = this.prevID++; + + const waiter = new Promise((resolve, reject) => { + const call: RawRPCCall = [name, id, arg]; + const callBuffer = ArrayBufferFromString(JSON.stringify(call)); + + this.waiters.set(id, [resolve, reject]); + this.port.postMessage(callBuffer, [callBuffer]); + }); + + await delay(0); + + return waiter as Promise>; + } +} + +export type InitPort = { + oppose: MessagePort; +}; \ No newline at end of file