/** * Inbound queue ordering stress / benchmark harness (same-language TypeScript baseline). * * Milestone `inbound-queue-ordering`의 verify Epic 기준선을 만든다. 절대 성능 합격선을 고정하지 않고 * local 환경 baseline(throughput, p50/p95/p99 latency, peak memory)을 남기며, 안정성 합격선만 hard fail로 * 둔다: nonce mismatch 0, response type mismatch 0, unexpected timeout 0, connection별 FIFO 위반 0, * 종료 후 pending/queue leak 0. * * 출력 계약(stdout): 각 측정 축마다 한 줄의 `ROW|...` 레코드를 내고, 마지막에 `SUMMARY|...` 한 줄을 낸다. * 사람이 읽는 진행 로그는 stderr로 보낸다. 안정성 위반이 하나라도 있으면 process exit code 1. * * 측정 축(profile): * - roundtrip : 동일 connection request-response, concurrency 1/16/64/256, latency 분포. * - burst : 단일 connection fire-and-forget burst, dispatch 순서/누락/leak. * - sustained : 일정 시간 지속 request-response, throughput/p95/p99/peak memory. * - parallel : 다중 connection 동시 request-response, connection별 FIFO/nonce 독립성. * - slow-mix : 빠른 handler와 느린 handler가 섞인 다중 connection FIFO/nonce 독립성. * - gateway : gateway receive path를 세 mode로 분리한 baseline. * off = gateway 미설정 inline control(legacy 대조군). * on = default coordinator gateway(worker_threads hop 없는 inline decode + reorder). * worker_threads = explicit Node `worker_threads` off-thread decode 실험 row. * frame-ingest in-process row(`onReceivedFrame` 직접 주입)에 더해, 실제 TCP/WS * read loop를 통과하는 transport row(`tcp`/`ws`)를 위 세 mode로 측정한다. 이렇게 * default coordinator overhead와 worker_threads overhead를 한 record에서 분리한다. */ import { spawn, type ChildProcess } from "node:child_process"; import * as net from "node:net"; import { availableParallelism } from "node:os"; import * as readline from "node:readline"; import { addListenerTyped, Communicator, parserFromSchema, type ParserMap, type Transport, } from "../src/communicator.js"; import { WorkerGateway, type DecodedEnvelope, type InboundGateway } from "../src/inbound_gateway.js"; import { createNodeWorkerGateway } from "../src/node_inbound_gateway.js"; import { create, PacketBaseSchema, TestDataSchema, toBinary, type TestData, } from "../src/packets/message_common_pb.js"; import { connectNodeWs, NodeWsClient } from "../src/node_ws_client.js"; import { NodeWsServer } from "../src/node_ws_server.js"; import { connectTcp, TcpClient } from "../src/tcp_client.js"; import { TcpServer } from "../src/tcp_server.js"; /** stress harness 전용 TCP client factory. heartbeat 비활성(interval=0)으로 측정 noise를 줄인다. */ function newTcpClient(socket: net.Socket): TcpClient { return new TcpClient(socket, 0, 0, parserMap()); } const HOST = "127.0.0.1"; const WS_PATH = "/"; const ALL_PROFILES = ["roundtrip", "burst", "sustained", "parallel", "slow-mix", "payload", "gateway"] as const; type Profile = (typeof ALL_PROFILES)[number]; type Mode = "quick" | "full"; /** same-language transport baseline 축. gateway profile은 transport-agnostic frame-ingest 경로라 별도다. */ const ALL_TRANSPORTS = ["tcp", "ws"] as const; type WireTransport = (typeof ALL_TRANSPORTS)[number]; const TRANSPORT_PROFILES: ReadonlySet = new Set(["roundtrip", "burst", "sustained", "parallel", "slow-mix", "payload"]); const SERVER_CHILD_PREFIX = "SERVER_CHILD"; /** payload size matrix 축. quick은 작은 샘플(1KB/64KB), full은 Milestone 기준(1KB/64KB/1MB)을 측정한다. */ const PAYLOAD_SIZES_QUICK = [1_024, 65_536] as const; const PAYLOAD_SIZES_FULL = [1_024, 65_536, 1_048_576] as const; /** deterministic filler seed. message field를 이 16자 패턴으로 채워 재현 가능한 payload를 만든다. */ const PAYLOAD_SEED = "0123456789abcdef"; /** 안정성 위반 카운터. 0이 아니면 해당 축은 FAIL이고 프로세스도 non-zero exit한다. */ interface Stability { timeouts: number; nonceMismatch: number; typeMismatch: number; fifoViolations: number; pendingLeak: number; } interface Row { profile: Profile; axis: string; language: string; transport: string; payloadBytes: number; clientCount: number; requests: number; throughputRps: number; p50: number; p95: number; p99: number; queueBacklog: number; gatewayBacklog: number; gateway: string; memMb: number; stability: Stability; } const rows: Row[] = []; const LANGUAGE = "TypeScript"; const TRANSPORT_TCP = "tcp"; const TRANSPORT_WS = "ws"; const TRANSPORT_FRAME_INGEST = "frame-ingest"; /** 측정 축에서 server/client 생성만 transport별로 분기한다. profile 로직은 transport-agnostic하다. */ type BenchServer = TcpServer | NodeWsServer; type BenchClient = TcpClient | NodeWsClient; function makeServer(transport: WireTransport): BenchServer { if (transport === "ws") { return new NodeWsServer(HOST, 0, WS_PATH, (ws) => new NodeWsClient(ws, 0, 0, parserMap())); } return new TcpServer(HOST, 0, newTcpClient); } async function connectClient(transport: WireTransport, port: number): Promise { if (transport === "ws") { return connectNodeWs(HOST, port, WS_PATH, 0, 0, parserMap()); } return connectTcp(HOST, port, 0, 0, parserMap()); } function transportName(transport: WireTransport): string { return transport === "ws" ? TRANSPORT_WS : TRANSPORT_TCP; } /** * gateway cleanup 계측: queued/reorder/sink/worker pending 합을 backlog로 본다. in-flight decode loop가 * close 후 잠시 drain되는 `active`는 leak 판정 대상이 아니므로 제외한다. stats 미구현 gateway는 0이다. */ function gatewayBacklogOf(gateway: InboundGateway): number { const s = gateway.stats?.(); if (s === undefined) { return 0; } return s.queued + s.reorderPending + s.sinkPending + s.workerPending; } function newStability(): Stability { return { timeouts: 0, nonceMismatch: 0, typeMismatch: 0, fifoViolations: 0, pendingLeak: 0 }; } function stabilityViolations(s: Stability): number { return s.timeouts + s.nonceMismatch + s.typeMismatch + s.fifoViolations + s.pendingLeak; } function rowStatus(row: Row): "PASS" | "FAIL" { return stabilityViolations(row.stability) === 0 ? "PASS" : "FAIL"; } function emitRow(row: Row): void { rows.push(row); const s = row.stability; process.stdout.write( [ "ROW", row.profile, row.axis, row.language, row.transport, String(row.payloadBytes), String(row.clientCount), String(row.requests), row.throughputRps.toFixed(1), row.p50.toFixed(3), row.p95.toFixed(3), row.p99.toFixed(3), String(s.timeouts), String(s.nonceMismatch), String(s.typeMismatch), String(s.fifoViolations), String(s.pendingLeak), String(row.queueBacklog), String(row.gatewayBacklog), row.gateway, row.memMb.toFixed(1), rowStatus(row), ].join("|") + "\n", ); // fixed-latency 진단용 stderr diagnostic. stdout ROW 계약은 그대로 두고, transport/gateway mode별 // p50/p99를 grep 가능한 형태로 남겨 ws fixed latency 분석에서 row를 쉽게 추출한다. log( `DIAG fixed-latency profile=${row.profile} transport=${row.transport} gateway=${row.gateway} ` + `p50=${row.p50.toFixed(3)} p99=${row.p99.toFixed(3)}`, ); } function log(line: string): void { process.stderr.write(line + "\n"); } function parserMap(): ParserMap { return new Map([[TestDataSchema.typeName, parserFromSchema(TestDataSchema)]]); } function percentile(sorted: number[], p: number): number { if (sorted.length === 0) { return 0; } const rank = Math.ceil((p / 100) * sorted.length) - 1; const idx = Math.min(sorted.length - 1, Math.max(0, rank)); return sorted[idx]; } function memMb(): number { return process.memoryUsage().rss / (1024 * 1024); } function testDataPayloadBytes(message: string): number { return toBinary(TestDataSchema, create(TestDataSchema, { index: 0, message })).byteLength; } /** size 바이트 길이의 deterministic filler string을 만든다. payload matrix의 message field를 채운다. */ function payloadFiller(size: number): string { if (size <= 0) { return ""; } return PAYLOAD_SEED.repeat(Math.ceil(size / PAYLOAD_SEED.length)).slice(0, size); } /** payload 크기를 사람이 읽는 축 이름으로 바꾼다(1024 -> 1KB, 1048576 -> 1MB). */ function payloadLabel(bytes: number): string { if (bytes >= 1_048_576 && bytes % 1_048_576 === 0) { return `${bytes / 1_048_576}MB`; } if (bytes >= 1_024 && bytes % 1_024 === 0) { return `${bytes / 1_024}KB`; } return `${bytes}B`; } function sleep(ms: number): Promise { return new Promise((resolve) => setTimeout(resolve, ms)); } function errorMessage(err: unknown): string { return err instanceof Error ? err.message : String(err); } /** request 실패를 안정성 카운터의 적절한 항목으로 분류한다. */ function classifyRequestError(message: string, stability: Stability): void { if (message.includes("timeout")) { stability.timeouts += 1; } else if (message.includes("type mismatch")) { stability.typeMismatch += 1; } else { stability.nonceMismatch += 1; } } /** request-response echo + per-connection FIFO 감시를 붙인 서버를 띄우고 body를 실행한다. */ async function withRequestServer( transport: WireTransport, stability: Stability, body: (port: number) => Promise, ): Promise { const server = makeServer(transport); server.onClientConnected = (client: { communicator: Communicator }) => { addEchoRequestHandler(client, () => { stability.fifoViolations += 1; }); }; await server.start(); try { await body(server.port); } finally { await server.stop(); } } function addEchoRequestHandler( client: { communicator: Communicator }, onFifoViolation: () => void, ): void { let lastNonce = 0; // raw addRequestListener로 nonce를 받아 connection별 FIFO(단조 증가) 위반을 감시한다. client.communicator.addRequestListener(TestDataSchema.typeName, (msg, nonce) => { if (nonce < lastNonce) { onFifoViolation(); } lastNonce = nonce; const data = msg as TestData; return create(TestDataSchema, { index: data.index * 2, message: `echo:${data.message}` }); }); } async function withSplitProcessRequestServer( transport: WireTransport, stability: Stability, body: (port: number) => Promise, ): Promise { const child = spawn("./node_modules/.bin/tsx", ["bench/stress.ts", `--server-child=${transport}`], { cwd: process.cwd(), stdio: ["ignore", "pipe", "pipe"], }); const lines = readline.createInterface({ input: child.stdout }); let ready = false; const readyPort = new Promise((resolve, reject) => { const rejectIfNotReady = (err: Error): void => { if (!ready) { reject(err); } }; const readyTimer = setTimeout(() => { rejectIfNotReady(new Error("server child ready timeout")); }, 10_000); lines.on("line", (line) => { if (line.startsWith(`${SERVER_CHILD_PREFIX}|ready|`)) { ready = true; clearTimeout(readyTimer); resolve(Number(line.split("|")[2])); return; } if (line === `${SERVER_CHILD_PREFIX}|fifo-violation`) { stability.fifoViolations += 1; return; } log(`[parallel-child:${transport}] ${line}`); }); child.once("error", (err) => { clearTimeout(readyTimer); rejectIfNotReady(err); }); child.once("exit", (code, signal) => { clearTimeout(readyTimer); rejectIfNotReady(new Error(`server child exited before ready code=${code ?? "null"} signal=${signal ?? "null"}`)); }); }); child.stderr.on("data", (chunk: Buffer) => { process.stderr.write(`[parallel-child:${transport}] ${chunk.toString()}`); }); try { const port = await readyPort; await body(port); } finally { lines.close(); await stopServerChild(child); } } async function stopServerChild(child: ChildProcess): Promise { if (child.exitCode !== null || child.signalCode !== null) { return; } const exited = new Promise((resolve) => { child.once("exit", () => resolve()); }); child.kill("SIGTERM"); await Promise.race([ exited, sleep(2_000).then(() => { if (child.exitCode === null && child.signalCode === null) { child.kill("SIGKILL"); } }), ]); } async function runServerChild(transport: WireTransport): Promise { const server = makeServer(transport); server.onClientConnected = (client: { communicator: Communicator }) => { addEchoRequestHandler(client, () => { process.stdout.write(`${SERVER_CHILD_PREFIX}|fifo-violation\n`); }); }; await server.start(); process.stdout.write(`${SERVER_CHILD_PREFIX}|ready|${server.port}\n`); let stopped = false; const stop = async (): Promise => { if (stopped) { return; } stopped = true; await server.stop(); }; process.once("SIGTERM", () => { void stop().finally(() => process.exit(0)); }); process.once("SIGINT", () => { void stop().finally(() => process.exit(0)); }); await new Promise(() => {}); } /** 한 connection에서 concurrency C로 request batch를 보내고 latency를 모은다. */ async function runRequestLoad( client: BenchClient, total: number, concurrency: number, timeoutMs: number, stability: Stability, ): Promise { const latencies: number[] = []; let issued = 0; let nextIndex = 0; async function oneRequest(index: number): Promise { const started = performance.now(); try { const res = await client.communicator.sendRequest( create(TestDataSchema, { index, message: `req-${index}` }), TestDataSchema, timeoutMs, ); const elapsed = performance.now() - started; latencies.push(elapsed); if (res.index !== index * 2 || res.message !== `echo:req-${index}`) { stability.nonceMismatch += 1; } } catch (err) { classifyRequestError(errorMessage(err), stability); } } while (issued < total) { const batchSize = Math.min(concurrency, total - issued); const batch: Array> = []; for (let i = 0; i < batchSize; i += 1) { batch.push(oneRequest(nextIndex)); nextIndex += 1; } issued += batchSize; await Promise.all(batch); } return latencies; } /** latency 분포가 없는 throughput 중심 축(burst/gateway)을 위한 row emit. */ function emitThroughputRow( profile: Profile, axis: string, count: number, elapsedMs: number, stability: Stability, gateway: string, observedMemMb: number, metadata: { transport: string; payloadBytes: number; clientCount: number; queueBacklog?: number; gatewayBacklog?: number; }, ): void { emitRow({ profile, axis, language: LANGUAGE, transport: metadata.transport, payloadBytes: metadata.payloadBytes, clientCount: metadata.clientCount, requests: count, throughputRps: elapsedMs > 0 ? (count / elapsedMs) * 1000 : 0, p50: 0, p95: 0, p99: 0, queueBacklog: metadata.queueBacklog ?? 0, gatewayBacklog: metadata.gatewayBacklog ?? 0, gateway, memMb: observedMemMb, stability, }); } function summarizeLatencies( profile: Profile, axis: string, latencies: number[], elapsedMs: number, stability: Stability, gateway: string, observedMemMb: number, metadata: { transport: string; payloadBytes: number; clientCount: number; queueBacklog?: number; gatewayBacklog?: number; }, ): void { const sorted = [...latencies].sort((a, b) => a - b); const throughput = elapsedMs > 0 ? (latencies.length / elapsedMs) * 1000 : 0; emitRow({ profile, axis, language: LANGUAGE, transport: metadata.transport, payloadBytes: metadata.payloadBytes, clientCount: metadata.clientCount, requests: latencies.length, throughputRps: throughput, p50: percentile(sorted, 50), p95: percentile(sorted, 95), p99: percentile(sorted, 99), queueBacklog: metadata.queueBacklog ?? 0, gatewayBacklog: metadata.gatewayBacklog ?? 0, gateway, memMb: observedMemMb, stability, }); } async function profileRoundtrip(transport: WireTransport, mode: Mode): Promise { const concurrencies = [1, 16, 64, 256]; const batches = mode === "quick" ? 2 : 20; const timeoutMs = mode === "quick" ? 5_000 : 15_000; log(`[roundtrip] transport=${transport} mode=${mode} concurrencies=${concurrencies.join(",")} batches/level=${batches}`); for (const concurrency of concurrencies) { const stability = newStability(); const total = concurrency * batches; await withRequestServer(transport, stability, async (port) => { const client = await connectClient(transport, port); try { const started = performance.now(); const latencies = await runRequestLoad(client, total, concurrency, timeoutMs, stability); const elapsed = performance.now() - started; summarizeLatencies("roundtrip", `concurrency=${concurrency}`, latencies, elapsed, stability, "off", memMb(), { transport: transportName(transport), payloadBytes: testDataPayloadBytes("req-0"), clientCount: 1, }); } finally { await client.close(); if (!leakClear(client)) { stability.pendingLeak += 1; } } }); log(`[roundtrip] transport=${transport} concurrency=${concurrency} done violations=${stabilityViolations(stability)}`); } } /** close 후 communicator가 더 이상 살아있지 않은지로 pending/queue leak을 근사한다. */ function leakClear(client: BenchClient): boolean { return !client.communicator.isAlive(); } async function profileBurst(transport: WireTransport, mode: Mode): Promise { const counts = mode === "quick" ? [200] : [1_000, 10_000, 100_000]; log(`[burst] transport=${transport} mode=${mode} counts=${counts.join(",")}`); for (const count of counts) { const stability = newStability(); let received = 0; let lastIndex = -1; let resolveDone: (() => void) | null = null; const done = new Promise((resolve) => { resolveDone = resolve; }); const server = makeServer(transport); server.onClientConnected = (client: { communicator: Communicator }) => { addListenerTyped(client.communicator, TestDataSchema, (msg) => { if (msg.index <= lastIndex) { stability.fifoViolations += 1; } lastIndex = msg.index; received += 1; if (received >= count && resolveDone !== null) { resolveDone(); resolveDone = null; } }); }; await server.start(); const client = await connectClient(transport, server.port); const started = performance.now(); try { // fire-and-forget를 순차 await하여 등록 순서를 보존한다(writeQueue backpressure도 함께 검증). for (let i = 0; i < count; i += 1) { await client.communicator.send(create(TestDataSchema, { index: i, message: `b-${i}` })); } await withTimeout(done, mode === "quick" ? 10_000 : 60_000, `burst ${count} dispatch timeout`); const elapsed = performance.now() - started; if (received !== count) { stability.pendingLeak += 1; } emitThroughputRow("burst", `count=${count}`, count, elapsed, stability, "off", memMb(), { transport: transportName(transport), payloadBytes: testDataPayloadBytes("b-0"), clientCount: 1, }); } catch (err) { stability.timeouts += 1; log(`[burst] count=${count} error=${errorMessage(err)}`); emitThroughputRow("burst", `count=${count}`, received, performance.now() - started, stability, "off", memMb(), { transport: transportName(transport), payloadBytes: testDataPayloadBytes("b-0"), clientCount: 1, }); } finally { await client.close(); await server.stop(); if (!leakClear(client)) { stability.pendingLeak += 1; } } log(`[burst] transport=${transport} count=${count} received=${received} violations=${stabilityViolations(stability)}`); } } async function profileSustained(transport: WireTransport, mode: Mode): Promise { const durationsMs = mode === "quick" ? [2_000] : [30_000, 300_000, 1_800_000]; const concurrency = 16; const timeoutMs = 15_000; log(`[sustained] transport=${transport} mode=${mode} durations=${durationsMs.join(",")}ms concurrency=${concurrency}`); for (const durationMs of durationsMs) { const stability = newStability(); let peakMem = memMb(); await withRequestServer(transport, stability, async (port) => { const client = await connectClient(transport, port); const latencies: number[] = []; let nextIndex = 0; const deadline = performance.now() + durationMs; try { while (performance.now() < deadline) { const batch: Array> = []; for (let i = 0; i < concurrency; i += 1) { const index = nextIndex; nextIndex += 1; const started = performance.now(); batch.push( client.communicator .sendRequest( create(TestDataSchema, { index, message: `s-${index}` }), TestDataSchema, timeoutMs, ) .then((res) => { latencies.push(performance.now() - started); if (res.index !== index * 2 || res.message !== `echo:s-${index}`) { stability.nonceMismatch += 1; } }) .catch((err: unknown) => { classifyRequestError(errorMessage(err), stability); }), ); } await Promise.all(batch); const current = memMb(); if (current > peakMem) { peakMem = current; } } const elapsed = durationMs; summarizeLatencies("sustained", `duration=${durationMs}ms`, latencies, elapsed, stability, "off", peakMem, { transport: transportName(transport), payloadBytes: testDataPayloadBytes("s-0"), clientCount: 1, }); } finally { await client.close(); if (!leakClear(client)) { stability.pendingLeak += 1; } } }); log(`[sustained] transport=${transport} duration=${durationMs}ms done peakMemMb=${peakMem.toFixed(1)} violations=${stabilityViolations(stability)}`); } } async function profileParallel(transport: WireTransport, mode: Mode): Promise { const clientCounts = mode === "quick" ? [4] : [16, 128, 512, 1_024]; const perClient = mode === "quick" ? 50 : 500; const timeoutMs = 15_000; log(`[parallel] transport=${transport} mode=${mode} clients=${clientCounts.join(",")} perClient=${perClient}`); for (const clientCount of clientCounts) { const concurrency = mode === "full" && clientCount >= 512 ? 2 : 8; const splitProcess = transport === "tcp" && mode === "full" && clientCount >= 1_024; const axis = splitProcess ? `clients=${clientCount},topology=split-process` : `clients=${clientCount}`; const stability = newStability(); log( `[parallel] transport=${transport} clients=${clientCount} ` + `perClientConcurrency=${concurrency} topology=${splitProcess ? "split-process" : "in-process"}`, ); const withServer = splitProcess ? withSplitProcessRequestServer : withRequestServer; await withServer(transport, stability, async (port) => { const clients = await Promise.all( Array.from({ length: clientCount }, () => connectClient(transport, port)), ); const allLatencies: number[] = []; const started = performance.now(); try { await Promise.all( clients.map(async (client) => { const latencies = await runRequestLoad(client, perClient, concurrency, timeoutMs, stability); allLatencies.push(...latencies); }), ); const elapsed = performance.now() - started; summarizeLatencies( "parallel", axis, allLatencies, elapsed, stability, "off", memMb(), { transport: transportName(transport), payloadBytes: testDataPayloadBytes("req-0"), clientCount, }, ); } finally { await Promise.all(clients.map(async (client) => client.close())); for (const client of clients) { if (!leakClear(client)) { stability.pendingLeak += 1; } } } }); log( `[parallel] transport=${transport} clients=${clientCount} ` + `topology=${splitProcess ? "split-process" : "in-process"} done violations=${stabilityViolations(stability)}`, ); } } /** 빠른/느린 request handler를 섞어 same-connection FIFO delay와 connection 간 nonce 독립성을 검증한다. */ async function profileSlowMix(transport: WireTransport, mode: Mode): Promise { const clientCount = mode === "quick" ? 4 : 16; const perClient = mode === "quick" ? 20 : 200; const slowEvery = mode === "quick" ? 5 : 10; const delayMs = mode === "quick" ? 50 : 100; const timeoutMs = mode === "quick" ? 10_000 : 30_000; const axis = `clients=${clientCount},slowEvery=${slowEvery},delayMs=${delayMs}`; const stability = newStability(); let peakMem = memMb(); const latencies: number[] = []; let elapsedMs = 0; log(`[slow-mix] transport=${transport} mode=${mode} ${axis} perClient=${perClient}`); await withSlowMixServer(transport, slowEvery, delayMs, async (port) => { const clients: Array<{ client: BenchClient; clientId: number }> = []; try { for (let clientId = 0; clientId < clientCount; clientId += 1) { try { clients.push({ client: await connectClient(transport, port), clientId }); } catch (err) { stability.timeouts += 1; log(`[slow-mix] client=${clientId} dial error=${errorMessage(err)}`); } } const completedByClient = Array(clientCount).fill(-1) as number[]; const started = performance.now(); await Promise.all( clients.map(({ client, clientId }) => Promise.all( Array.from({ length: perClient }, async (_, index) => { const message = `sm-${clientId}-${index}`; const reqStart = performance.now(); try { const res = await client.communicator.sendRequest( create(TestDataSchema, { index, message }), TestDataSchema, timeoutMs, ); latencies.push(performance.now() - reqStart); if (res.index !== index * 2 || res.message !== `echo:${message}`) { stability.nonceMismatch += 1; } if (index <= completedByClient[clientId]) { stability.fifoViolations += 1; } completedByClient[clientId] = index; } catch (err) { classifyRequestError(errorMessage(err), stability); } }), ), ), ); peakMem = Math.max(peakMem, memMb()); elapsedMs = performance.now() - started; } finally { await Promise.all(clients.map(async ({ client }) => client.close())); for (const { client } of clients) { if (!leakClear(client)) { stability.pendingLeak += 1; } } } }); summarizeLatencies("slow-mix", axis, latencies, elapsedMs, stability, "off", peakMem, { transport: transportName(transport), payloadBytes: testDataPayloadBytes("sm-0-0"), clientCount, }); log(`[slow-mix] transport=${transport} clients=${clientCount} done peakMemMb=${peakMem.toFixed(1)} violations=${stabilityViolations(stability)}`); } async function withSlowMixServer( transport: WireTransport, slowEvery: number, delayMs: number, body: (port: number) => Promise, ): Promise { const server = makeServer(transport); server.onClientConnected = (client: { communicator: Communicator }) => { client.communicator.addRequestListener(TestDataSchema.typeName, async (msg) => { const data = msg as TestData; if ((data.index + 1) % slowEvery === 0) { await sleep(delayMs); } return create(TestDataSchema, { index: data.index * 2, message: `echo:${data.message}` }); }); }; await server.start(); try { await body(server.port); } finally { await server.stop(); } } /** * payload size matrix. 동일 connection request-response를 1KB/64KB/1MB payload에서 측정한다. * message field를 deterministic filler로 채우고 실제 serialized payload bytes를 기록하며, payload별 * p50/p95/p99 latency, throughput, peak memory, stability counters를 결과 row에 남긴다. */ async function profilePayload(transport: WireTransport, mode: Mode): Promise { const sizes = mode === "quick" ? PAYLOAD_SIZES_QUICK : PAYLOAD_SIZES_FULL; const concurrency = mode === "quick" ? 4 : 8; const requests = mode === "quick" ? 20 : 100; const timeoutMs = mode === "quick" ? 10_000 : 30_000; log( `[payload] transport=${transport} mode=${mode} sizes=${sizes.map(payloadLabel).join(",")} ` + `concurrency=${concurrency} requests/size=${requests}`, ); for (const size of sizes) { const stability = newStability(); const filler = payloadFiller(size); const expectedEcho = `echo:${filler}`; const payloadBytes = testDataPayloadBytes(filler); let peakMem = memMb(); await withRequestServer(transport, stability, async (port) => { const client = await connectClient(transport, port); try { const latencies: number[] = []; let issued = 0; let nextIndex = 0; const started = performance.now(); while (issued < requests) { const batchSize = Math.min(concurrency, requests - issued); const batch: Array> = []; for (let i = 0; i < batchSize; i += 1) { const index = nextIndex; nextIndex += 1; const reqStart = performance.now(); batch.push( client.communicator .sendRequest(create(TestDataSchema, { index, message: filler }), TestDataSchema, timeoutMs) .then((res) => { latencies.push(performance.now() - reqStart); if (res.index !== index * 2 || res.message !== expectedEcho) { stability.nonceMismatch += 1; } }) .catch((err: unknown) => { classifyRequestError(errorMessage(err), stability); }), ); } issued += batchSize; await Promise.all(batch); const current = memMb(); if (current > peakMem) { peakMem = current; } } const elapsed = performance.now() - started; summarizeLatencies("payload", `payload=${payloadLabel(size)}`, latencies, elapsed, stability, "off", peakMem, { transport: transportName(transport), payloadBytes, clientCount: 1, }); } finally { await client.close(); if (!leakClear(client)) { stability.pendingLeak += 1; } } }); log( `[payload] transport=${transport} size=${payloadLabel(size)} bytes=${payloadBytes} ` + `peakMemMb=${peakMem.toFixed(1)} violations=${stabilityViolations(stability)}`, ); } } const NOOP_TRANSPORT: Transport = { writePacket: async () => {}, close: async () => {}, }; /** * gateway receive path 측정 mode. off/on/worker_threads를 row label로 그대로 쓴다. * - off : gateway 미설정 inline control(legacy 대조군). * - on : default coordinator gateway. `Communicator.enableInboundGateway`가 내부에서 만드는 것과 * 같은 inline-decode {@link WorkerGateway}이며, worker_threads hop이 없다. * - worker_threads : explicit Node `worker_threads` off-thread decode 실험 row. */ type GatewayMode = "off" | "on" | "worker_threads"; const GATEWAY_MODES = ["off", "on", "worker_threads"] as const; /** * gateway mode별 server inbound gateway를 만든다. off는 gateway 없이 inline decode를 쓰므로 null을 반환한다. * on/worker_threads는 동일 sink/onError/workers로 만들어 backlog 계측과 결과 해석을 같은 기준으로 둔다. * `enableInboundGateway` 대신 {@link WorkerGateway}를 직접 만들어 default coordinator gateway 참조를 잡고, * 부하 중 peak backlog와 close 후 residual backlog를 동일하게 계측한다. */ function createGatewayForMode( mode: GatewayMode, sink: (frame: DecodedEnvelope) => void | Promise, onError: (err: Error) => void, ): InboundGateway | null { const workers = Math.max(2, Math.min(4, availableParallelism())); if (mode === "on") { return new WorkerGateway({ sink, onError, workers }); } if (mode === "worker_threads") { return createNodeWorkerGateway({ sink, onError, workers }); } return null; } /** * 실제 TCP/WS 수신 경로의 gateway baseline 1행을 mode별로 측정한다. server inbound path를 inline control(off), * default coordinator gateway(on), explicit worker_threads gateway(worker_threads)로 두고, 동일 concurrency/ * payload의 request-response를 흘려 throughput/p95/p99/peak memory를 기록한다. frame-ingest in-process 경로와 * 달리 실제 transport read loop(`onReceivedFrame`)를 통과한다. gateway backlog는 부하 중 peak과 close 후 * residual을 계측하며, residual이 0이 아니면 cleanup leak으로 본다. */ async function runGatewayTransportBaseline( transport: WireTransport, mode: Mode, gatewayMode: GatewayMode, ): Promise { const concurrency = 16; const batches = mode === "quick" ? 4 : 40; const total = concurrency * batches; const timeoutMs = mode === "quick" ? 8_000 : 20_000; const stability = newStability(); const server = makeServer(transport); let lastNonce = 0; // server inbound path에 attach된 gateway를 잡아 부하 중 peak backlog와 close 후 잔여 backlog를 계측한다. let serverGateway: InboundGateway | null = null; server.onClientConnected = (client: { communicator: Communicator }) => { const comm = client.communicator; // 효과 해석을 흐리지 않도록 server inbound path 한쪽만 gateway를 켠다. off mode는 gateway 없이 inline decode. const gw = createGatewayForMode( gatewayMode, (frame) => comm.enqueueInbound(frame.typeName, frame.data, frame.incomingNonce, frame.responseNonce), () => { stability.typeMismatch += 1; }, ); if (gw !== null) { comm.attachInboundGateway(gw); serverGateway = gw; } // raw addRequestListener로 nonce를 받아 gateway reorder 이후에도 connection별 FIFO(단조 증가)를 감시한다. comm.addRequestListener(TestDataSchema.typeName, (msg, nonce) => { if (nonce < lastNonce) { stability.fifoViolations += 1; } lastNonce = nonce; const data = msg as TestData; return create(TestDataSchema, { index: data.index * 2, message: `echo:${data.message}` }); }); }; await server.start(); let latencies: number[] = []; let elapsed = 0; let peakBacklog = 0; // 부하 중 gateway backlog peak을 샘플링한다. gateway off면 sampler 없이 backlog는 0으로 남는다. const sampler = gatewayMode !== "off" ? setInterval(() => { if (serverGateway !== null) { peakBacklog = Math.max(peakBacklog, gatewayBacklogOf(serverGateway)); } }, 2) : undefined; try { const client = await connectClient(transport, server.port); try { const started = performance.now(); latencies = await runRequestLoad(client, total, concurrency, timeoutMs, stability); elapsed = performance.now() - started; } finally { await client.close(); if (!leakClear(client)) { stability.pendingLeak += 1; } } } finally { if (sampler !== undefined) { clearInterval(sampler); } await server.stop(); } // server.stop()이 communicator/gateway를 close한다. close 후 잔여 backlog가 0이 아니면 cleanup leak이다. let residualBacklog = 0; if (serverGateway !== null) { residualBacklog = gatewayBacklogOf(serverGateway); if (residualBacklog > 0) { stability.pendingLeak += 1; } } const gatewayBacklog = Math.max(peakBacklog, residualBacklog); summarizeLatencies( "gateway", `transport=${transport} concurrency=${concurrency}`, latencies, elapsed, stability, gatewayMode, memMb(), { transport: transportName(transport), payloadBytes: testDataPayloadBytes("req-0"), clientCount: 1, gatewayBacklog, }, ); log( `[gateway] transport=${transport} mode=${gatewayMode} requests=${total} ` + `backlog=${gatewayBacklog} residual=${residualBacklog} violations=${stabilityViolations(stability)}`, ); } /** * gateway baseline. off/on/worker_threads 세 mode로 두 축을 측정한다: * 1) frame-ingest: communicator `onReceivedFrame` in-process 경로에 raw PacketBase frame을 흘려 * off=inline decode, on=default coordinator gateway, worker_threads=`worker_threads` decode pool의 * dispatch 순서/throughput을 비교한다. * 2) tcp/ws: 실제 transport read loop를 통과하는 request-response baseline으로, server inbound path의 * mode별 throughput/p95/p99/memory를 기록한다. * 작은 payload의 측정이므로 성능 이득이 불명확할 수 있고, 안정성(순서/leak) 0 위반만 합격선으로 둔다. */ async function profileGateway(mode: Mode, transports: readonly WireTransport[]): Promise { const count = mode === "quick" ? 500 : 5_000; log(`[gateway] mode=${mode} frames=${count}`); const frames: Uint8Array[] = []; for (let i = 0; i < count; i += 1) { const base = create(PacketBaseSchema, { typeName: TestDataSchema.typeName, nonce: i + 1, data: toBinary(TestDataSchema, create(TestDataSchema, { index: i, message: `g-${i}` })), }); frames.push(toBinary(PacketBaseSchema, base)); } const framePayloadBytes = frames[0]?.byteLength ?? 0; for (const gatewayMode of GATEWAY_MODES) { const stability = newStability(); const comm = new Communicator(); comm.initialize(NOOP_TRANSPORT, parserMap()); comm.setFrameErrorHandler(() => { stability.typeMismatch += 1; }); let dispatched = 0; let lastIndex = -1; let resolveDone: (() => void) | null = null; const done = new Promise((resolve) => { resolveDone = resolve; }); addListenerTyped(comm, TestDataSchema, (msg) => { if (msg.index <= lastIndex) { stability.fifoViolations += 1; } lastIndex = msg.index; dispatched += 1; if (dispatched >= count && resolveDone !== null) { resolveDone(); resolveDone = null; } }); const gateway = createGatewayForMode( gatewayMode, (frame) => comm.enqueueInbound(frame.typeName, frame.data, frame.incomingNonce, frame.responseNonce), () => { stability.typeMismatch += 1; }, ); if (gateway !== null) { comm.attachInboundGateway(gateway); } const started = performance.now(); let elapsed = 0; let peakBacklog = 0; let dispatchError: unknown = null; // ingest 중 gateway backlog peak을 샘플링한다. gateway off면 backlog는 0으로 남는다. const sampler = gateway !== null ? setInterval(() => { if (gateway !== null) { peakBacklog = Math.max(peakBacklog, gatewayBacklogOf(gateway)); } }, 2) : undefined; try { for (const raw of frames) { comm.onReceivedFrame(raw); } await withTimeout(done, mode === "quick" ? 15_000 : 60_000, `gateway dispatch timeout`); elapsed = performance.now() - started; if (dispatched !== count) { stability.pendingLeak += 1; } } catch (err) { dispatchError = err; elapsed = performance.now() - started; stability.timeouts += 1; log(`[gateway] mode=${gatewayMode} error=${errorMessage(err)}`); } finally { if (sampler !== undefined) { clearInterval(sampler); } if (gateway !== null) { peakBacklog = Math.max(peakBacklog, gatewayBacklogOf(gateway)); } await comm.close(); } // comm.close()가 gateway를 close한다. close 후 잔여 backlog가 0이 아니면 cleanup leak이다. let residualBacklog = 0; if (gateway !== null) { residualBacklog = gatewayBacklogOf(gateway); if (residualBacklog > 0) { stability.pendingLeak += 1; } } const gatewayBacklog = Math.max(peakBacklog, residualBacklog); emitThroughputRow( "gateway", `frames=${count}`, dispatchError !== null ? dispatched : count, elapsed, stability, gatewayMode, memMb(), { transport: TRANSPORT_FRAME_INGEST, payloadBytes: framePayloadBytes, clientCount: 0, gatewayBacklog, }, ); log( `[gateway] mode=${gatewayMode} dispatched=${dispatched} backlog=${gatewayBacklog} ` + `residual=${residualBacklog} violations=${stabilityViolations(stability)}`, ); } // 실제 transport 수신 경로의 gateway rows. 선택된 transport마다 off, on, worker_threads를 동일 조건으로 측정한다. for (const transport of transports) { for (const gatewayMode of GATEWAY_MODES) { await runGatewayTransportBaseline(transport, mode, gatewayMode); } } } async function withTimeout(promise: Promise, timeoutMs: number, label: string): Promise { let timer: ReturnType | undefined; try { return await Promise.race([ promise, new Promise((_, reject) => { timer = setTimeout(() => reject(new Error(label)), timeoutMs); }), ]); } finally { if (timer !== undefined) { clearTimeout(timer); } } } function parseArgs(argv: string[]): { mode: Mode; profiles: Profile[]; transports: WireTransport[]; serverChildTransport: WireTransport | null; } { let mode: Mode = "quick"; let profiles: Profile[] = [...ALL_PROFILES]; let transports: WireTransport[] = [...ALL_TRANSPORTS]; let serverChildTransport: WireTransport | null = null; for (const arg of argv) { if (arg === "--quick") { mode = "quick"; } else if (arg === "--full") { mode = "full"; } else if (arg.startsWith("--mode=")) { const value = arg.slice("--mode=".length); mode = value === "full" ? "full" : "quick"; } else if (arg.startsWith("--profiles=") || arg.startsWith("--profile=")) { const value = arg.slice(arg.indexOf("=") + 1); const requested = value .split(",") .map((p) => p.trim()) .filter((p) => p.length > 0); const selected = requested.filter((p): p is Profile => (ALL_PROFILES as readonly string[]).includes(p)); if (selected.length > 0) { profiles = selected; } } else if (arg.startsWith("--transports=") || arg.startsWith("--transport=")) { const value = arg.slice(arg.indexOf("=") + 1); const requested = value .split(",") .map((t) => t.trim()) .filter((t) => t.length > 0); const selected = requested.filter((t): t is WireTransport => (ALL_TRANSPORTS as readonly string[]).includes(t)); if (selected.length > 0) { transports = selected; } } else if (arg.startsWith("--server-child=")) { const value = arg.slice("--server-child=".length); if ((ALL_TRANSPORTS as readonly string[]).includes(value)) { serverChildTransport = value as WireTransport; } } } return { mode, profiles, transports, serverChildTransport }; } async function main(): Promise { const { mode, profiles, transports, serverChildTransport } = parseArgs(process.argv.slice(2)); if (serverChildTransport !== null) { await runServerChild(serverChildTransport); return; } log( `INFO stress harness language=${LANGUAGE} mode=${mode} transports=${transports.join(",")} ` + `profiles=${profiles.join(",")} typeName=${TestDataSchema.typeName}`, ); const transportRunners: Record Promise> = { roundtrip: profileRoundtrip, burst: profileBurst, sustained: profileSustained, parallel: profileParallel, "slow-mix": profileSlowMix, payload: profilePayload, }; // transport별 same-language 축은 transport를 바깥 루프로 돈다. gateway는 transport-agnostic in-process // frame-ingest 경로라 transport와 무관하게 한 번만 실행한다. for (const transport of transports) { for (const profile of profiles) { if (!TRANSPORT_PROFILES.has(profile)) { continue; } await transportRunners[profile](transport, mode); } } if (profiles.includes("gateway")) { await profileGateway(mode, transports); } const totalViolations = rows.reduce((acc, row) => acc + stabilityViolations(row.stability), 0); const status = totalViolations === 0 ? "PASS" : "FAIL"; process.stdout.write( `SUMMARY|status=${status}|language=${LANGUAGE}|mode=${mode}|transports=${transports.join(",")}|profiles=${profiles.join(",")}|rows=${rows.length}|stability_violations=${totalViolations}\n`, ); if (status !== "PASS") { process.exitCode = 1; } } void main().catch((err: unknown) => { process.stdout.write(`SUMMARY|status=FAIL|error=${errorMessage(err)}\n`); log(`FATAL ${errorMessage(err)}`); process.exitCode = 1; });