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,