- Refactor RowParser to use pooled buffers for batch row processing - Add parallel row processing with Web Worker offloading in browser - Optimize TypeScript Node.js gateway TCP/WS handlers - Update tests for communicator, tcp, ws with performance improvements
316 lines
10 KiB
TypeScript
316 lines
10 KiB
TypeScript
import * as fs from "node:fs";
|
|
import * as path from "node:path";
|
|
import { fileURLToPath } from "node:url";
|
|
|
|
import { afterEach, describe, expect, test } from "vitest";
|
|
|
|
import {
|
|
addListenerTyped,
|
|
addRequestListenerTyped,
|
|
parserFromSchema,
|
|
sendRequestTyped,
|
|
} from "../src/communicator.js";
|
|
|
|
import { create, TestDataSchema, type TestData } from "../src/packets/message_common_pb.js";
|
|
import { connectTcp, connectTcpTls, TcpClient } from "../src/tcp_client.js";
|
|
import { TcpServer } from "../src/tcp_server.js";
|
|
|
|
const certsDir = path.join(path.dirname(fileURLToPath(import.meta.url)), "certs");
|
|
const cert = fs.readFileSync(path.join(certsDir, "server.crt"));
|
|
const key = fs.readFileSync(path.join(certsDir, "server.key"));
|
|
|
|
function parserMap() {
|
|
return new Map([[TestDataSchema.typeName, parserFromSchema(TestDataSchema)]]);
|
|
}
|
|
|
|
async function waitForClientMessage(client: TcpClient): Promise<TestData> {
|
|
return new Promise<TestData>((resolve) => {
|
|
addListenerTyped(client.communicator, TestDataSchema, (msg) => resolve(msg));
|
|
});
|
|
}
|
|
|
|
describe("TCP", () => {
|
|
const cleanup: Array<() => Promise<void>> = [];
|
|
|
|
afterEach(async () => {
|
|
while (cleanup.length > 0) {
|
|
const fn = cleanup.pop();
|
|
if (fn) {
|
|
await fn();
|
|
}
|
|
}
|
|
});
|
|
|
|
test("connectTcp sends and receives TestData", async () => {
|
|
const server = new TcpServer("127.0.0.1", 0, (socket) => new TcpClient(socket, 0, 0, parserMap()));
|
|
cleanup.push(async () => server.stop());
|
|
await server.start();
|
|
|
|
server.onClientConnected = (client) => {
|
|
addListenerTyped(client.communicator, TestDataSchema, async (msg) => {
|
|
if (msg.index === 11 && msg.message === "hello over tcp") {
|
|
await client.communicator.send(
|
|
create(TestDataSchema, { index: 200, message: "push from tcp server" }),
|
|
);
|
|
}
|
|
});
|
|
};
|
|
|
|
const client = await connectTcp("127.0.0.1", server.port, 0, 0, parserMap());
|
|
cleanup.push(async () => client.close());
|
|
|
|
const pushed = waitForClientMessage(client);
|
|
await client.communicator.send(create(TestDataSchema, { index: 11, message: "hello over tcp" }));
|
|
|
|
await expect(pushed).resolves.toMatchObject({
|
|
index: 200,
|
|
message: "push from tcp server",
|
|
});
|
|
});
|
|
|
|
test("sendRequest/response roundtrip", async () => {
|
|
const server = new TcpServer("127.0.0.1", 0, (socket) => new TcpClient(socket, 0, 0, parserMap()));
|
|
cleanup.push(async () => server.stop());
|
|
await server.start();
|
|
|
|
server.onClientConnected = (client) => {
|
|
addRequestListenerTyped(client.communicator, TestDataSchema, (req) =>
|
|
create(TestDataSchema, {
|
|
index: req.index * 2,
|
|
message: `echo: ${req.message}`,
|
|
}),
|
|
);
|
|
};
|
|
|
|
const client = await connectTcp("127.0.0.1", server.port, 0, 0, parserMap());
|
|
cleanup.push(async () => client.close());
|
|
|
|
await expect(
|
|
sendRequestTyped(
|
|
client.communicator,
|
|
create(TestDataSchema, { index: 21, message: "request over tcp" }),
|
|
TestDataSchema,
|
|
500,
|
|
),
|
|
).resolves.toMatchObject({
|
|
index: 42,
|
|
message: "echo: request over tcp",
|
|
});
|
|
});
|
|
|
|
test("connectTcpTls sends and receives TestData", async () => {
|
|
const server = new TcpServer(
|
|
"127.0.0.1",
|
|
0,
|
|
(socket) => new TcpClient(socket, 0, 0, parserMap()),
|
|
{ cert, key },
|
|
);
|
|
cleanup.push(async () => server.stop());
|
|
await server.start();
|
|
|
|
server.onClientConnected = (client) => {
|
|
addListenerTyped(client.communicator, TestDataSchema, async (msg) => {
|
|
if (msg.index === 12 && msg.message === "hello over tls") {
|
|
await client.communicator.send(
|
|
create(TestDataSchema, { index: 201, message: "push from tls server" }),
|
|
);
|
|
}
|
|
});
|
|
};
|
|
|
|
const client = await connectTcpTls(
|
|
"127.0.0.1",
|
|
server.port,
|
|
{ ca: cert, servername: "localhost" },
|
|
0,
|
|
0,
|
|
parserMap(),
|
|
);
|
|
cleanup.push(async () => client.close());
|
|
|
|
const pushed = waitForClientMessage(client);
|
|
await client.communicator.send(create(TestDataSchema, { index: 12, message: "hello over tls" }));
|
|
|
|
await expect(pushed).resolves.toMatchObject({
|
|
index: 201,
|
|
message: "push from tls server",
|
|
});
|
|
});
|
|
|
|
test("TLS sendRequest/response roundtrip", async () => {
|
|
const server = new TcpServer(
|
|
"127.0.0.1",
|
|
0,
|
|
(socket) => new TcpClient(socket, 0, 0, parserMap()),
|
|
{ cert, key },
|
|
);
|
|
cleanup.push(async () => server.stop());
|
|
await server.start();
|
|
|
|
server.onClientConnected = (client) => {
|
|
addRequestListenerTyped(client.communicator, TestDataSchema, (req) =>
|
|
create(TestDataSchema, {
|
|
index: req.index * 2,
|
|
message: `tls echo: ${req.message}`,
|
|
}),
|
|
);
|
|
};
|
|
|
|
const client = await connectTcpTls(
|
|
"127.0.0.1",
|
|
server.port,
|
|
{ ca: cert, servername: "localhost" },
|
|
0,
|
|
0,
|
|
parserMap(),
|
|
);
|
|
cleanup.push(async () => client.close());
|
|
|
|
await expect(
|
|
sendRequestTyped(
|
|
client.communicator,
|
|
create(TestDataSchema, { index: 22, message: "request over tls" }),
|
|
TestDataSchema,
|
|
500,
|
|
),
|
|
).resolves.toMatchObject({
|
|
index: 44,
|
|
message: "tls echo: request over tls",
|
|
});
|
|
});
|
|
|
|
test("concurrent sendRequest/response roundtrip over TCP", async () => {
|
|
const server = new TcpServer("127.0.0.1", 0, (socket) => new TcpClient(socket, 0, 0, parserMap()));
|
|
cleanup.push(async () => server.stop());
|
|
await server.start();
|
|
|
|
server.onClientConnected = (client) => {
|
|
addRequestListenerTyped(client.communicator, TestDataSchema, (req) =>
|
|
create(TestDataSchema, {
|
|
index: req.index * 2,
|
|
message: `echo: ${req.message}`,
|
|
}),
|
|
);
|
|
};
|
|
|
|
const client = await connectTcp("127.0.0.1", server.port, 0, 0, parserMap());
|
|
cleanup.push(async () => client.close());
|
|
|
|
const results = await Promise.all(
|
|
Array.from({ length: 5 }, (_, i) =>
|
|
sendRequestTyped(
|
|
client.communicator,
|
|
create(TestDataSchema, { index: 30 + i, message: `request ${i}` }),
|
|
TestDataSchema,
|
|
2000,
|
|
),
|
|
),
|
|
);
|
|
|
|
results.forEach((res, i) => {
|
|
expect(res.index).toBe((30 + i) * 2);
|
|
expect(res.message).toBe(`echo: request ${i}`);
|
|
});
|
|
});
|
|
|
|
test(
|
|
"TCP read loop drives server inbound gateway: concurrent request-response keeps FIFO/nonce",
|
|
async () => {
|
|
const server = new TcpServer("127.0.0.1", 0, (socket) => new TcpClient(socket, 0, 0, parserMap()));
|
|
cleanup.push(async () => server.stop());
|
|
await server.start();
|
|
|
|
let lastNonce = 0;
|
|
let fifoViolation = false;
|
|
server.onClientConnected = (client) => {
|
|
const comm = client.communicator;
|
|
// 실제 TCP read loop가 onReceivedFrame으로 흘린 raw frame을 기본 receive gateway가
|
|
// decode하고, reorder 후 coordinator가 FIFO로 dispatch하는지 검증한다.
|
|
comm.enableInboundGateway({ workers: 2 });
|
|
// raw addRequestListener로 nonce를 받아 gateway reorder 이후에도 connection별 nonce가
|
|
// 단조 증가(FIFO)하는지 확인한다.
|
|
comm.addRequestListener(TestDataSchema.typeName, (msg, nonce) => {
|
|
if (nonce < lastNonce) {
|
|
fifoViolation = true;
|
|
}
|
|
lastNonce = nonce;
|
|
const req = msg as TestData;
|
|
return create(TestDataSchema, { index: req.index * 2, message: `gw echo: ${req.message}` });
|
|
});
|
|
};
|
|
|
|
const client = await connectTcp("127.0.0.1", server.port, 0, 0, parserMap());
|
|
cleanup.push(async () => client.close());
|
|
|
|
// 동시 요청이 기본 receive gateway(inline decode + reorder)를 거쳐도 nonce 상관관계가 유지되어야 한다.
|
|
// 많은 동시 요청이 몰려도 순차 처리가 안전하게 보장되는지 검증하기 위해 요청 timeout에 충분한 여유를 둔다.
|
|
const results = await Promise.all(
|
|
Array.from({ length: 16 }, (_, i) =>
|
|
sendRequestTyped(
|
|
client.communicator,
|
|
create(TestDataSchema, { index: i, message: `req ${i}` }),
|
|
TestDataSchema,
|
|
8000,
|
|
),
|
|
),
|
|
);
|
|
|
|
results.forEach((res, i) => {
|
|
expect(res.index).toBe(i * 2);
|
|
expect(res.message).toBe(`gw echo: req ${i}`);
|
|
});
|
|
expect(fifoViolation).toBe(false);
|
|
},
|
|
15000,
|
|
);
|
|
|
|
test("TCP inbound backpressure pauses socket read on server", async () => {
|
|
let serverClient: TcpClient | null = null;
|
|
const server = new TcpServer("127.0.0.1", 0, (socket) => {
|
|
serverClient = new TcpClient(socket, 0, 0, parserMap());
|
|
return serverClient;
|
|
});
|
|
cleanup.push(async () => server.stop());
|
|
await server.start();
|
|
|
|
let resolveHandler: (() => void) | null = null;
|
|
const handlerPromise = new Promise<void>((r) => {
|
|
resolveHandler = r;
|
|
});
|
|
|
|
server.onClientConnected = (client) => {
|
|
addRequestListenerTyped(client.communicator, TestDataSchema, async () => {
|
|
await handlerPromise;
|
|
return create(TestDataSchema, { index: 100, message: "done" });
|
|
});
|
|
};
|
|
|
|
const client = await connectTcp("127.0.0.1", server.port, 0, 0, parserMap());
|
|
cleanup.push(async () => client.close());
|
|
|
|
const firstReq = sendRequestTyped(
|
|
client.communicator,
|
|
create(TestDataSchema, { index: 1, message: "first" }),
|
|
TestDataSchema,
|
|
5000,
|
|
);
|
|
|
|
await new Promise<void>((r) => setTimeout(r, 20));
|
|
|
|
for (let i = 2; i <= 65; i++) {
|
|
await client.communicator.send(create(TestDataSchema, { index: i, message: "heavy" }));
|
|
}
|
|
|
|
await client.communicator.send(create(TestDataSchema, { index: 66, message: "overflow" }));
|
|
await new Promise<void>((r) => setTimeout(r, 100));
|
|
|
|
expect(serverClient).not.toBeNull();
|
|
expect(((serverClient as any).socket as any).isPaused()).toBe(true);
|
|
|
|
resolveHandler!();
|
|
await firstReq;
|
|
|
|
await new Promise<void>((r) => setTimeout(r, 50));
|
|
expect(((serverClient as any).socket as any).isPaused()).toBe(false);
|
|
});
|
|
});
|