Skip to content
Merged
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 src/build/vite/_dev-worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,13 +21,14 @@ export async function writeDevWorkerEntry(nitro: Nitro): Promise<string> {
// 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();

Expand Down
12 changes: 2 additions & 10 deletions src/presets/_nitro/runtime/nitro-dev.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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) => {
Expand Down
5 changes: 1 addition & 4 deletions src/runtime/internal/vite/dev-entry.mjs
Original file line number Diff line number Diff line change
@@ -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";
Expand All @@ -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");
35 changes: 22 additions & 13 deletions src/runtime/internal/vite/dev-worker.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand All @@ -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;
}
}

Expand Down Expand Up @@ -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) {
Expand Down
10 changes: 10 additions & 0 deletions test/vite/websocket-fixture/routes/ws.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
import { defineWebSocketHandler } from "nitro";

export default defineWebSocketHandler({
open(peer) {
peer.send("ready");
},
message(peer, message) {
peer.send(`echo:${message.text()}`);
},
});
6 changes: 6 additions & 0 deletions test/vite/websocket-fixture/vite.config.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
import { nitro } from "nitro/vite";
import { defineConfig } from "vite";

export default defineConfig({
plugins: [nitro({ serverDir: "./", features: { websocket: true } })],
});
111 changes: 111 additions & 0 deletions test/vite/websocket.test.ts
Original file line number Diff line number Diff line change
@@ -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<string[]> {
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();
}
});
});
}
Loading