test
This commit is contained in:
@@ -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;
|
||||||
|
}
|
||||||
@@ -1 +1,34 @@
|
|||||||
console.log("Hello from worker_test.ts!");
|
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<void>((resolve) => {
|
||||||
|
worker_a.once('message', async (cnt: number) => {
|
||||||
|
await worker_b.terminate();
|
||||||
|
resolve();
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|||||||
@@ -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<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++;
|
||||||
|
|
||||||
|
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<ReturnType<RPCL[T]>>;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
export type InitPort = {
|
||||||
|
oppose: MessagePort;
|
||||||
|
};
|
||||||
Reference in New Issue
Block a user