- Kotlin: TCP parallel test, WebSocket latency fix, stress test updates - TypeScript: TCP client, WebSocket client/server fixes, stress test, WS test updates - Archive completed subtasks (03_kotlin_tcp_parallel, 04_ws_fixed_latency) - Update roadmap milestone
540 lines
17 KiB
TypeScript
540 lines
17 KiB
TypeScript
import * as fs from "node:fs";
|
|
import * as net from "node:net";
|
|
import * as path from "node:path";
|
|
import * as tls from "node:tls";
|
|
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, toBinary, PacketBaseSchema } from "../src/packets/message_common_pb.js";
|
|
import { connectNodeWs, connectNodeWss, NodeWebSocket, NodeWsClient } from "../src/node_ws_client.js";
|
|
import { NodeWsServer } from "../src/node_ws_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: NodeWsClient): Promise<TestData> {
|
|
return new Promise<TestData>((resolve) => {
|
|
addListenerTyped(client.communicator, TestDataSchema, (msg) => resolve(msg));
|
|
});
|
|
}
|
|
|
|
describe("WS", () => {
|
|
const cleanup: Array<() => Promise<void>> = [];
|
|
|
|
afterEach(async () => {
|
|
while (cleanup.length > 0) {
|
|
const fn = cleanup.pop();
|
|
if (fn) {
|
|
await fn();
|
|
}
|
|
}
|
|
});
|
|
|
|
test("connectWs sends and receives TestData", async () => {
|
|
const server = new NodeWsServer("127.0.0.1", 0, "/", (ws) => new NodeWsClient(ws, 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 ws") {
|
|
await client.communicator.send(
|
|
create(TestDataSchema, { index: 200, message: "push from ws server" }),
|
|
);
|
|
}
|
|
});
|
|
};
|
|
|
|
const client = await connectNodeWs("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 ws" }));
|
|
|
|
await expect(pushed).resolves.toMatchObject({
|
|
index: 200,
|
|
message: "push from ws server",
|
|
});
|
|
});
|
|
|
|
test("sendRequest/response roundtrip over WS", async () => {
|
|
const server = new NodeWsServer("127.0.0.1", 0, "/", (ws) => new NodeWsClient(ws, 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 connectNodeWs("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 ws" }),
|
|
TestDataSchema,
|
|
500,
|
|
),
|
|
).resolves.toMatchObject({
|
|
index: 42,
|
|
message: "echo: request over ws",
|
|
});
|
|
});
|
|
|
|
test("connectWss sends and receives TestData", async () => {
|
|
const server = new NodeWsServer("127.0.0.1", 0, "/", (ws) => new NodeWsClient(ws, 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 wss") {
|
|
await client.communicator.send(
|
|
create(TestDataSchema, { index: 201, message: "push from wss server" }),
|
|
);
|
|
}
|
|
});
|
|
};
|
|
|
|
const client = await connectNodeWss(
|
|
"127.0.0.1",
|
|
server.port,
|
|
"/",
|
|
{ ca: cert },
|
|
0,
|
|
0,
|
|
parserMap(),
|
|
);
|
|
cleanup.push(async () => client.close());
|
|
|
|
const pushed = waitForClientMessage(client);
|
|
await client.communicator.send(create(TestDataSchema, { index: 12, message: "hello over wss" }));
|
|
|
|
await expect(pushed).resolves.toMatchObject({
|
|
index: 201,
|
|
message: "push from wss server",
|
|
});
|
|
});
|
|
|
|
test("WSS sendRequest/response roundtrip", async () => {
|
|
const server = new NodeWsServer("127.0.0.1", 0, "/", (ws) => new NodeWsClient(ws, 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: `wss echo: ${req.message}`,
|
|
}),
|
|
);
|
|
};
|
|
|
|
const client = await connectNodeWss(
|
|
"127.0.0.1",
|
|
server.port,
|
|
"/",
|
|
{ ca: cert },
|
|
0,
|
|
0,
|
|
parserMap(),
|
|
);
|
|
cleanup.push(async () => client.close());
|
|
|
|
await expect(
|
|
sendRequestTyped(
|
|
client.communicator,
|
|
create(TestDataSchema, { index: 22, message: "request over wss" }),
|
|
TestDataSchema,
|
|
500,
|
|
),
|
|
).resolves.toMatchObject({
|
|
index: 44,
|
|
message: "wss echo: request over wss",
|
|
});
|
|
});
|
|
|
|
test("concurrent sendRequest/response roundtrip over WS", async () => {
|
|
const server = new NodeWsServer("127.0.0.1", 0, "/", (ws) => new NodeWsClient(ws, 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 connectNodeWs("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(
|
|
"WS read loop drives server inbound gateway: concurrent request-response keeps FIFO/nonce",
|
|
async () => {
|
|
const server = new NodeWsServer("127.0.0.1", 0, "/", (ws) => new NodeWsClient(ws, 0, 0, parserMap()));
|
|
cleanup.push(async () => server.stop());
|
|
await server.start();
|
|
|
|
let lastNonce = 0;
|
|
let fifoViolation = false;
|
|
server.onClientConnected = (client) => {
|
|
const comm = client.communicator;
|
|
// 실제 WS read loop가 onReceivedFrame으로 흘린 raw binary message를 기본 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 connectNodeWs("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("connectNodeWss rejects when peer never completes WebSocket handshake", async () => {
|
|
const server = tls.createServer({ cert, key }, (socket) => {
|
|
socket.on("data", () => {});
|
|
});
|
|
cleanup.push(
|
|
() =>
|
|
new Promise<void>((resolve) => {
|
|
server.close(() => resolve());
|
|
}),
|
|
);
|
|
await new Promise<void>((resolve) => server.listen(0, "127.0.0.1", resolve));
|
|
const port = (server.address() as net.AddressInfo).port;
|
|
|
|
await expect(
|
|
connectNodeWss("127.0.0.1", port, "/", { ca: cert }, 0, 0, parserMap()),
|
|
).rejects.toThrow(/websocket handshake timed out/);
|
|
});
|
|
test("1MB message WS roundtrip", async () => {
|
|
const server = new NodeWsServer("127.0.0.1", 0, "/", (ws) => new NodeWsClient(ws, 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: req.message,
|
|
}),
|
|
);
|
|
};
|
|
|
|
const client = await connectNodeWs("127.0.0.1", server.port, "/", 0, 0, parserMap());
|
|
cleanup.push(async () => client.close());
|
|
|
|
const largeMessage = "A".repeat(1024 * 1024); // 1MB
|
|
const res = await sendRequestTyped(
|
|
client.communicator,
|
|
create(TestDataSchema, { index: 50, message: largeMessage }),
|
|
TestDataSchema,
|
|
10000,
|
|
);
|
|
|
|
expect(res.index).toBe(100);
|
|
expect(res.message.length).toBe(largeMessage.length);
|
|
expect(res.message).toBe(largeMessage);
|
|
}, 15000);
|
|
|
|
test("fragmented large binary frame regression test using raw net.Socket", async () => {
|
|
const { randomBytes } = await import("node:crypto");
|
|
const server = new NodeWsServer("127.0.0.1", 0, "/", (ws) => new NodeWsClient(ws, 0, 0, parserMap()));
|
|
cleanup.push(async () => server.stop());
|
|
await server.start();
|
|
|
|
const receivedPromise = new Promise<TestData>((resolve) => {
|
|
server.onClientConnected = (client) => {
|
|
addListenerTyped(client.communicator, TestDataSchema, (msg) => {
|
|
resolve(msg);
|
|
});
|
|
};
|
|
});
|
|
|
|
const socket = new net.Socket();
|
|
await new Promise<void>((resolve, reject) => {
|
|
socket.connect(server.port, "127.0.0.1", () => resolve());
|
|
socket.on("error", reject);
|
|
});
|
|
cleanup.push(async () => {
|
|
socket.destroy();
|
|
});
|
|
|
|
const key = randomBytes(16).toString("base64");
|
|
socket.write(
|
|
[
|
|
`GET / HTTP/1.1`,
|
|
`Host: 127.0.0.1:${server.port}`,
|
|
"Upgrade: websocket",
|
|
"Connection: Upgrade",
|
|
`Sec-WebSocket-Key: ${key}`,
|
|
"Sec-WebSocket-Version: 13",
|
|
"\r\n",
|
|
].join("\r\n"),
|
|
);
|
|
|
|
await new Promise<void>((resolve, reject) => {
|
|
let response = Buffer.alloc(0);
|
|
const onData = (chunk: Buffer) => {
|
|
response = Buffer.concat([response, chunk]);
|
|
if (response.indexOf("\r\n\r\n") >= 0) {
|
|
socket.off("data", onData);
|
|
socket.off("error", onError);
|
|
resolve();
|
|
}
|
|
};
|
|
const onError = (err: Error) => reject(err);
|
|
socket.on("data", onData);
|
|
socket.on("error", onError);
|
|
});
|
|
|
|
const largeMessage = "B".repeat(1024 * 1024); // 1MB
|
|
const testData = create(TestDataSchema, { index: 12345, message: largeMessage });
|
|
|
|
const encodedTestData = toBinary(TestDataSchema, testData);
|
|
|
|
const base = create(PacketBaseSchema, {
|
|
typeName: TestDataSchema.typeName,
|
|
nonce: 42,
|
|
data: encodedTestData,
|
|
});
|
|
|
|
const rawPayload = toBinary(PacketBaseSchema, base);
|
|
|
|
const payloadLength = rawPayload.byteLength;
|
|
const lengthBytes = payloadLength < 126 ? 0 : payloadLength <= 0xffff ? 2 : 8;
|
|
const header = Buffer.allocUnsafe(2 + lengthBytes + 4);
|
|
let offset = 0;
|
|
header[offset] = 0x80 | 0x2;
|
|
offset += 1;
|
|
if (payloadLength < 126) {
|
|
header[offset] = 0x80 | payloadLength;
|
|
offset += 1;
|
|
} else if (payloadLength <= 0xffff) {
|
|
header[offset] = 0x80 | 126;
|
|
offset += 1;
|
|
header.writeUInt16BE(payloadLength, offset);
|
|
offset += 2;
|
|
} else {
|
|
header[offset] = 0x80 | 127;
|
|
offset += 1;
|
|
header.writeBigUInt64BE(BigInt(payloadLength), offset);
|
|
offset += 8;
|
|
}
|
|
|
|
const maskingKey = randomBytes(4);
|
|
maskingKey.copy(header, offset);
|
|
|
|
const maskedPayload = Buffer.allocUnsafe(payloadLength);
|
|
for (let i = 0; i < payloadLength; i++) {
|
|
maskedPayload[i] = rawPayload[i] ^ maskingKey[i % 4];
|
|
}
|
|
|
|
const fullFrame = Buffer.concat([header, maskedPayload]);
|
|
|
|
const chunkSize = 8192;
|
|
for (let i = 0; i < fullFrame.byteLength; i += chunkSize) {
|
|
const chunk = fullFrame.subarray(i, i + chunkSize);
|
|
socket.write(chunk);
|
|
await new Promise((resolve) => setTimeout(resolve, 1));
|
|
}
|
|
|
|
const received = await receivedPromise;
|
|
expect(received.index).toBe(12345);
|
|
expect(received.message.length).toBe(largeMessage.length);
|
|
expect(received.message).toBe(largeMessage);
|
|
}, 15000);
|
|
|
|
test("empty WS control frame (ping/pong) regression test", async () => {
|
|
const { randomBytes } = await import("node:crypto");
|
|
const server = new NodeWsServer("127.0.0.1", 0, "/", (ws) => new NodeWsClient(ws, 0, 0, parserMap()));
|
|
cleanup.push(async () => server.stop());
|
|
await server.start();
|
|
|
|
const socket = new net.Socket();
|
|
await new Promise<void>((resolve, reject) => {
|
|
socket.connect(server.port, "127.0.0.1", () => resolve());
|
|
socket.on("error", reject);
|
|
});
|
|
cleanup.push(async () => {
|
|
socket.destroy();
|
|
});
|
|
|
|
const key = randomBytes(16).toString("base64");
|
|
socket.write(
|
|
[
|
|
`GET / HTTP/1.1`,
|
|
`Host: 127.0.0.1:${server.port}`,
|
|
"Upgrade: websocket",
|
|
"Connection: Upgrade",
|
|
`Sec-WebSocket-Key: ${key}`,
|
|
"Sec-WebSocket-Version: 13",
|
|
"\r\n",
|
|
].join("\r\n"),
|
|
);
|
|
|
|
// Wait for handshake response
|
|
await new Promise<void>((resolve, reject) => {
|
|
let response = Buffer.alloc(0);
|
|
const onData = (chunk: Buffer) => {
|
|
response = Buffer.concat([response, chunk]);
|
|
if (response.indexOf("\r\n\r\n") >= 0) {
|
|
socket.off("data", onData);
|
|
socket.off("error", onError);
|
|
resolve();
|
|
}
|
|
};
|
|
const onError = (err: Error) => reject(err);
|
|
socket.on("data", onData);
|
|
socket.on("error", onError);
|
|
});
|
|
|
|
// Send empty masked ping frame
|
|
const pingFrame = Buffer.alloc(6);
|
|
pingFrame[0] = 0x80 | 0x9; // FIN | PING
|
|
pingFrame[1] = 0x80 | 0; // MASKED | payload length 0
|
|
const maskingKey = randomBytes(4);
|
|
maskingKey.copy(pingFrame, 2);
|
|
|
|
const pongReceivedPromise = new Promise<Buffer>((resolve, reject) => {
|
|
socket.once("data", (data) => {
|
|
resolve(data);
|
|
});
|
|
socket.once("error", reject);
|
|
});
|
|
|
|
socket.write(pingFrame);
|
|
|
|
const pongData = await pongReceivedPromise;
|
|
// Expected pong: FIN | PONG (0x8A), payload length 0 (0x00)
|
|
expect(pongData[0]).toBe(0x80 | 0xA);
|
|
expect(pongData[1]).toBe(0x00);
|
|
}, 10000);
|
|
|
|
test("NodeWebSocket enables TCP no-delay on its socket", () => {
|
|
// Nagle's algorithm + delayed-ACK 조합으로 인한 ~40ms 고정 지연을 막기 위해
|
|
// NodeWebSocket이 생성 시 socket no-delay를 켜는지 fake socket으로 검증한다.
|
|
let noDelay: boolean | undefined;
|
|
const fakeSocket = {
|
|
setNoDelay(value: boolean) {
|
|
noDelay = value;
|
|
},
|
|
on() {
|
|
return this;
|
|
},
|
|
once() {
|
|
return this;
|
|
},
|
|
destroyed: false,
|
|
};
|
|
|
|
new NodeWebSocket(fakeSocket as unknown as net.Socket, true);
|
|
expect(noDelay).toBe(true);
|
|
});
|
|
|
|
test("WS masking tail bytes equivalence unit test for various lengths", async () => {
|
|
const lengths = [0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 15, 16, 17, 1023, 1024, 1025];
|
|
const server = new NodeWsServer("127.0.0.1", 0, "/", (ws) => new NodeWsClient(ws, 0, 0, parserMap()));
|
|
cleanup.push(async () => server.stop());
|
|
await server.start();
|
|
|
|
server.onClientConnected = (client) => {
|
|
addRequestListenerTyped(client.communicator, TestDataSchema, (req) =>
|
|
create(TestDataSchema, {
|
|
index: req.index,
|
|
message: req.message,
|
|
}),
|
|
);
|
|
};
|
|
|
|
const client = await connectNodeWs("127.0.0.1", server.port, "/", 0, 0, parserMap());
|
|
cleanup.push(async () => client.close());
|
|
|
|
for (const len of lengths) {
|
|
const msg = "A".repeat(len);
|
|
const res = await sendRequestTyped(
|
|
client.communicator,
|
|
create(TestDataSchema, { index: len, message: msg }),
|
|
TestDataSchema,
|
|
2000,
|
|
);
|
|
expect(res.index).toBe(len);
|
|
expect(res.message.length).toBe(len);
|
|
expect(res.message).toBe(msg);
|
|
}
|
|
});
|
|
});
|