Skip to content

Commit 1d7e5f4

Browse files
committed
fix(dev): close stalled gateway websocket handshakes
1 parent 1fd2259 commit 1d7e5f4

2 files changed

Lines changed: 72 additions & 10 deletions

File tree

scripts/dev/gateway-ws-client.ts

Lines changed: 17 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -108,18 +108,26 @@ export function createGatewayWsClient(params: {
108108

109109
const waitOpen = () =>
110110
new Promise<void>((resolve, reject) => {
111-
const t = setTimeout(
112-
() => reject(new Error("ws open timeout")),
113-
params.openTimeoutMs ?? 8000,
114-
);
115-
ws.once("open", () => {
111+
const cleanup = () => {
116112
clearTimeout(t);
113+
ws.off("open", onOpen);
114+
ws.off("error", onError);
115+
};
116+
const onOpen = () => {
117+
cleanup();
117118
resolve();
118-
});
119-
ws.once("error", (err) => {
120-
clearTimeout(t);
119+
};
120+
const onError = (err: Error) => {
121+
cleanup();
121122
reject(err instanceof Error ? err : new Error(String(err)));
122-
});
123+
};
124+
const t = setTimeout(() => {
125+
cleanup();
126+
ws.terminate();
127+
reject(new Error("ws open timeout"));
128+
}, params.openTimeoutMs ?? 8000);
129+
ws.once("open", onOpen);
130+
ws.once("error", onError);
123131
});
124132

125133
ws.on("message", (data) => {

test/scripts/gateway-ws-client.test.ts

Lines changed: 55 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
import { createServer, type Server } from "node:http";
22
import { afterEach, describe, expect, it } from "vitest";
3-
import { WebSocketServer, type WebSocket } from "ws";
3+
import { WebSocket, WebSocketServer } from "ws";
44
import { createGatewayWsClient } from "../../scripts/dev/gateway-ws-client.js";
55

66
let server: Server | undefined;
@@ -44,6 +44,36 @@ async function listen(handler: (ws: WebSocket) => void): Promise<string> {
4444
return `ws://127.0.0.1:${address.port}`;
4545
}
4646

47+
async function listenStalledUpgrade(): Promise<{ close: () => Promise<void>; url: string }> {
48+
const stalledServer = createServer();
49+
const sockets = new Set<import("node:net").Socket>();
50+
stalledServer.on("upgrade", (_req, socket) => {
51+
// Keep the socket open without completing the websocket handshake.
52+
sockets.add(socket);
53+
socket.once("close", () => {
54+
sockets.delete(socket);
55+
});
56+
});
57+
await new Promise<void>((resolve) => {
58+
stalledServer.listen(0, "127.0.0.1", resolve);
59+
});
60+
const address = stalledServer.address();
61+
if (!address || typeof address === "string") {
62+
throw new Error("test websocket server did not get a TCP address");
63+
}
64+
return {
65+
close: async () => {
66+
for (const socket of sockets) {
67+
socket.destroy();
68+
}
69+
await new Promise<void>((resolve, reject) => {
70+
stalledServer.close((error) => (error ? reject(error) : resolve()));
71+
});
72+
},
73+
url: `ws://127.0.0.1:${address.port}`,
74+
};
75+
}
76+
4777
describe("createGatewayWsClient", () => {
4878
it("rejects pending RPC requests when the client closes", async () => {
4979
const url = await listen(() => {});
@@ -70,4 +100,28 @@ describe("createGatewayWsClient", () => {
70100
);
71101
client.close();
72102
});
103+
104+
it("terminates stalled websocket handshakes after the open timeout", async () => {
105+
const stalled = await listenStalledUpgrade();
106+
const client = createGatewayWsClient({ openTimeoutMs: 5, url: stalled.url });
107+
try {
108+
await expect(client.waitOpen()).rejects.toThrow("ws open timeout");
109+
await waitFor(() => client.ws.readyState === WebSocket.CLOSED);
110+
} finally {
111+
client.close();
112+
await stalled.close();
113+
}
114+
});
73115
});
116+
117+
async function waitFor(condition: () => boolean, timeoutMs = 1_000) {
118+
const startedAt = Date.now();
119+
while (!condition()) {
120+
if (Date.now() - startedAt > timeoutMs) {
121+
throw new Error("timed out waiting for condition");
122+
}
123+
await new Promise((resolve) => {
124+
setTimeout(resolve, 10);
125+
});
126+
}
127+
}

0 commit comments

Comments
 (0)