Compare commits
5
Commits
master
...
worker_test
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
2e23be1585 | ||
|
|
ca0bf296a6 | ||
|
|
c713b858c1 | ||
|
|
7885c8c8d7 | ||
|
|
a8ddaac1ed |
@@ -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",
|
||||
|
||||
@@ -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;
|
||||
});
|
||||
});
|
||||
@@ -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;
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,5 @@
|
||||
import type { RandomMsg } from "./worker_util.js";
|
||||
|
||||
export type RPCSampleList = {
|
||||
do: (msg: RandomMsg) => string;
|
||||
}
|
||||
@@ -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();
|
||||
});
|
||||
});
|
||||
@@ -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;
|
||||
};
|
||||
@@ -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);
|
||||
|
||||
Reference in New Issue
Block a user