5 Commits
Author SHA1 Message Date
Hide_D 2e23be1585 확장자 제거 2024-03-08 11:35:19 +00:00
Hide_D ca0bf296a6 simple test 2024-03-08 11:24:35 +09:00
Hide_D c713b858c1 delay 0 2024-03-08 11:08:05 +09:00
Hide_D 7885c8c8d7 test 2024-03-07 16:58:03 +00:00
Hide_D a8ddaac1ed worker_test init 2024-03-07 15:06:49 +00:00
7 changed files with 274 additions and 1 deletions
+2 -1
View File
@@ -5,7 +5,8 @@
"main": "dist/index.js",
"scripts": {
"build": "tsc --build",
"dev": "tsc --watch"
"dev": "tsc --watch",
"worker_test": "tsx src/worker_test.ts"
},
"author": "Hide_D <hided62@gmail.com>",
"type": "module",
+49
View File
@@ -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<RPCSampleList>(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<void>((resolve) => {
parentPort.once('message', async (ports: InitPort) => {
await worker_a(ports.oppose);
resolve();
});
parentPort.once('close', () => {
run = false;
});
});
+40
View File
@@ -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<void>((resolve) => {
parentPort.once('message', async (ports: InitPort) => {
await worker_b(ports.oppose);
resolve();
});
parentPort.once('close', () => {
run = false;
});
});
+5
View File
@@ -0,0 +1,5 @@
import type { RandomMsg } from "./worker_util.js";
export type RPCSampleList = {
do: (msg: RandomMsg) => string;
}
+32
View File
@@ -0,0 +1,32 @@
import { Worker, MessageChannel } from "node:worker_threads";
import * as nodePath from "node:path";
import { fileURLToPath } from "node:url";
console.log("Hello from worker_test.ts!");
const __filename = fileURLToPath(new URL(import.meta.url));
const __dirname = nodePath.dirname(__filename);
const worker_b_ts = nodePath.join(__dirname, "worker_b");
const worker_a_ts = nodePath.join(__dirname, "worker_a");
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<void>((resolve) => {
worker_a.once('message', async (cnt: number) => {
await worker_b.terminate();
resolve();
});
});
+143
View File
@@ -0,0 +1,143 @@
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<any>;
}
export class PortRPCServer<RPCL extends RPCLists>{
private listeners: Map<string, (arg: unknown) => unknown | Promise<unknown>> = 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<RPCL extends RPCLists> {
private prevID = 0;
private waiters = new Map<number, [resolve: (result: unknown)=>void, 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<T extends keyof RPCL & string>(name: T, arg: Parameters<RPCL[T]>[0]): Promise<ReturnType<RPCL[T]>> {
const id = this.prevID++;
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: callFunction failed");
}
return waiter as Promise<ReturnType<RPCL[T]>>;
}
}
export type InitPort = {
oppose: MessagePort;
};
+3
View File
@@ -20,6 +20,9 @@ export function calcBase64Len(length: number) {
export function delay(time: number): Promise<void>;
export function delay<T>(time: number, result: T): Promise<T>;
export function delay<T = undefined>(time: number, result?: T): Promise<T | void> {
if(time <= 0){
return new Promise(r => r(result));
}
return new Promise(resolve =>
setTimeout(() => {
resolve(result);