diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index 66b9823afb3..7c59e7ac7fc 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -13,6 +13,7 @@ import { browserApiCorsLayer, } from "./http.ts"; import { fixPath } from "./os-jank.ts"; +import { guardUpgradeSockets } from "./upgradeSocketGuard.ts"; import { websocketRpcRouteLayer } from "./ws.ts"; import * as ExternalLauncher from "./process/externalLauncher.ts"; import { layerConfig as SqlitePersistenceLayerLive } from "./persistence/Layers/Sqlite.ts"; @@ -137,7 +138,7 @@ const HttpServerLive = Layer.unwrap( Effect.promise(() => import("@effect/platform-node/NodeHttpServer")), Effect.promise(() => import("node:http")), ]); - return NodeHttpServer.layer(NodeHttp.createServer, { + return NodeHttpServer.layer(() => guardUpgradeSockets(NodeHttp.createServer()), { host: config.host ?? "127.0.0.1", port: config.port, gracefulShutdownTimeout: HTTP_PREEMPTIVE_SHUTDOWN_GRACE_MS, diff --git a/apps/server/src/upgradeSocketGuard.test.ts b/apps/server/src/upgradeSocketGuard.test.ts new file mode 100644 index 00000000000..27812590d3e --- /dev/null +++ b/apps/server/src/upgradeSocketGuard.test.ts @@ -0,0 +1,96 @@ +import * as NodeHttp from "node:http"; +import * as NodeNet from "node:net"; + +import { expect, it } from "@effect/vitest"; + +import { guardUpgradeSockets } from "./upgradeSocketGuard.ts"; + +const UPGRADE_REQUEST = + "GET /ws HTTP/1.1\r\n" + + "Host: 127.0.0.1\r\n" + + "Upgrade: websocket\r\n" + + "Connection: Upgrade\r\n" + + "Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==\r\n" + + "Sec-WebSocket-Version: 13\r\n" + + "\r\n"; + +const listen = (server: NodeHttp.Server): Promise => + new Promise((resolve) => { + server.listen(0, "127.0.0.1", () => resolve((server.address() as NodeNet.AddressInfo).port)); + }); + +const close = (server: NodeHttp.Server): Promise => + new Promise((resolve) => server.close(() => resolve())); + +const resetUpgradeConnection = (port: number): Promise => + new Promise((resolve) => { + const socket = NodeNet.connect(port, "127.0.0.1"); + socket.on("error", () => resolve()); + socket.on("connect", () => { + socket.write(UPGRADE_REQUEST); + setTimeout(() => { + socket.resetAndDestroy(); + resolve(); + }, 50); + }); + }); + +const makeUpgradeServer = () => { + const server = NodeHttp.createServer((_request, response) => { + response.writeHead(200); + response.end("ok"); + }); + server.on("upgrade", (_request, socket) => { + socket.write( + "HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n\r\n", + ); + }); + return server; +}; + +const httpGet = (port: number): Promise => + new Promise((resolve, reject) => { + NodeHttp.get({ host: "127.0.0.1", port, path: "/" }, (response) => { + response.resume(); + resolve(response.statusCode ?? 0); + }).on("error", reject); + }); + +it("attaches an error listener to every upgrade socket", async () => { + const server = guardUpgradeSockets(makeUpgradeServer()); + // Registered after the guard, so it runs second and observes the guard's listener. + const observed = new Promise((resolve) => { + server.on("upgrade", (_request, socket) => resolve(socket.listenerCount("error"))); + }); + const port = await listen(server); + try { + const client = NodeNet.connect(port, "127.0.0.1"); + client.on("error", () => {}); + client.on("connect", () => client.write(UPGRADE_REQUEST)); + expect(await observed).toBe(1); + client.destroy(); + } finally { + await close(server); + } +}); + +it("keeps the server alive when a client RSTs a /ws upgrade connection", async () => { + const uncaught: Array = []; + const onUncaught = (error: unknown) => uncaught.push(error); + process.on("uncaughtException", onUncaught); + + const server = guardUpgradeSockets(makeUpgradeServer()); + const port = await listen(server); + try { + for (let i = 0; i < 10; i++) { + await resetUpgradeConnection(port); + } + await new Promise((resolve) => setTimeout(resolve, 100)); + + expect(uncaught).toEqual([]); + expect(await httpGet(port)).toBe(200); + } finally { + process.off("uncaughtException", onUncaught); + await close(server); + } +}); diff --git a/apps/server/src/upgradeSocketGuard.ts b/apps/server/src/upgradeSocketGuard.ts new file mode 100644 index 00000000000..1399e18b6ad --- /dev/null +++ b/apps/server/src/upgradeSocketGuard.ts @@ -0,0 +1,15 @@ +import type * as NodeHttp from "node:http"; +import type * as NodeNet from "node:net"; + +// Node throws an uncaught exception — taking the whole process down — when a +// socket emits 'error' with no listener attached. @effect/platform-node's +// upgrade handler wires a 'close' listener on the raw upgrade socket but not an +// 'error' one, so a client that RSTs a /ws connection during or right after the +// handshake kills the server. Attaching a no-op 'error' listener on every +// upgrade socket turns the reset into a normal disconnect. +export const guardUpgradeSockets = (server: NodeHttp.Server): NodeHttp.Server => { + server.on("upgrade", (_request, socket: NodeNet.Socket) => { + socket.on("error", () => {}); + }); + return server; +};