From 6319a57a7656231bdee8b9ba4ed1ff65ad5d9d98 Mon Sep 17 00:00:00 2001 From: Pooya Parsa Date: Sat, 3 Oct 2026 14:32:40 +0000 Subject: [PATCH] fix(dev): use the websocket adapter of the runtime the dev worker runs in Supersedes #4376 and #4597. Co-authored-by: productdevbook --- src/build/vite/_dev-worker.ts | 3 +- src/presets/_nitro/runtime/nitro-dev.ts | 12 +-- src/runtime/internal/vite/dev-entry.mjs | 5 +- src/runtime/internal/vite/dev-worker.mjs | 35 ++++--- test/vite/websocket-fixture/routes/ws.ts | 10 ++ test/vite/websocket-fixture/vite.config.ts | 6 ++ test/vite/websocket.test.ts | 111 +++++++++++++++++++++ 7 files changed, 154 insertions(+), 28 deletions(-) create mode 100644 test/vite/websocket-fixture/routes/ws.ts create mode 100644 test/vite/websocket-fixture/vite.config.ts create mode 100644 test/vite/websocket.test.ts diff --git a/src/build/vite/_dev-worker.ts b/src/build/vite/_dev-worker.ts index 61098285d8..b74211473a 100644 --- a/src/build/vite/_dev-worker.ts +++ b/src/build/vite/_dev-worker.ts @@ -21,13 +21,14 @@ export async function writeDevWorkerEntry(nitro: Nitro): Promise { // specifiers cannot escape the entry directory in workerd (fails on windows). const devWorker = pathToFileURL(resolve(runtimeDir, "internal/vite/dev-worker.mjs")).href; const moduleRunner = pathToFileURL(await resolveViteModuleRunner(viteImportOptions(nitro))).href; + const websocket = nitro.options.features.websocket ?? nitro.options.experimental.websocket; const contents = /* js */ ` // Generated by Nitro import * as moduleRunner from ${JSON.stringify(moduleRunner)}; import { setModuleRunner } from ${JSON.stringify(devWorker)}; export * from ${JSON.stringify(devWorker)}; - +${websocket ? `export { websocketOptions as websocket } from ${JSON.stringify(devWorker)};\n` : ""} setModuleRunner(moduleRunner); `.trimStart(); diff --git a/src/presets/_nitro/runtime/nitro-dev.ts b/src/presets/_nitro/runtime/nitro-dev.ts index 498bc3717d..951c30941e 100644 --- a/src/presets/_nitro/runtime/nitro-dev.ts +++ b/src/presets/_nitro/runtime/nitro-dev.ts @@ -23,20 +23,12 @@ if (import.meta._tasks) { startScheduleRunner({}); } -const ws = import.meta._websocket - ? await import("crossws/adapters/node").then((m) => - (m.default || m)({ resolve: resolveWebsocketHooks }) - ) - : undefined; - export default { ...serverEntryOptions, fetch: nitroApp.fetch, plugins: [...tracingSrvxPlugins], - upgrade: ws - ? (context: { node: { req: any; socket: any; head: any } }) => { - ws.handleUpgrade(context.node.req, context.node.socket, context.node.head); - } + websocket: import.meta._websocket + ? ({ resolve: resolveWebsocketHooks } as AppEntry["websocket"]) : undefined, ipc: { onOpen: (ctx) => { diff --git a/src/runtime/internal/vite/dev-entry.mjs b/src/runtime/internal/vite/dev-entry.mjs index a7d6f75842..5b2ec762b3 100644 --- a/src/runtime/internal/vite/dev-entry.mjs +++ b/src/runtime/internal/vite/dev-entry.mjs @@ -1,5 +1,4 @@ import "#nitro/virtual/polyfills"; -import wsAdapter from "crossws/adapters/node"; import { useNitroApp, useNitroHooks } from "nitro/app"; import { resolveWebsocketHooks } from "#nitro/runtime/app"; @@ -9,13 +8,11 @@ const nitroApp = useNitroApp(); export const fetch = nitroApp.fetch; -const ws = import.meta._websocket ? wsAdapter({ resolve: resolveWebsocketHooks }) : undefined; +export const websocket = import.meta._websocket ? { resolve: resolveWebsocketHooks } : undefined; if (import.meta._tasks) { startScheduleRunner({}); } -export const handleUpgrade = ws?.handleUpgrade; - // Called by the dev worker when the runner shuts down (see `ipc.onClose` in `dev-worker.mjs`). export const close = () => useNitroHooks().callHook("close"); diff --git a/src/runtime/internal/vite/dev-worker.mjs b/src/runtime/internal/vite/dev-worker.mjs index a1028070fa..c730a30be2 100644 --- a/src/runtime/internal/vite/dev-worker.mjs +++ b/src/runtime/internal/vite/dev-worker.mjs @@ -172,8 +172,17 @@ class ViteEnvRunner { // they propagate to the caller (the nitro app's error handler or the // env-runner fetch boundary below). async fetch(req, init) { - // Wait until nothing is queued or in flight so requests never hit an entry - // that is about to be replaced. + const entry = await this.waitForEntry(); + const entryFetch = entry.default?.fetch || entry.fetch; + if (!entryFetch) { + throw httpError(500, `No fetch handler exported from ${this.entryPath}`); + } + return entryFetch(req, init); + } + + // Waits until nothing is queued or in flight so callers never reach an entry + // that is about to be replaced. + async waitForEntry() { const deadline = Date.now() + RELOAD_WAIT_TIMEOUT; let reloadPromise; while (reloadPromise !== this.reloadPromise) { @@ -191,11 +200,7 @@ class ViteEnvRunner { if (!this.entry) { throw httpError(503, `Vite environment "${this.name}" is unavailable`); } - const entryFetch = this.entry.default?.fetch || this.entry.fetch; - if (!entryFetch) { - throw httpError(500, `No fetch handler exported from ${this.entryPath}`); - } - return entryFetch(req, init); + return this.entry; } } @@ -273,12 +278,16 @@ export async function fetch(req) { } } -export function upgrade(context) { - const handleUpgrade = envs.nitro?.entry?.handleUpgrade; - if (handleUpgrade) { - handleUpgrade(context.node.req, context.node.socket, context.node.head); - } -} +// Re-exported as `websocket` by the generated entry only when the feature is enabled (see +// `build/vite/_dev-worker.ts`), so the runner installs its runtime's crossws adapter on demand. +export const websocketOptions = { + async resolve(request) { + const env = envs.nitro || (await waitForEnv("nitro")); + const entry = await env?.waitForEntry(); + const websocket = entry?.default?.websocket || entry?.websocket; + return (await websocket?.resolve(request)) || {}; + }, +}; export const ipc = { onOpen(ctx) { diff --git a/test/vite/websocket-fixture/routes/ws.ts b/test/vite/websocket-fixture/routes/ws.ts new file mode 100644 index 0000000000..0b2de3d6f2 --- /dev/null +++ b/test/vite/websocket-fixture/routes/ws.ts @@ -0,0 +1,10 @@ +import { defineWebSocketHandler } from "nitro"; + +export default defineWebSocketHandler({ + open(peer) { + peer.send("ready"); + }, + message(peer, message) { + peer.send(`echo:${message.text()}`); + }, +}); diff --git a/test/vite/websocket-fixture/vite.config.ts b/test/vite/websocket-fixture/vite.config.ts new file mode 100644 index 0000000000..b427835b63 --- /dev/null +++ b/test/vite/websocket-fixture/vite.config.ts @@ -0,0 +1,6 @@ +import { nitro } from "nitro/vite"; +import { defineConfig } from "vite"; + +export default defineConfig({ + plugins: [nitro({ serverDir: "./", features: { websocket: true } })], +}); diff --git a/test/vite/websocket.test.ts b/test/vite/websocket.test.ts new file mode 100644 index 0000000000..dbe2bf9acf --- /dev/null +++ b/test/vite/websocket.test.ts @@ -0,0 +1,111 @@ +import { mkdtempSync, rmSync, writeFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { fileURLToPath } from "node:url"; +import { execa, execaSync } from "execa"; +import { join } from "pathe"; +import { afterAll, beforeAll, describe, expect, test } from "vitest"; + +const rootDir = fileURLToPath(new URL("./websocket-fixture", import.meta.url)); +const vitePkg = process.env.NITRO_VITE_PKG || "vite"; + +const hasBun = execaSync("bun", ["--version"], { stdio: "ignore", reject: false }).exitCode === 0; + +// #3939: the dev worker loaded the Node `crossws` adapter unconditionally, so WebSocket upgrades +// broke as soon as it ran in another runtime (`bun --bun`). Each case starts a real dev server in a +// child process, because the runtime under test is the one that has to run it. +describe("dev: websocket", { concurrent: false }, () => { + let tmpDir: string; + const scripts = {} as Record<"vite" | "nitro", string>; + + beforeAll(() => { + tmpDir = mkdtempSync(join(tmpdir(), "nitro-dev-ws-")); + scripts.vite = join(tmpDir, "vite.mjs"); + writeFileSync( + scripts.vite, + /* js */ ` +import { createServer } from ${JSON.stringify(import.meta.resolve(vitePkg))}; +const server = await createServer({ root: ${JSON.stringify(rootDir)}, logLevel: "warn" }); +await server.listen(0); +console.log("ready:" + server.resolvedUrls.local[0]); +` + ); + scripts.nitro = join(tmpDir, "nitro.mjs"); + writeFileSync( + scripts.nitro, + /* js */ ` +import { build, createDevServer, createNitro, prepare } from ${JSON.stringify(import.meta.resolve("nitro/builder"))}; +const nitro = await createNitro({ + rootDir: ${JSON.stringify(rootDir)}, + buildDir: ${JSON.stringify(join(tmpDir, ".nitro"))}, + dev: true, + builder: "rolldown", + serverDir: "./", + features: { websocket: true }, +}); +const server = await createDevServer(nitro).listen({ port: 0, hostname: "127.0.0.1" }); +const ready = new Promise((resolve) => nitro.hooks.hook("dev:reload", resolve)); +await prepare(nitro); +await build(nitro); +await ready; +console.log("ready:" + server.url); +` + ); + }); + + afterAll(() => { + rmSync(tmpDir, { recursive: true, force: true }); + }); + + for (const dev of ["vite", "nitro"] as const) { + test(`${dev} dev (node)`, () => echo(process.execPath, [scripts[dev]]), 60_000); + test.runIf(hasBun)(`${dev} dev (bun)`, () => echo("bun", ["--bun", scripts[dev]]), 60_000); + } +}); + +async function echo(command: string, args: string[]) { + const child = execa(command, args, { cwd: rootDir, reject: false }); + let output = ""; + child.stdout!.on("data", (data) => (output += data)); + child.stderr!.on("data", (data) => (output += data)); + try { + const deadline = Date.now() + 30_000; + while (!/ready:\S+/.test(output) && Date.now() < deadline) { + await new Promise((resolve) => setTimeout(resolve, 100)); + } + const url = output.match(/ready:(\S+)/)?.[1]; + expect(url, output).toBeTruthy(); + // First request on a fresh server: the upgrade must wait for the app entry. + expect(await collectMessages(new URL("/ws", url!.replace(/^http/, "ws")).href), output).toEqual( + ["ready", "echo:hi"] + ); + } finally { + child.kill("SIGKILL"); + } +} + +// Connects, sends one message and resolves with everything the server sent back. +function collectMessages(url: string): Promise { + return new Promise((resolve, reject) => { + const messages: string[] = []; + const ws = new WebSocket(url); + const done = (error?: string) => { + clearTimeout(timer); + ws.close(); + if (error) { + reject(new Error(`${error} (received: ${JSON.stringify(messages)})`)); + } else { + resolve(messages); + } + }; + const timer = setTimeout(() => done("Timed out waiting for a WebSocket echo"), 20_000); + ws.addEventListener("open", () => ws.send("hi")); + ws.addEventListener("error", () => done("WebSocket connection failed")); + ws.addEventListener("close", () => done()); + ws.addEventListener("message", (event) => { + messages.push(String(event.data)); + if (messages.length >= 2) { + done(); + } + }); + }); +}