Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
104 changes: 104 additions & 0 deletions apps/server/src/httpResponseErrorGuard.test.ts
Original file line number Diff line number Diff line change
@@ -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<number> {
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<number> {
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<void>();
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([]);
});
});
39 changes: 39 additions & 0 deletions apps/server/src/httpResponseErrorGuard.ts
Original file line number Diff line number Diff line change
@@ -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<T extends NodeHttp.Server>(
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;
}
3 changes: 2 additions & 1 deletion apps/server/src/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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,
Expand Down
Loading