From f2341259fd8d82d2bea879be5edc8f254ceb868f Mon Sep 17 00:00:00 2001 From: Oskar Lebuda Date: Mon, 6 Jul 2026 23:33:45 +0200 Subject: [PATCH 1/6] feat: add cluster mode --- .github/workflows/ci.yml | 6 +- docs/1.guide/10.cluster.md | 70 ++++++ docs/1.guide/5.options.md | 10 + docs/1.guide/9.cli.md | 3 + src/_cluster.ts | 375 +++++++++++++++++++++++++++++++ src/_plugins.ts | 8 + src/adapters/bun.ts | 5 +- src/adapters/deno.ts | 5 +- src/adapters/node.ts | 3 +- src/cli/main.ts | 32 ++- src/cli/serve.ts | 5 +- src/cli/types.ts | 2 + src/cli/usage.ts | 3 + src/types.ts | 19 ++ test/cluster.test.ts | 183 +++++++++++++++ test/fixtures/cluster/server.mjs | 22 ++ 16 files changed, 740 insertions(+), 11 deletions(-) create mode 100644 docs/1.guide/10.cluster.md create mode 100644 src/_cluster.ts create mode 100644 test/cluster.test.ts create mode 100644 test/fixtures/cluster/server.mjs diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 85e04cca..8012827b 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -27,7 +27,7 @@ jobs: - uses: actions/setup-node@v5 with: { node-version: "${{ matrix.node-version }}", cache: pnpm } - run: pnpm install - - run: pnpm vitest --coverage test/node.test.ts test/node-adapters.test.ts test/url.test.ts + - run: pnpm vitest --coverage test/node.test.ts test/node-adapters.test.ts test/url.test.ts test/cluster.test.ts - uses: codecov/codecov-action@v3 if: matrix.node-version == 24 with: @@ -44,7 +44,7 @@ jobs: # - uses: denoland/setup-deno@v2 - run: mkdir -p "$HOME/.deno/bin" && curl -fsSL https://github.com/denoland/deno/releases/download/v2.7.12/deno-x86_64-unknown-linux-gnu.zip -o /tmp/deno.zip && unzip -q /tmp/deno.zip -d "$HOME/.deno/bin" && echo "$HOME/.deno/bin" >> "$GITHUB_PATH" - run: pnpm install - - run: pnpm vitest --coverage test/deno.test.ts test/url.test.ts + - run: pnpm vitest --coverage test/deno.test.ts test/url.test.ts test/cluster.test.ts - run: pnpm add -D undici@^7.25.0 - run: deno run test:node-compat:deno tests_bun: @@ -56,7 +56,7 @@ jobs: with: { node-version: lts/*, cache: pnpm } - uses: oven-sh/setup-bun@v2 - run: pnpm install - - run: pnpm vitest --coverage test/bun.test.ts test/url.test.ts + - run: pnpm vitest --coverage test/bun.test.ts test/url.test.ts test/cluster.test.ts - run: bun run test:node-compat:bun publish: runs-on: ubuntu-latest diff --git a/docs/1.guide/10.cluster.md b/docs/1.guide/10.cluster.md new file mode 100644 index 00000000..dc572f49 --- /dev/null +++ b/docs/1.guide/10.cluster.md @@ -0,0 +1,70 @@ +--- +icon: ri:stack-line +--- + +# Cluster Mode + +> Run multiple server processes sharing the same port + +A single JavaScript process uses one CPU core. In production, cluster mode lets srvx spawn multiple worker processes that share the same port, so incoming requests are load balanced across all CPU cores - no external process manager needed. + +The main process becomes a lightweight supervisor. It never handles requests itself; it spawns workers (re-executing the same server entry), restarts crashed workers (with exponential backoff) and forwards shutdown signals so each worker can close gracefully. + +## Usage + +**CLI:** + +```bash +# One worker per CPU core +srvx serve --prod --cluster + +# Exact number of workers +srvx serve --prod --cluster=4 +``` + +> [!NOTE] +> Cluster mode is production-only in the CLI (in dev mode a single process with watcher is used). + +**Programmatic:** + +```js +import { serve } from "srvx"; + +serve({ + cluster: true, // or an exact number of workers + fetch: () => new Response(`👋 Hello from worker ${process.env.SRVX_CLUSTER_WORKER}`), +}); +``` + +**Environment variable:** + +Setting `SRVX_WORKERS` enables cluster mode without touching code or CLI flags (handy in Docker / Kubernetes): + +```bash +SRVX_WORKERS=4 srvx serve --prod +``` + +An explicit `cluster: ` option takes precedence over `SRVX_WORKERS`. Use `cluster: false` (or `--cluster=false` in the CLI) to disable cluster mode entirely, including `SRVX_WORKERS`. + +## Worker processes + +Each worker re-executes the same server entry. Workers can be detected via the `SRVX_CLUSTER_WORKER` environment variable, which contains the worker index (starting at `"0"`): + +```js +if (process.env.SRVX_CLUSTER_WORKER) { + // Running as a cluster worker +} +``` + +If a worker crashes, the supervisor restarts it automatically (with exponential backoff). If a worker repeatedly fails during startup (e.g. port in use, invalid TLS config), the supervisor gives up after 3 attempts and exits with a non-zero code. + +On `SIGINT`/`SIGTERM`, the supervisor forwards the signal to all workers for graceful shutdown and exits once all of them have finished. + +## Runtime support + +- **Node.js**: workers share the listening socket via [`node:cluster`](https://nodejs.org/api/cluster.html) with round-robin load balancing (all platforms, including macOS and Windows). +- **Bun** and **Deno**: workers bind the port with `SO_REUSEPORT` and the kernel load balances between them. Kernel load balancing is **Linux only** - on other platforms a single supervised worker is started instead (crash restarts still work). +- **Serverless runtimes** (Cloudflare, AWS Lambda, ...): the platform scales processes itself, the `cluster` option is ignored. + +> [!IMPORTANT] +> Cluster mode requires a fixed port (`port: 0` is not supported). diff --git a/docs/1.guide/5.options.md b/docs/1.guide/5.options.md index 4d87eef6..838dc1ad 100644 --- a/docs/1.guide/5.options.md +++ b/docs/1.guide/5.options.md @@ -58,6 +58,16 @@ Enabling this option allows multiple processes to bind to the same port, which i > [!NOTE] > Despite Node.js built-in behavior that has `exclusive` flag enabled by default, srvx uses non-exclusive mode for consistency. +### `cluster` + +Run multiple server processes sharing the same port. + +- `true`: number of workers from the `SRVX_WORKERS` environment variable, or CPU cores. +- `number`: exact number of worker processes. +- `false`: explicitly disable cluster mode (also ignores `SRVX_WORKERS`). + +:read-more{to="/guide/cluster"} + ### `silent` If enabled, no server listening message will be printed (enabled by default when `TEST` environment variable is set). diff --git a/docs/1.guide/9.cli.md b/docs/1.guide/9.cli.md index 07004419..9ab61174 100644 --- a/docs/1.guide/9.cli.md +++ b/docs/1.guide/9.cli.md @@ -40,6 +40,7 @@ SERVE MODE # srvx serve [options] $ srvx serve --entry ./server.ts # Start development server $ srvx serve --prod # Start production server +$ srvx serve --prod --cluster # Production server with one worker per CPU core $ srvx serve --port=8080 # Listen on port 8080 $ srvx serve --host=localhost # Bind to localhost only $ srvx serve --static=./dist # Serve static files (no entry needed) @@ -71,6 +72,7 @@ SERVE OPTIONS --host Host to bind to (default: all interfaces) -s, --static Serve static files from the specified directory (default: public) --prod Run in production mode (no watch, no debug) + --cluster [N] Run N server processes (default: CPU cores, requires --prod) --import ES module to preload --tls Enable TLS (HTTPS/HTTP2) --cert TLS certificate file @@ -87,6 +89,7 @@ ENVIRONMENT PORT Override port HOST Override host + SRVX_WORKERS Number of cluster workers (enables cluster mode) NODE_ENV Set to production for production mode. ``` diff --git a/src/_cluster.ts b/src/_cluster.ts new file mode 100644 index 00000000..4e9aedd4 --- /dev/null +++ b/src/_cluster.ts @@ -0,0 +1,375 @@ +import nodeCluster from "node:cluster"; +import { fork, type ChildProcess } from "node:child_process"; +import { availableParallelism } from "node:os"; +import * as c from "./cli/_utils.ts"; +import { fmtURL, printListening, resolvePortAndHost } from "./_utils.ts"; +import type { Server, ServerHandler, ServerOptions } from "./types.ts"; + +/** + * Name of the environment variable set in cluster worker processes. + * The value is the worker index, starting from `"0"`. + */ +export const CLUSTER_WORKER_ENV = "SRVX_CLUSTER_WORKER"; + +const IS_BUN = !!globalThis.process?.versions?.bun; +const IS_DENO = !!globalThis.process?.versions?.deno; + +// Uptime after which a worker is considered stable and its crash counter resets. +const STABLE_UPTIME = 10_000; + +// Crash-loop protection: give up if a worker never becomes ready after this many attempts. +const MAX_START_ATTEMPTS = 3; + +type ReadyMessage = { srvx?: string; url?: string }; + +/** + * Adds cluster (multi-process) support to a runtime adapter's `serve()`. + * + * When cluster mode is enabled, the main process becomes a supervisor that only + * spawns and monitors workers, while worker processes (detected via + * `SRVX_CLUSTER_WORKER`) serve regularly on the shared port. Otherwise the + * call passes through to the adapter factory. + * + * @param options Server options as passed to `serve()`. + * @param factory Creates the runtime-specific server instance. + * @returns The cluster supervisor, or the server created by `factory`. + */ +export function withCluster( + options: ServerOptions, + factory: (options: ServerOptions) => Server, +): Server { + const env = globalThis.process?.env; + // Loader context (entry is being inspected, server won't listen) or non-process runtime + if (!env || (globalThis as any).__srvxLoader__) { + return factory(options); + } + // Worker process: start a regular server on the shared port + if (env[CLUSTER_WORKER_ENV]) { + // Deno supports SO_REUSEPORT on Linux only (non-Linux runs a single + // supervised worker that binds the port exclusively). + const reusePort = !(IS_DENO && process.platform !== "linux"); + const server = factory({ ...options, cluster: false, reusePort, silent: true }); + Promise.resolve(server.ready()).then( + () => process.send?.({ srvx: "cluster-worker-ready", url: server.url }), + (error) => { + console.error(error); + process.exit(1); + }, + ); + return server; + } + const workers = resolveClusterSize(options); + if (workers === undefined) { + return factory(options); + } + return new ClusterServer(options, workers); +} + +/** + * Resolves how many cluster workers should be spawned. + * + * An explicit numeric `cluster` option wins over the `SRVX_WORKERS` environment + * variable, which in turn wins over the CPU core count. + * + * @param options Server options as passed to `serve()`. + * @returns The worker count, or `undefined` when cluster mode is disabled. + */ +function resolveClusterSize(options: ServerOptions): number | undefined { + if (options.cluster === false || options.cluster === 0) { + return; + } + const envSize = Number.parseInt(globalThis.process?.env?.SRVX_WORKERS || "", 10) || undefined; + if (!options.cluster && !envSize) { + return; + } + const size = + typeof options.cluster === "number" ? options.cluster : (envSize ?? availableParallelism()); + return Math.max(1, Math.floor(size)); +} + +type WorkerState = { + child: ChildProcess; + ready: boolean; + startedAt: number; +}; + +/** + * Cluster supervisor implementing the `Server` interface. + * + * It never listens itself — it spawns worker processes that serve on the + * shared port, restarts them when they crash (with exponential backoff) and + * forwards `SIGINT`/`SIGTERM` for graceful shutdown. + */ +class ClusterServer implements Server { + readonly runtime = IS_BUN ? "bun" : IS_DENO ? "deno" : "node"; + readonly options: Server["options"]; + readonly fetch: ServerHandler; + + #size: number; + #workers = new Map(); + #respawnTimers = new Set>(); + #ready = Promise.withResolvers(); + #url?: string; + #fallbackURL?: string; + #started = false; + #announced = false; + #closing = false; + + constructor(options: ServerOptions, size: number) { + this.options = { ...options, middleware: [...(options.middleware || [])] }; + this.fetch = options.fetch; + this.#size = size; + this.#ready.promise.catch(() => {}); // avoid unhandled rejection when ready() is not awaited + + if (!process.argv[1]) { + throw new Error( + "Cluster mode requires a server entry file (cannot re-execute this process).", + ); + } + const { port, hostname } = resolvePortAndHost(options); + if (!port) { + throw new Error("Cluster mode requires a fixed port (port: 0 is not supported)."); + } + const secure = !!(options.tls?.cert || options.protocol === "https"); + this.#fallbackURL = fmtURL(hostname || "localhost", port, secure); + + if (!options.manual) { + this.serve().catch(() => {}); + } + } + + /** + * Spawns all workers (only once) and registers signal forwarding. + * + * @returns A promise that resolves once every worker is ready. + */ + serve(): Promise { + if (!this.#started) { + this.#started = true; + + // SO_REUSEPORT load balancing is Linux-only for Bun/Deno: fall back to a + // single supervised worker instead of spawning processes that would never + // receive connections (Node uses node:cluster round-robin on all platforms). + if (process.platform !== "linux" && (IS_BUN || IS_DENO) && this.#size > 1) { + this.#log( + c.yellow, + `Cluster load balancing requires Linux on ${IS_DENO ? "Deno" : "Bun"} (starting 1 worker)`, + ); + this.#size = 1; + } + + this.#log(c.gray, `Starting ${this.#size} cluster worker${this.#size > 1 ? "s" : ""}...`); + for (let slot = 0; slot < this.#size; slot++) { + this.#spawn(slot, 0); + } + for (const signal of ["SIGINT", "SIGTERM"] as const) { + process.on(signal, this.#onSignal); + } + } + return this.ready(); + } + + get url(): string | undefined { + return this.#url || this.#fallbackURL; + } + + ready(): Promise { + return this.#ready.promise.then(() => this); + } + + /** + * Stops all workers and waits for them to exit. + * + * @param closeAll Escalate to `SIGKILL` for workers still alive after 1s. + */ + async close(closeAll?: boolean): Promise { + for (const signal of ["SIGINT", "SIGTERM"] as const) { + process.off(signal, this.#onSignal); + } + + await this.#killAll(closeAll); + + // Unblock pending ready() awaiters: reject when closed before the cluster + // ever became ready, so callers can tell "listening" from "shut down". + if (this.#announced) { + this.#ready.resolve(); + } else { + this.#ready.reject(new Error("Cluster server closed before becoming ready.")); + } + } + + /** + * Sends `SIGTERM` to all live workers and resolves once they exited. + * + * @param forceClose `SIGKILL` workers still alive after 1s. + */ + #killAll(forceClose?: boolean): Promise { + this.#closing = true; + this.#clearRespawns(); + + const children = [...this.#workers.values()].map((w) => w.child); + const alive = children.filter((child) => child.exitCode === null && child.signalCode === null); + const exits = alive.map( + (child) => new Promise((resolve) => child.once("exit", () => resolve())), + ); + + for (const child of alive) { + child.kill("SIGTERM"); + } + + const killTimer = forceClose + ? setTimeout(() => { + for (const child of alive) { + if (child.exitCode === null && child.signalCode === null) { + child.kill("SIGKILL"); + } + } + }, 1000) + : undefined; + + killTimer?.unref?.(); + + return Promise.all(exits).then(() => { + if (killTimer) { + clearTimeout(killTimer); + } + }); + } + + /** + * Starts a worker for the given slot and wires up its readiness and restart handling. + * + * @param slot Worker slot index (stable across restarts). + * @param restarts Consecutive restarts of this slot so far (drives the backoff). + */ + #spawn(slot: number, restarts: number): void { + if (this.#closing) { + return; + } + + const child = this.#forkWorker(slot); + const state: WorkerState = { child, ready: false, startedAt: Date.now() }; + this.#workers.set(slot, state); + + child.on("message", (message: ReadyMessage) => { + if (message?.srvx === "cluster-worker-ready" && !state.ready) { + state.ready = true; + this.#url ||= message.url; + + if (!this.#announced && this.#allReady()) { + this.#announced = true; + printListening(this.options, this.url); + this.#ready.resolve(); + } + } + }); + + child.on("error", (error) => { + console.error(`Cluster worker ${slot} error:`, error); + }); + + child.on("exit", (code, signal) => { + if (this.#closing) { + return; + } + + this.#workers.delete(slot); + + if (!state.ready && restarts >= MAX_START_ATTEMPTS - 1) { + const error = new Error( + `Cluster worker ${slot} failed to start after ${MAX_START_ATTEMPTS} attempts (exited with ${signal || code}).`, + ); + + this.#fatal(error); + return; + } + + const uptime = Date.now() - state.startedAt; + const nextRestarts = uptime > STABLE_UPTIME ? 1 : restarts + 1; + const delay = Math.min(100 * 2 ** nextRestarts, 30_000); + + this.#log( + c.yellow, + `Cluster worker ${slot} exited unexpectedly (${signal || code}), restarting in ${delay}ms...`, + ); + + const timer = setTimeout(() => { + this.#respawnTimers.delete(timer); + this.#spawn(slot, nextRestarts); + }, delay); + + this.#respawnTimers.add(timer); + }); + } + + /** + * Forks a worker process that re-executes the current entry. + * + * Node workers are created with `node:cluster` so they share the listening + * handle (round-robin on all platforms). Bun and Deno workers are plain forks + * that bind the port themselves with `SO_REUSEPORT`. + * + * @param slot Worker slot index, exposed to the worker via `SRVX_CLUSTER_WORKER`. + */ + #forkWorker(slot: number): ChildProcess { + const workerEnv = { [CLUSTER_WORKER_ENV]: String(slot) }; + + if (this.runtime === "node") { + return nodeCluster.fork(workerEnv).process; + } + + return fork(process.argv[1], process.argv.slice(2), { + env: { ...process.env, ...workerEnv }, + // fork's default, passed explicitly: runtime flags like --import must reach the workers + execArgv: process.execArgv, + }); + } + + #allReady(): boolean { + if (this.#workers.size < this.#size) { + return false; + } + + for (const worker of this.#workers.values()) { + if (!worker.ready) { + return false; + } + } + + return true; + } + + #onSignal = (signal: NodeJS.Signals): void => { + this.#closing = true; + this.#clearRespawns(); + + for (const { child } of this.#workers.values()) { + child.kill(signal); + } + }; + + /** + * Handles an unrecoverable startup failure: stops all workers and exits with + * a non-zero code, so process managers (Docker, K8s, ...) see the failed + * start instead of an idle supervisor. + */ + #fatal(error: Error): void { + console.error(c.red(error.message)); + this.#ready.reject(error); + this.#killAll(true).then(() => process.exit(1)); + } + + #clearRespawns(): void { + for (const timer of this.#respawnTimers) { + clearTimeout(timer); + } + + this.#respawnTimers.clear(); + } + + #log(color: (t: string) => string, message: string): void { + if (!(this.options.silent ?? globalThis.process?.env?.TEST)) { + console.log(color(message)); + } + } +} diff --git a/src/_plugins.ts b/src/_plugins.ts index 5206d7a1..437dabb5 100644 --- a/src/_plugins.ts +++ b/src/_plugins.ts @@ -33,11 +33,18 @@ export const gracefulShutdownPlugin: ServerPlugin = (server) => { const w = server.options.silent ? () => {} : process.stderr.write.bind(process.stderr); + const exitClusterWorker = () => { + if (process.env.SRVX_CLUSTER_WORKER) { + process.exit(0); + } + }; + const forceClose = async () => { if (isClosed) return; w(c.red("\x1b[2K\rForcibly closing connections...\n")); isClosed = true; await server.close(true); + exitClusterWorker(); }; const shutdown = async () => { @@ -68,6 +75,7 @@ export const gracefulShutdownPlugin: ServerPlugin = (server) => { if (closed) { w("\x1b[2K\r" + c.green("Server closed successfully.\n")); isClosed = true; + exitClusterWorker(); return; } } diff --git a/src/adapters/bun.ts b/src/adapters/bun.ts index 23b8b7b0..bbc533c9 100644 --- a/src/adapters/bun.ts +++ b/src/adapters/bun.ts @@ -10,12 +10,13 @@ import { } from "../_utils.ts"; import { wrapFetch } from "../_middleware.ts"; import { gracefulShutdownPlugin } from "../_plugins.ts"; +import { withCluster } from "../_cluster.ts"; export { FastURL } from "../_url.ts"; export const FastResponse: typeof globalThis.Response = Response; -export function serve(options: ServerOptions): BunServer { - return new BunServer(options); +export function serve(options: ServerOptions): Server { + return withCluster(options, (options) => new BunServer(options)); } // https://bun.sh/docs/api/http diff --git a/src/adapters/deno.ts b/src/adapters/deno.ts index 2cfe06f4..786552ae 100644 --- a/src/adapters/deno.ts +++ b/src/adapters/deno.ts @@ -10,12 +10,13 @@ import { import { wrapFetch } from "../_middleware.ts"; import { gracefulShutdownPlugin } from "../_plugins.ts"; import { limitRequestBody } from "../_body-limit.ts"; +import { withCluster } from "../_cluster.ts"; export { FastURL } from "../_url.ts"; export const FastResponse: typeof globalThis.Response = Response; -export function serve(options: ServerOptions): DenoServer { - return new DenoServer(options); +export function serve(options: ServerOptions): Server { + return withCluster(options, (options) => new DenoServer(options)); } // https://docs.deno.com/api/deno/~/Deno.serve diff --git a/src/adapters/node.ts b/src/adapters/node.ts index 55004a94..1c637959 100644 --- a/src/adapters/node.ts +++ b/src/adapters/node.ts @@ -9,6 +9,7 @@ import { } from "../_utils.ts"; import { wrapFetch } from "../_middleware.ts"; import { errorPlugin, gracefulShutdownPlugin } from "../_plugins.ts"; +import { withCluster } from "../_cluster.ts"; import nodeHTTP from "node:http"; import nodeHTTPS from "node:https"; @@ -36,7 +37,7 @@ export { toNodeHandler, toFetchHandler } from "./_node/adapter.ts"; export type { AdapterMeta } from "./_node/adapter.ts"; export function serve(options: ServerOptions): Server { - return new NodeServer(options); + return withCluster(options, (options) => new NodeServer(options)); } // https://nodejs.org/api/http.html diff --git a/src/cli/main.ts b/src/cli/main.ts index d3f94dc4..bc137aa7 100644 --- a/src/cli/main.ts +++ b/src/cli/main.ts @@ -44,6 +44,11 @@ export async function main(mainOpts: MainOptions): Promise { // Log versions console.log(c.gray([...versions(mainOpts), cliOpts.prod ? "prod" : "dev"].join(" · "))); + // Cluster mode is production-only (dev mode uses a single process with watcher) + if (!cliOpts.prod && (cliOpts.cluster || process.env.SRVX_WORKERS)) { + console.log(c.yellow("Cluster mode is only available in production mode (--prod), ignoring.")); + } + // Resolve .env files const envFiles = [".env", cliOpts.prod ? ".env.production" : ".env.local"].filter((f) => existsSync(f), @@ -96,13 +101,18 @@ function parseArgs(args: string[]): CLIOptions { if (mode === "serve") { // Serve mode + // Support both `--cluster` (worker count = CPU cores) and `--cluster ` / `--cluster=` + const serveArgs = args.map((arg, i) => + arg === "--cluster" && !/^(\d+|true|false)$/.test(args[i + 1] || "") ? "--cluster=true" : arg, + ); const { values, positionals } = parseNodeArgs({ - args, + args: serveArgs, allowPositionals: true, options: { ...commonArgs, url: { type: "string" }, prod: { type: "boolean" }, + cluster: { type: "string" }, port: { type: "string", short: "p" }, static: { type: "string", short: "s" }, import: { type: "string" }, @@ -130,7 +140,25 @@ function parseArgs(args: string[]): CLIOptions { } } - return { mode, ...values }; + // Convert `--cluster` value: "true" enables with default size, a number sets + // the size, "false" disables cluster mode entirely (including SRVX_WORKERS) + let cluster: boolean | number | undefined; + if (values.cluster !== undefined) { + if (values.cluster === "true") { + cluster = true; + } else if (values.cluster === "false") { + cluster = false; + } else { + cluster = Number(values.cluster); + if (!Number.isInteger(cluster) || cluster <= 0) { + throw new Error( + `Invalid --cluster value: "${values.cluster}" (expected a positive integer, "true" or "false")`, + ); + } + } + } + + return { mode, ...values, cluster }; } // Fetch mode diff --git a/src/cli/serve.ts b/src/cli/serve.ts index 190b6497..d540fdbe 100644 --- a/src/cli/serve.ts +++ b/src/cli/serve.ts @@ -49,10 +49,13 @@ export async function cliServe(cliOpts: CLIOptions): Promise { ...loaded.module, } as Partial; - printInfo(cliOpts, loaded); + if (!process.env.SRVX_CLUSTER_WORKER) { + printInfo(cliOpts, loaded); + } server = srvxServe({ ...serverOptions, gracefulShutdown: !!cliOpts.prod, + cluster: cliOpts.prod ? (cliOpts.cluster ?? serverOptions.cluster) : false, port: cliOpts.port ?? serverOptions.port, hostname: cliOpts.hostname ?? cliOpts.host ?? serverOptions.hostname, tls: cliOpts.tls ? { cert: cliOpts.cert, key: cliOpts.key } : undefined, diff --git a/src/cli/types.ts b/src/cli/types.ts index 0bcb8661..7cd02018 100644 --- a/src/cli/types.ts +++ b/src/cli/types.ts @@ -35,6 +35,8 @@ export type CLIOptions = { /** Run in production mode (no watch, no debug) */ prod?: boolean; + /** Cluster mode: number of worker processes (true = CPU count), production only */ + cluster?: boolean | number; /** Serve static files from the specified directory (default: "public") */ static?: string; /** ES module to preload */ diff --git a/src/cli/usage.ts b/src/cli/usage.ts index ab6aafb7..5ff1f057 100644 --- a/src/cli/usage.ts +++ b/src/cli/usage.ts @@ -17,6 +17,7 @@ ${c.bold("SERVE MODE")} ${c.bold(c.green(`# ${command} serve [options]`))} ${c.gray("$")} ${c.cyan(command)} serve --entry ${c.gray("./server.ts")} ${c.gray("# Start development server")} ${c.gray("$")} ${c.cyan(command)} serve --prod ${c.gray("# Start production server")} +${c.gray("$")} ${c.cyan(command)} serve --prod --cluster ${c.gray("# Production server with one worker per CPU core")} ${c.gray("$")} ${c.cyan(command)} serve --port=8080 ${c.gray("# Listen on port 8080")} ${c.gray("$")} ${c.cyan(command)} serve --host=localhost ${c.gray("# Bind to localhost only")} ${c.gray("$")} ${c.cyan(command)} serve --static=./dist ${c.gray("# Serve static files (no entry needed)")} @@ -48,6 +49,7 @@ ${c.bold("SERVE OPTIONS")} ${c.green("--host")} ${c.yellow("")} Host to bind to (default: all interfaces) ${c.green("-s, --static")} ${c.yellow("")} Serve static files from the specified directory (default: ${c.yellow("public")}) ${c.green("--prod")} Run in production mode (no watch, no debug) + ${c.green("--cluster")} ${c.yellow("[N]")} Run N server processes (default: CPU cores, requires --prod) ${c.green("--import")} ${c.yellow("")} ES module to preload ${c.green("--tls")} Enable TLS (HTTPS/HTTP2) ${c.green("--cert")} ${c.yellow("")} TLS certificate file @@ -64,6 +66,7 @@ ${c.bold("ENVIRONMENT")} ${c.green("PORT")} Override port ${c.green("HOST")} Override host + ${c.green("SRVX_WORKERS")} Number of cluster workers (enables cluster mode) ${c.green("NODE_ENV")} Set to ${c.yellow("production")} for production mode. ${mainOpts.usage?.docs ? `➤ ${c.url("Documentation", mainOpts.usage.docs)}` : ""} diff --git a/src/types.ts b/src/types.ts index 0ca9b24f..6cbe7b2c 100644 --- a/src/types.ts +++ b/src/types.ts @@ -99,6 +99,25 @@ export interface ServerOptions { */ reusePort?: boolean; + /** + * Run multiple server processes sharing the same port (cluster mode). + * + * The main process becomes a small supervisor that spawns workers (re-executing the same entry), + * restarts them when they crash and forwards shutdown signals. Set to a number for an exact + * worker count, or `true` to read it from the `SRVX_WORKERS` environment variable (number of + * CPU cores when unset). `SRVX_WORKERS` alone also enables cluster mode; `false` disables it entirely. + * + * Worker processes can be detected via the `SRVX_CLUSTER_WORKER` environment variable (worker index, starting at `"0"`). + * + * Supported on Node.js (`node:cluster`, all platforms) and on Bun and Deno (`SO_REUSEPORT`, load balancing on Linux only). + * Serverless runtimes scale processes themselves and ignore this option. + * + * **Note:** Cluster mode requires a fixed port (`port: 0` is not supported). + * + * @see https://srvx.h3.dev/guide/cluster + */ + cluster?: boolean | number; + /** * The protocol to use for the server. * diff --git a/test/cluster.test.ts b/test/cluster.test.ts new file mode 100644 index 00000000..0b4caedb --- /dev/null +++ b/test/cluster.test.ts @@ -0,0 +1,183 @@ +import { describe, it, expect, beforeAll, afterAll } from "vitest"; +import { fileURLToPath } from "node:url"; +import { existsSync } from "node:fs"; +import { spawnSync } from "node:child_process"; +import { get as httpGet } from "node:http"; +import { execa, type ResultPromise as ExecaRes } from "execa"; +import { getRandomPort, waitForPort } from "get-port-please"; + +const fixtureEntry = fileURLToPath(new URL("fixtures/cluster/server.mjs", import.meta.url)); +const cliBin = fileURLToPath(new URL("../bin/srvx.mjs", import.meta.url)); + +const isLinux = process.platform === "linux"; + +// The CLI bin runs from dist, runtime suites need their binary available +const hasDist = existsSync(fileURLToPath(new URL("../dist/cli.mjs", import.meta.url))); +const hasBin = (bin: string) => { + try { + return spawnSync(bin, ["--version"], { stdio: "ignore" }).status === 0; + } catch { + return false; + } +}; + +// Plain HTTP GET without keep-alive: a pooled fetch() connection would stick to +// a single worker and hide the round-robin distribution. +function fetchWorker(url: string): Promise<{ pid: number; worker: string | null }> { + return new Promise((resolve, reject) => { + const req = httpGet(url, { agent: false, timeout: 2000 }, (res) => { + let body = ""; + res.on("data", (chunk) => (body += chunk)); + res.on("end", () => { + try { + expect(res.statusCode).toBe(200); + resolve(JSON.parse(body)); + } catch (error) { + reject(error as Error); + } + }); + }); + req.on("error", reject); + req.on("timeout", () => req.destroy(new Error("request timeout"))); + }); +} + +// Collect worker pids until `expected` distinct ones are seen (or timeout). +async function collectPids(url: string, expected: number, timeout = 8000): Promise> { + const pids = new Set(); + const deadline = Date.now() + timeout; + while (pids.size < expected && Date.now() < deadline) { + try { + pids.add((await fetchWorker(url)).pid); + } catch { + // worker may be restarting; retry + await new Promise((resolve) => setTimeout(resolve, 100)); + } + } + return pids; +} + +function testClusterExec(cmd: string[], opts: { workers: number; lb: boolean }) { + let childProc: ExecaRes; + let url: string; + + beforeAll(async () => { + const port = await getRandomPort("localhost"); + url = `http://localhost:${port}/`; + childProc = execa(cmd[0], cmd.slice(1), { + env: { PORT: port.toString(), SRVX_TEST_CLUSTER: String(opts.workers) }, + }); + childProc.catch(() => {}); // killed with SIGTERM on teardown + if (process.env.TEST_DEBUG) { + childProc.stdout!.on("data", (chunk) => console.log(chunk.toString())); + childProc.stderr!.on("data", (chunk) => console.log(chunk.toString())); + } + await waitForPort(port, { host: "localhost", delay: 50, retries: 100 }); + }); + + afterAll(async () => { + await childProc.kill(); + }); + + it("serves requests from cluster workers", async () => { + const { worker } = await fetchWorker(url); + expect(worker).not.toBeNull(); + }); + + if (opts.lb) { + it("load balances across all workers", async () => { + const pids = await collectPids(url, opts.workers); + expect(pids.size).toBe(opts.workers); + }); + + it("restarts crashed workers", async () => { + const { pid: killedPid } = await fetchWorker(url); + process.kill(killedPid, "SIGKILL"); + // Eventually all worker slots serve again, including a fresh process + const deadline = Date.now() + 10_000; + let pids = new Set(); + while (Date.now() < deadline) { + pids = await collectPids(url, opts.workers, 2000); + if (pids.size === opts.workers && !pids.has(killedPid)) { + break; + } + } + expect(pids.size).toBe(opts.workers); + expect(pids.has(killedPid)).toBe(false); + }); + } + + it("shuts down all workers on SIGTERM", async () => { + childProc.kill("SIGTERM"); + await childProc.catch(() => {}); + await expect(fetch(url, { signal: AbortSignal.timeout(1000) })).rejects.toThrow(); + }); +} + +describe("cluster (node, programmatic)", () => { + // node:cluster round-robin load balances on all platforms + testClusterExec(["node", fixtureEntry], { workers: 2, lb: true }); +}); + +describe("cluster (node, options)", () => { + async function spawnFixture(env: Record) { + const port = await getRandomPort("localhost"); + const proc = execa("node", [fixtureEntry], { env: { PORT: String(port), ...env } }); + proc.catch(() => {}); // killed with SIGTERM on teardown + return { proc, port, url: `http://localhost:${port}/` }; + } + + it("SRVX_WORKERS env enables cluster mode", async () => { + const { proc, port, url } = await spawnFixture({ SRVX_WORKERS: "2" }); + try { + await waitForPort(port, { host: "localhost", delay: 50, retries: 100 }); + const pids = await collectPids(url, 2); + expect(pids.size).toBe(2); + } finally { + await proc.kill(); + } + }); + + it("cluster: false disables cluster mode including SRVX_WORKERS", async () => { + const { proc, port, url } = await spawnFixture({ + SRVX_WORKERS: "2", + SRVX_TEST_CLUSTER: "false", + }); + try { + await waitForPort(port, { host: "localhost", delay: 50, retries: 100 }); + const { worker } = await fetchWorker(url); + expect(worker).toBeNull(); + } finally { + await proc.kill(); + } + }); + + it("supervisor exits with a non-zero code when workers fail to start", async () => { + const port = await getRandomPort("localhost"); + const result = await execa("node", [fixtureEntry], { + env: { PORT: String(port), SRVX_TEST_CLUSTER: "2", SRVX_TEST_CRASH: "1" }, + reject: false, + timeout: 12_000, + }); + expect(result.timedOut).toBe(false); + expect(result.exitCode).toBe(1); + }); +}); + +describe.skipIf(!hasDist)("cluster (node, cli)", () => { + // --host=localhost: waitForPort can not detect a wildcard listener on macOS + testClusterExec( + ["node", cliBin, "--prod", "--cluster=2", "--host=localhost", "--entry", fixtureEntry], + { workers: 2, lb: true }, + ); +}); + +describe.skipIf(!hasBin("bun"))("cluster (bun)", () => { + // SO_REUSEPORT load balancing is Linux-only for Bun + testClusterExec(["bun", "run", fixtureEntry], { workers: 2, lb: isLinux }); +}); + +describe.skipIf(!hasBin("deno"))("cluster (deno)", () => { + // SO_REUSEPORT load balancing is Linux-only for Deno (single worker elsewhere) + testClusterExec(["deno", "run", "-A", fixtureEntry], { workers: 2, lb: isLinux }); +}); diff --git a/test/fixtures/cluster/server.mjs b/test/fixtures/cluster/server.mjs new file mode 100644 index 00000000..459ba334 --- /dev/null +++ b/test/fixtures/cluster/server.mjs @@ -0,0 +1,22 @@ +// Minimal cluster entry: the supervisor re-executes this file for each worker. +const runtime = globalThis.Deno ? "deno" : globalThis.Bun ? "bun" : "node"; +const { serve } = await import(`../../../src/adapters/${runtime}.ts`); + +// Simulate a worker that can never start (crash-loop / fatal startup tests) +if (process.env.SRVX_TEST_CRASH && process.env.SRVX_CLUSTER_WORKER) { + process.exit(7); +} + +// `SRVX_TEST_CLUSTER`: worker count, "false" to disable, unset to leave the +// `cluster` option out (e.g. to test `SRVX_WORKERS` env activation). +const testCluster = process.env.SRVX_TEST_CLUSTER; + +serve({ + cluster: testCluster === "false" ? false : testCluster ? Number(testCluster) : undefined, + hostname: "localhost", + fetch: () => + Response.json({ + pid: process.pid, + worker: process.env.SRVX_CLUSTER_WORKER ?? null, + }), +}); From 7fc21ac3041368ccadeb2b29706bbabe6285d2f0 Mon Sep 17 00:00:00 2001 From: Oskar Lebuda Date: Mon, 6 Jul 2026 23:37:30 +0200 Subject: [PATCH 2/6] feat: add cluster mode --- src/_cluster.ts | 26 +++++++++++++++++++------- 1 file changed, 19 insertions(+), 7 deletions(-) diff --git a/src/_cluster.ts b/src/_cluster.ts index 4e9aedd4..8b80ab17 100644 --- a/src/_cluster.ts +++ b/src/_cluster.ts @@ -22,6 +22,12 @@ const MAX_START_ATTEMPTS = 3; type ReadyMessage = { srvx?: string; url?: string }; +type WorkerState = { + child: ChildProcess; + ready: boolean; + startedAt: number; +}; + /** * Adds cluster (multi-process) support to a runtime adapter's `serve()`. * @@ -43,6 +49,7 @@ export function withCluster( if (!env || (globalThis as any).__srvxLoader__) { return factory(options); } + // Worker process: start a regular server on the shared port if (env[CLUSTER_WORKER_ENV]) { // Deno supports SO_REUSEPORT on Linux only (non-Linux runs a single @@ -58,10 +65,12 @@ export function withCluster( ); return server; } + const workers = resolveClusterSize(options); if (workers === undefined) { return factory(options); } + return new ClusterServer(options, workers); } @@ -75,24 +84,23 @@ export function withCluster( * @returns The worker count, or `undefined` when cluster mode is disabled. */ function resolveClusterSize(options: ServerOptions): number | undefined { + // Cluster mode disabled if (options.cluster === false || options.cluster === 0) { return; } + + // Cluster size from environment variable const envSize = Number.parseInt(globalThis.process?.env?.SRVX_WORKERS || "", 10) || undefined; if (!options.cluster && !envSize) { return; } + + // Cluster size from options const size = typeof options.cluster === "number" ? options.cluster : (envSize ?? availableParallelism()); return Math.max(1, Math.floor(size)); } -type WorkerState = { - child: ChildProcess; - ready: boolean; - startedAt: number; -}; - /** * Cluster supervisor implementing the `Server` interface. * @@ -126,10 +134,12 @@ class ClusterServer implements Server { "Cluster mode requires a server entry file (cannot re-execute this process).", ); } + const { port, hostname } = resolvePortAndHost(options); if (!port) { throw new Error("Cluster mode requires a fixed port (port: 0 is not supported)."); } + const secure = !!(options.tls?.cert || options.protocol === "https"); this.#fallbackURL = fmtURL(hostname || "localhost", port, secure); @@ -159,13 +169,16 @@ class ClusterServer implements Server { } this.#log(c.gray, `Starting ${this.#size} cluster worker${this.#size > 1 ? "s" : ""}...`); + for (let slot = 0; slot < this.#size; slot++) { this.#spawn(slot, 0); } + for (const signal of ["SIGINT", "SIGTERM"] as const) { process.on(signal, this.#onSignal); } } + return this.ready(); } @@ -320,7 +333,6 @@ class ClusterServer implements Server { return fork(process.argv[1], process.argv.slice(2), { env: { ...process.env, ...workerEnv }, - // fork's default, passed explicitly: runtime flags like --import must reach the workers execArgv: process.execArgv, }); } From b1e16eedd447aa239baadc0b608062796da3febf Mon Sep 17 00:00:00 2001 From: "autofix-ci[bot]" <114827586+autofix-ci[bot]@users.noreply.github.com> Date: Mon, 6 Jul 2026 21:40:56 +0000 Subject: [PATCH 3/6] chore: apply automated updates --- src/_cluster.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/_cluster.ts b/src/_cluster.ts index 8b80ab17..78e1d79d 100644 --- a/src/_cluster.ts +++ b/src/_cluster.ts @@ -173,7 +173,7 @@ class ClusterServer implements Server { for (let slot = 0; slot < this.#size; slot++) { this.#spawn(slot, 0); } - + for (const signal of ["SIGINT", "SIGTERM"] as const) { process.on(signal, this.#onSignal); } From 1df76b0993316eb2833af9de8c658113df9e3d3d Mon Sep 17 00:00:00 2001 From: Oskar Lebuda Date: Tue, 7 Jul 2026 22:37:01 +0200 Subject: [PATCH 4/6] feat: add cluster mode --- src/cli/main.ts | 16 ++++++++++------ test/cluster.test.ts | 3 ++- 2 files changed, 12 insertions(+), 7 deletions(-) diff --git a/src/cli/main.ts b/src/cli/main.ts index bc137aa7..9cee323f 100644 --- a/src/cli/main.ts +++ b/src/cli/main.ts @@ -1,7 +1,7 @@ import { parseArgs as parseNodeArgs } from "node:util"; import { fileURLToPath } from "node:url"; import { fork } from "node:child_process"; -import { existsSync, statSync } from "node:fs"; +import { existsSync, readFileSync, statSync } from "node:fs"; import * as c from "./_utils.ts"; import type { CLIOptions, MainOptions } from "./types.ts"; import { cliServe, NO_ENTRY_ERROR } from "./serve.ts"; @@ -44,11 +44,6 @@ export async function main(mainOpts: MainOptions): Promise { // Log versions console.log(c.gray([...versions(mainOpts), cliOpts.prod ? "prod" : "dev"].join(" · "))); - // Cluster mode is production-only (dev mode uses a single process with watcher) - if (!cliOpts.prod && (cliOpts.cluster || process.env.SRVX_WORKERS)) { - console.log(c.yellow("Cluster mode is only available in production mode (--prod), ignoring.")); - } - // Resolve .env files const envFiles = [".env", cliOpts.prod ? ".env.production" : ".env.local"].filter((f) => existsSync(f), @@ -59,6 +54,15 @@ export async function main(mainOpts: MainOptions): Promise { ); } + if ( + !cliOpts.prod && + (cliOpts.cluster || + process.env.SRVX_WORKERS || + envFiles.some((f) => /^\s*SRVX_WORKERS\s*=/m.test(readFileSync(f, "utf8")))) + ) { + console.log(c.yellow("Cluster mode is only available in production mode (--prod), ignoring.")); + } + // In prod mode without --import, run directly in current process (no fork needed) if (cliOpts.prod && !cliOpts.import) { // Load env files manually since we're not forking with --env-file args diff --git a/test/cluster.test.ts b/test/cluster.test.ts index 0b4caedb..225a340a 100644 --- a/test/cluster.test.ts +++ b/test/cluster.test.ts @@ -76,7 +76,8 @@ function testClusterExec(cmd: string[], opts: { workers: number; lb: boolean }) }); afterAll(async () => { - await childProc.kill(); + childProc.kill("SIGTERM"); + await childProc.catch(() => {}); }); it("serves requests from cluster workers", async () => { From 0beeab46518e94197239f57bb6569796de2ac8d5 Mon Sep 17 00:00:00 2001 From: Oskar Lebuda Date: Tue, 7 Jul 2026 22:51:21 +0200 Subject: [PATCH 5/6] feat: add cluster mode --- src/_cluster.ts | 19 ++++++++++++++++++- 1 file changed, 18 insertions(+), 1 deletion(-) diff --git a/src/_cluster.ts b/src/_cluster.ts index 78e1d79d..7a1287ea 100644 --- a/src/_cluster.ts +++ b/src/_cluster.ts @@ -1,5 +1,5 @@ import nodeCluster from "node:cluster"; -import { fork, type ChildProcess } from "node:child_process"; +import { fork, spawn, type ChildProcess } from "node:child_process"; import { availableParallelism } from "node:os"; import * as c from "./cli/_utils.ts"; import { fmtURL, printListening, resolvePortAndHost } from "./_utils.ts"; @@ -331,6 +331,23 @@ class ClusterServer implements Server { return nodeCluster.fork(workerEnv).process; } + // Deno workers bind with SO_REUSEPORT (Linux-only), which Deno gates behind + // `--unstable-net` and enforces with a hard process exit. fork() cannot pass + // Deno CLI flags (execArgv is translated as Node/V8 flags), so spawn the + // worker with explicit `deno run` arguments and an IPC channel instead. + // Compiled (standalone) binaries have unstable config baked in and accept + // no `run` subcommand, so they keep using fork(). + if (IS_DENO && process.platform === "linux" && !(globalThis as any).Deno?.build?.standalone) { + return spawn( + process.execPath, + ["run", "--unstable-net", "-A", process.argv[1], ...process.argv.slice(2)], + { + env: { ...process.env, ...workerEnv }, + stdio: ["inherit", "inherit", "inherit", "ipc"], + }, + ); + } + return fork(process.argv[1], process.argv.slice(2), { env: { ...process.env, ...workerEnv }, execArgv: process.execArgv, From 7856edbde04c62e20d5f652579d4f407c33f8b07 Mon Sep 17 00:00:00 2001 From: Oskar Lebuda Date: Sun, 6 Sep 2026 02:06:04 +0200 Subject: [PATCH 6/6] fix(cluster): address review feedback - workers force `manual: false`: with `manual` left on, `ready()` resolves without a listener and the supervisor announces a cluster nobody can reach - node workers pin `exclusive: false` explicitly, so an entry's own `node` options cannot opt a worker out of the `node:cluster` shared handle - CLI honors `cluster` from an entry's own intercepted `serve()` call - docs/types: Node load balances on all platforms but round-robin is not the Windows default (`SCHED_NONE`); `cluster: ` is a positive integer and `0` disables cluster mode --- docs/1.guide/05.options.md | 2 +- docs/1.guide/12.cluster.md | 2 +- src/_cluster.ts | 27 ++++++++++++++++++--------- src/cli/serve.ts | 2 +- src/types.ts | 5 +++-- test/cluster.test.ts | 2 +- 6 files changed, 25 insertions(+), 15 deletions(-) diff --git a/docs/1.guide/05.options.md b/docs/1.guide/05.options.md index 35e670ed..a8099120 100644 --- a/docs/1.guide/05.options.md +++ b/docs/1.guide/05.options.md @@ -63,7 +63,7 @@ Enabling this option allows multiple processes to bind to the same port, which i Run multiple server processes sharing the same port. - `true`: number of workers from the `SRVX_WORKERS` environment variable, or CPU cores. -- `number`: exact number of worker processes. +- `number`: number of worker processes (a positive integer; `0` disables cluster mode). - `false`: explicitly disable cluster mode (also ignores `SRVX_WORKERS`). :read-more{to="/guide/cluster"} diff --git a/docs/1.guide/12.cluster.md b/docs/1.guide/12.cluster.md index dc572f49..9a24dd52 100644 --- a/docs/1.guide/12.cluster.md +++ b/docs/1.guide/12.cluster.md @@ -62,7 +62,7 @@ On `SIGINT`/`SIGTERM`, the supervisor forwards the signal to all workers for gra ## Runtime support -- **Node.js**: workers share the listening socket via [`node:cluster`](https://nodejs.org/api/cluster.html) with round-robin load balancing (all platforms, including macOS and Windows). +- **Node.js**: workers share the listening socket via [`node:cluster`](https://nodejs.org/api/cluster.html), so load balancing works on every platform. The distribution follows Node's own scheduling policy: round-robin everywhere except Windows, where Node defaults to `SCHED_NONE` and lets the OS pick the worker (override with `NODE_CLUSTER_SCHED_POLICY=rr`). - **Bun** and **Deno**: workers bind the port with `SO_REUSEPORT` and the kernel load balances between them. Kernel load balancing is **Linux only** - on other platforms a single supervised worker is started instead (crash restarts still work). - **Serverless runtimes** (Cloudflare, AWS Lambda, ...): the platform scales processes itself, the `cluster` option is ignored. diff --git a/src/_cluster.ts b/src/_cluster.ts index 716ada2f..3e89efe7 100644 --- a/src/_cluster.ts +++ b/src/_cluster.ts @@ -52,19 +52,27 @@ export function withCluster( // Worker process: start a regular server on the shared port if (env[CLUSTER_WORKER_ENV]) { - // Deno supports SO_REUSEPORT on Linux only (non-Linux runs a single - // supervised worker that binds the port exclusively). - const reusePort = !(IS_DENO && process.platform !== "linux"); + const isNodeRuntime = !IS_BUN && !IS_DENO; + // Bun and Deno workers bind the shared port themselves with SO_REUSEPORT. + // Deno supports it on Linux only (non-Linux runs a single supervised worker + // that binds the port exclusively). + const reusePort = !isNodeRuntime && !(IS_DENO && process.platform !== "linux"); const server = factory({ ...options, cluster: false, reusePort, + // A worker must start listening as soon as the supervisor spawns it: with + // `manual` left on it would report ready without ever accepting a request. + manual: false, silent: true, // Node workers get the listening handle from the `node:cluster` primary, - // so they only need a non-exclusive bind: asking for SO_REUSEPORT on top - // is pointless and throws ENOTSUP on platforms that reject it for the - // resolved address (e.g. `::1` on macOS). - ...(IS_BUN || IS_DENO ? undefined : { node: { ...options.node, reusePort: false } }), + // so the bind has to stay non-exclusive; SO_REUSEPORT on top is pointless + // and throws ENOTSUP on platforms that reject it for the resolved address + // (e.g. `::1` on macOS). Set explicitly so an entry's own `node` options + // cannot opt a worker out of the shared handle. + ...(isNodeRuntime + ? { node: { ...options.node, exclusive: false, reusePort: false } } + : undefined), }); Promise.resolve(server.ready()).then( () => process.send?.({ srvx: "cluster-worker-ready", url: server.url }), @@ -169,7 +177,7 @@ class ClusterServer implements Server { // SO_REUSEPORT load balancing is Linux-only for Bun/Deno: fall back to a // single supervised worker instead of spawning processes that would never - // receive connections (Node uses node:cluster round-robin on all platforms). + // receive connections (Node load balances via node:cluster on all platforms). if (process.platform !== "linux" && (IS_BUN || IS_DENO) && this.#size > 1) { this.#log( c.yellow, @@ -329,7 +337,8 @@ class ClusterServer implements Server { * Forks a worker process that re-executes the current entry. * * Node workers are created with `node:cluster` so they share the listening - * handle (round-robin on all platforms). Bun and Deno workers are plain forks + * handle, which load balances on all platforms (round-robin except on Windows, + * where Node defaults to `SCHED_NONE`). Bun and Deno workers are plain forks * that bind the port themselves with `SO_REUSEPORT`. * * @param slot Worker slot index, exposed to the worker via `SRVX_CLUSTER_WORKER`. diff --git a/src/cli/serve.ts b/src/cli/serve.ts index 70f0538b..5e368811 100644 --- a/src/cli/serve.ts +++ b/src/cli/serve.ts @@ -109,7 +109,7 @@ export async function cliServe(cliOpts: CLIOptions): Promise { server = srvxServe({ ...serverOptions, gracefulShutdown: !!cliOpts.prod, - cluster: cliOpts.prod ? (cliOpts.cluster ?? serverOptions.cluster) : false, + cluster: cliOpts.prod ? (cliOpts.cluster ?? serverOptions.cluster ?? nested?.cluster) : false, port: cliOpts.port ?? serverOptions.port, hostname: cliOpts.hostname ?? cliOpts.host ?? serverOptions.hostname, tls, diff --git a/src/types.ts b/src/types.ts index 20e4b141..501e44eb 100644 --- a/src/types.ts +++ b/src/types.ts @@ -110,9 +110,10 @@ export interface ServerOptions { * Run multiple server processes sharing the same port (cluster mode). * * The main process becomes a small supervisor that spawns workers (re-executing the same entry), - * restarts them when they crash and forwards shutdown signals. Set to a number for an exact + * restarts them when they crash and forwards shutdown signals. Set to a positive integer for the * worker count, or `true` to read it from the `SRVX_WORKERS` environment variable (number of - * CPU cores when unset). `SRVX_WORKERS` alone also enables cluster mode; `false` disables it entirely. + * CPU cores when unset). `SRVX_WORKERS` alone also enables cluster mode; `false` (or `0`) + * disables it entirely. * * Worker processes can be detected via the `SRVX_CLUSTER_WORKER` environment variable (worker index, starting at `"0"`). * diff --git a/test/cluster.test.ts b/test/cluster.test.ts index 225a340a..da8d3058 100644 --- a/test/cluster.test.ts +++ b/test/cluster.test.ts @@ -116,7 +116,7 @@ function testClusterExec(cmd: string[], opts: { workers: number; lb: boolean }) } describe("cluster (node, programmatic)", () => { - // node:cluster round-robin load balances on all platforms + // node:cluster load balances on all platforms (round-robin except on Windows) testClusterExec(["node", fixtureEntry], { workers: 2, lb: true }); });