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
3 changes: 2 additions & 1 deletion apps/server/src/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
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(() => guardUpgradeSockets(NodeHttp.createServer()), {
host: config.host ?? "127.0.0.1",
port: config.port,
gracefulShutdownTimeout: HTTP_PREEMPTIVE_SHUTDOWN_GRACE_MS,
Expand Down
96 changes: 96 additions & 0 deletions apps/server/src/upgradeSocketGuard.test.ts
Original file line number Diff line number Diff line change
@@ -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<number> =>
new Promise((resolve) => {
server.listen(0, "127.0.0.1", () => resolve((server.address() as NodeNet.AddressInfo).port));
});

const close = (server: NodeHttp.Server): Promise<void> =>
new Promise((resolve) => server.close(() => resolve()));

const resetUpgradeConnection = (port: number): Promise<void> =>
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<number> =>
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<number>((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<unknown> = [];
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);
}
});
15 changes: 15 additions & 0 deletions apps/server/src/upgradeSocketGuard.ts
Original file line number Diff line number Diff line change
@@ -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;
};
Loading