From 5970c0ff61ce61a3843a41a5d37182bf846ca90f Mon Sep 17 00:00:00 2001 From: Aditya Garud Date: Fri, 24 Jul 2026 23:54:41 +0530 Subject: [PATCH] Keep the server alive when a response write hits a dead socket A client that disconnects while an error response is being written crashes the whole Node server with an unhandled EPIPE, disconnecting every other client and abandoning in-flight provider work. Upgrade sockets are the worst case: once a connection upgrades for websocket RPC, Node's http server detaches its own socket error handling, so an auth rejection written to a vanished client emits an error event with no listener at all. Arm every response and upgrade socket with an error listener when the Node server is created. The failed request is already interrupted through its close event, so the write failure only needs to be observed instead of taking down the process. Fixes #4410 --- .../server/src/httpResponseErrorGuard.test.ts | 104 ++++++++++++++++++ apps/server/src/httpResponseErrorGuard.ts | 39 +++++++ apps/server/src/server.ts | 3 +- 3 files changed, 145 insertions(+), 1 deletion(-) create mode 100644 apps/server/src/httpResponseErrorGuard.test.ts create mode 100644 apps/server/src/httpResponseErrorGuard.ts diff --git a/apps/server/src/httpResponseErrorGuard.test.ts b/apps/server/src/httpResponseErrorGuard.test.ts new file mode 100644 index 00000000000..983547b5bcc --- /dev/null +++ b/apps/server/src/httpResponseErrorGuard.test.ts @@ -0,0 +1,104 @@ +// @effect-diagnostics nodeBuiltinImport:off +import * as NodeHttp from "node:http"; +import * as NodeNet from "node:net"; +import { afterEach, describe, expect, it } from "vite-plus/test"; + +import { guardHttpResponseWriteErrors } from "./httpResponseErrorGuard.ts"; + +const servers: NodeHttp.Server[] = []; + +function listen(server: NodeHttp.Server): Promise { + servers.push(server); + return new Promise((resolve) => { + server.listen(0, "127.0.0.1", () => { + resolve((server.address() as NodeNet.AddressInfo).port); + }); + }); +} + +function fetchStatus(port: number, path: string): Promise { + return new Promise((resolve, reject) => { + const request = NodeHttp.get({ host: "127.0.0.1", port, path }, (response) => { + response.resume(); + resolve(response.statusCode ?? 0); + }); + request.on("error", reject); + request.setTimeout(5_000, () => reject(new Error("request timed out"))); + }); +} + +afterEach(() => { + for (const server of servers.splice(0)) { + server.close(); + } +}); + +describe("guardHttpResponseWriteErrors", () => { + it("contains an upgrade socket write failure instead of crashing the process", async () => { + const writeErrors: unknown[] = []; + const failureObserved = Promise.withResolvers(); + const server = guardHttpResponseWriteErrors(NodeHttp.createServer(), (error) => { + writeErrors.push(error); + failureObserved.resolve(); + }); + + server.on("upgrade", (_request, socket) => { + // Simulate the client vanishing while the auth rejection response is + // written to the upgrade socket: the write failure surfaces as an + // "error" event on a socket Node's http server no longer listens to. + socket.destroy(Object.assign(new Error("write EPIPE"), { code: "EPIPE" })); + }); + + const port = await listen(server); + + const client = NodeNet.connect(port, "127.0.0.1", () => { + client.write( + [ + "GET /rpc HTTP/1.1", + "Host: 127.0.0.1", + "Connection: Upgrade", + "Upgrade: websocket", + "Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==", + "Sec-WebSocket-Version: 13", + "", + "", + ].join("\r\n"), + ); + }); + client.on("error", () => {}); + + await failureObserved.promise; + client.destroy(); + + expect(writeErrors).toHaveLength(1); + expect(writeErrors[0]).toBeInstanceOf(Error); + expect((writeErrors[0] as NodeJS.ErrnoException).code).toBe("EPIPE"); + + // The process survived the failed write and the server keeps serving. + server.on("request", (_request, response) => { + response.writeHead(200, { "content-type": "text/plain" }); + response.end("ok"); + }); + await expect(fetchStatus(port, "/")).resolves.toBe(200); + }); + + it("arms every response with an error listener without disturbing normal traffic", async () => { + const writeErrors: unknown[] = []; + let responseErrorListeners = -1; + const server = guardHttpResponseWriteErrors(NodeHttp.createServer(), (error) => { + writeErrors.push(error); + }); + + server.on("request", (_request, response) => { + responseErrorListeners = response.listenerCount("error"); + response.writeHead(200, { "content-type": "text/plain" }); + response.end("ok"); + }); + + const port = await listen(server); + + await expect(fetchStatus(port, "/")).resolves.toBe(200); + expect(responseErrorListeners).toBeGreaterThan(0); + expect(writeErrors).toEqual([]); + }); +}); diff --git a/apps/server/src/httpResponseErrorGuard.ts b/apps/server/src/httpResponseErrorGuard.ts new file mode 100644 index 00000000000..d673eaa42be --- /dev/null +++ b/apps/server/src/httpResponseErrorGuard.ts @@ -0,0 +1,39 @@ +// @effect-diagnostics nodeBuiltinImport:off +import type * as NodeHttp from "node:http"; + +/** + * Node surfaces late socket write failures (EPIPE, ECONNRESET, + * ERR_STREAM_DESTROYED) as "error" events. An "error" event without a + * listener escalates into an uncaught exception and terminates the whole + * server process, taking every other client and all in-flight provider + * work with it. + * + * Two emitters need coverage: + * + * - Upgrade sockets. Once a connection upgrades (the websocket RPC path, + * including its auth rejection responses), Node's http server detaches + * its own socket error handling, so the raw socket has no listener at + * all until the websocket server adopts it. + * - Server responses. Response streams have no default error listener + * either. + * + * A disconnected client only affects its own request: the request fiber is + * already interrupted through the response "close" event, so the write + * failure needs no handling beyond being observed. + */ +export function guardHttpResponseWriteErrors( + server: T, + onError?: (error: unknown) => void, +): T { + server.on("request", (_request, response) => { + response.on("error", (error) => { + onError?.(error); + }); + }); + server.on("upgrade", (_request, socket) => { + socket.on("error", (error) => { + onError?.(error); + }); + }); + return server; +} diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index 66b9823afb3..21418c06024 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -12,6 +12,7 @@ import { staticAndDevRouteLayer, browserApiCorsLayer, } from "./http.ts"; +import { guardHttpResponseWriteErrors } from "./httpResponseErrorGuard.ts"; import { fixPath } from "./os-jank.ts"; import { websocketRpcRouteLayer } from "./ws.ts"; import * as ExternalLauncher from "./process/externalLauncher.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(() => guardHttpResponseWriteErrors(NodeHttp.createServer()), { host: config.host ?? "127.0.0.1", port: config.port, gracefulShutdownTimeout: HTTP_PREEMPTIVE_SHUTDOWN_GRACE_MS,