Skip to content

Commit cd32e1b

Browse files
committed
fix(dev): proxy websocket upgrades for devProxy routes with ws
1 parent 77b77ff commit cd32e1b

3 files changed

Lines changed: 154 additions & 0 deletions

File tree

src/dev/app.ts

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ import type { Nitro } from "nitro/types";
22
import type { H3Event, HTTPHandler } from "h3";
33
import { createProxyServer, type ProxyServerOptions } from "httpxy";
44
import type { IncomingMessage, ServerResponse } from "node:http";
5+
import type { Socket } from "node:net";
56
import { H3, toEventHandler, serveStatic, fromNodeHandler, HTTPError } from "h3";
67
import { joinURL } from "ufo";
78
import mime from "mime";
@@ -21,6 +22,8 @@ export class NitroDevApp {
2122
nitro: Nitro;
2223
fetch: (req: Request) => Response | Promise<Response>;
2324

25+
#wsProxies: { route: string; proxy: ReturnType<typeof createHTTPProxy> }[] = [];
26+
2427
constructor(nitro: Nitro, catchAllHandler?: HTTPHandler) {
2528
this.nitro = nitro;
2629
const app = this.#createApp(catchAllHandler);
@@ -97,6 +100,9 @@ export class NitroDevApp {
97100
}
98101
const proxy = createHTTPProxy(opts);
99102
app.all(route, proxy.handleEvent);
103+
if (opts.ws) {
104+
this.#wsProxies.push({ route, proxy });
105+
}
100106
}
101107

102108
// Main handler
@@ -106,6 +112,29 @@ export class NitroDevApp {
106112

107113
return app;
108114
}
115+
116+
/**
117+
* Proxy a WebSocket upgrade request if it matches a `devProxy` rule with `ws` enabled.
118+
*
119+
* @returns `true` if the socket was handed to a proxy, `false` if the caller should handle it.
120+
*/
121+
proxyUpgrade(req: IncomingMessage, socket: Socket, head: any): boolean {
122+
if (this.#wsProxies.length === 0) {
123+
return false;
124+
}
125+
const path = (req.url || "/").split("?")[0]!;
126+
const match = this.#wsProxies.find(
127+
({ route }) => path === route || path.startsWith(route.endsWith("/") ? route : `${route}/`)
128+
);
129+
if (!match) {
130+
return false;
131+
}
132+
match.proxy.proxy.ws(req, socket, {}, head).catch((error) => {
133+
this.nitro.logger.error(`Failed to proxy WebSocket upgrade for \`${path}\`:`, error);
134+
socket.destroy();
135+
});
136+
return true;
137+
}
109138
}
110139

111140
// TODO: upstream to h3/node

src/dev/server.ts

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -123,6 +123,9 @@ export class NitroDevServer extends NitroDevApp implements RunnerRPCHooks {
123123
// #region Public Methods
124124

125125
async upgrade(req: IncomingMessage, socket: Socket, head: any) {
126+
if (this.proxyUpgrade(req, socket, head)) {
127+
return;
128+
}
126129
if (!this.#manager.upgrade) {
127130
throw new HTTPError({
128131
status: 501,

test/unit/dev-proxy-ws.test.ts

Lines changed: 122 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,122 @@
1+
import { createServer, request, type Server } from "node:http";
2+
import type { AddressInfo, Socket } from "node:net";
3+
import type { Nitro } from "nitro/types";
4+
import { afterAll, beforeAll, describe, expect, it } from "vitest";
5+
import { NitroDevApp } from "../../src/dev/app.ts";
6+
7+
const sockets = new Set<Socket>();
8+
9+
function listen(server: Server): Promise<number> {
10+
server.on("connection", (socket) => {
11+
sockets.add(socket);
12+
socket.on("close", () => sockets.delete(socket));
13+
});
14+
return new Promise((resolve, reject) => {
15+
server.on("error", reject);
16+
server.listen(0, "127.0.0.1", () => resolve((server.address() as AddressInfo).port));
17+
});
18+
}
19+
20+
/** Minimal upgrade-aware target: replies `101` then echoes raw frames. */
21+
function createUpgradeTarget(upgrades: string[]): Server {
22+
const server = createServer((_req, res) => res.end("http"));
23+
server.on("upgrade", (req, socket) => {
24+
upgrades.push(req.url!);
25+
socket.write(
26+
"HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n\r\n"
27+
);
28+
socket.on("data", (chunk) => socket.write(chunk));
29+
});
30+
return server;
31+
}
32+
33+
function upgradeRequest(port: number, path: string) {
34+
return new Promise<{ status: number; echo: string }>((resolve, reject) => {
35+
const req = request({
36+
port,
37+
host: "127.0.0.1",
38+
path,
39+
headers: { Connection: "Upgrade", Upgrade: "websocket" },
40+
});
41+
const timeout = setTimeout(() => reject(new Error("upgrade timed out")), 5000);
42+
req.on("error", reject);
43+
req.on("response", (res) => {
44+
clearTimeout(timeout);
45+
reject(new Error(`unexpected HTTP response: ${res.statusCode}`));
46+
});
47+
req.on("upgrade", (res, socket) => {
48+
socket.once("data", (chunk) => {
49+
clearTimeout(timeout);
50+
socket.destroy();
51+
resolve({ status: res.statusCode!, echo: chunk.toString() });
52+
});
53+
socket.write("ping");
54+
});
55+
req.end();
56+
});
57+
}
58+
59+
describe("dev server devProxy websocket upgrades", () => {
60+
const upgrades: string[] = [];
61+
const target = createUpgradeTarget(upgrades);
62+
let devServer: Server;
63+
let devPort: number;
64+
let workerUpgrades = 0;
65+
66+
beforeAll(async () => {
67+
const targetPort = await listen(target);
68+
const app = new NitroDevApp({
69+
logger: console,
70+
options: {
71+
baseURL: "/",
72+
devHandlers: [],
73+
publicAssets: [],
74+
devProxy: {
75+
"/proxy/ws": { target: `http://127.0.0.1:${targetPort}`, ws: true },
76+
"/proxy/http": { target: `http://127.0.0.1:${targetPort}` },
77+
},
78+
},
79+
} as unknown as Nitro);
80+
devServer = createServer(async (req, res) => {
81+
const response = await app.fetch(new Request(`http://127.0.0.1${req.url}`));
82+
res.end(await response.text());
83+
});
84+
devServer.on("upgrade", (req, socket, head) => {
85+
if (!app.proxyUpgrade(req, socket as Socket, head)) {
86+
workerUpgrades++;
87+
socket.destroy();
88+
}
89+
});
90+
devPort = await listen(devServer);
91+
});
92+
93+
afterAll(async () => {
94+
for (const socket of sockets) {
95+
socket.destroy();
96+
}
97+
await Promise.all([
98+
new Promise((resolve) => devServer.close(resolve)),
99+
new Promise((resolve) => target.close(resolve)),
100+
]);
101+
});
102+
103+
it("proxies upgrades for a matching rule with `ws` enabled", async () => {
104+
const res = await upgradeRequest(devPort, "/proxy/ws");
105+
expect(res.status).toBe(101);
106+
expect(res.echo).toBe("ping");
107+
expect(upgrades).toContain("/proxy/ws");
108+
});
109+
110+
it("proxies upgrades for subpaths of a matching rule", async () => {
111+
const res = await upgradeRequest(devPort, "/proxy/ws/nested?foo=1");
112+
expect(res.status).toBe(101);
113+
expect(upgrades).toContain("/proxy/ws/nested?foo=1");
114+
});
115+
116+
it("forwards upgrades to the worker when no rule enables `ws`", async () => {
117+
await expect(upgradeRequest(devPort, "/proxy/http")).rejects.toThrow();
118+
await expect(upgradeRequest(devPort, "/other")).rejects.toThrow();
119+
expect(workerUpgrades).toBe(2);
120+
expect(upgrades).not.toContain("/proxy/http");
121+
});
122+
});

0 commit comments

Comments
 (0)