diff --git a/src/js/internal/cluster/RoundRobinHandle.ts b/src/js/internal/cluster/RoundRobinHandle.ts index c5e85176bea7..5a82a25b79b2 100644 --- a/src/js/internal/cluster/RoundRobinHandle.ts +++ b/src/js/internal/cluster/RoundRobinHandle.ts @@ -110,6 +110,12 @@ export default class RoundRobinHandle { return this.all.has(worker.id); } + // The descriptor that the listener of this handle holds, -1 before it listens and after remove() closed it. + get fd() { + const fd = this.server?._handle?.fd; + return typeof fd === "number" ? fd : -1; + } + // With the channel still up the unacked newconn is settled by its ack; once it is gone, a crashed worker's goes to another worker and a disconnected worker's (already settled by it) is dropped. remove(worker, channelGone = false) { if (channelGone) { diff --git a/src/js/internal/cluster/SharedHandle.ts b/src/js/internal/cluster/SharedHandle.ts index 1f2c2a3c1c6b..76933aee8eff 100644 --- a/src/js/internal/cluster/SharedHandle.ts +++ b/src/js/internal/cluster/SharedHandle.ts @@ -1,6 +1,6 @@ const clusterRawBind = $newRustFunction("node_cluster_binding.rs", "clusterRawBind", 4); const closeRawHandle = $newRustFunction("node_cluster_binding.rs", "clusterCloseHandle", 1); -const validateFd = $newRustFunction("node_cluster_binding.rs", "clusterValidateFd", 1); +const validateFd = $newRustFunction("node_cluster_binding.rs", "clusterValidateFd", 2); export default class SharedHandle { key; @@ -19,7 +19,7 @@ export default class SharedHandle { this.sharedOnly = sharedOnly === true; if (typeof fd === "number" && fd >= 0) { - const err = validateFd(fd); + const err = validateFd(fd, addressType === "udp4" || addressType === "udp6"); if (err !== 0) { this.errno = err; } else { @@ -47,6 +47,12 @@ export default class SharedHandle { return this.workers.has(worker.id); } + // The descriptor that remove() closes: the number a worker named, or the socket that clusterRawBind made. + get fd() { + const handle = this.handle; + return handle ? handle.fd : -1; + } + remove(worker) { const workers = this.workers; if (!workers.has(worker.id)) return false; diff --git a/src/js/internal/cluster/primary.ts b/src/js/internal/cluster/primary.ts index 9165c71762f8..7b881bdae6af 100644 --- a/src/js/internal/cluster/primary.ts +++ b/src/js/internal/cluster/primary.ts @@ -7,11 +7,13 @@ const { kHandle } = require("internal/shared"); const sendHelper = $newRustFunction("node_cluster_binding.rs", "sendHelperPrimary", 4); const onInternalMessage = $newRustFunction("node_cluster_binding.rs", "onInternalMessagePrimary", 3); -const { UV_EINVAL, UV_ENOBUFS } = process.binding("uv"); +const { UV_EEXIST, UV_EINVAL, UV_ENOBUFS } = process.binding("uv"); let child_process; let RoundRobinHandle; let SharedHandle; +let validateFd; +let isFdOfDgramSocket; const ArrayPrototypeSlice = Array.prototype.slice; const ObjectValues = Object.values; @@ -45,7 +47,7 @@ const SCHED_RR = 2; export default cluster; -const handles = new Map(); +const handles = new Map(); cluster.isWorker = false; cluster.isMaster = true; // Deprecated alias. Must be same as isPrimary. cluster.isPrimary = true; @@ -240,6 +242,21 @@ function exitedAfterDisconnect(worker, message) { send(worker, { ack: message.seq }); } +// A handle closes its descriptor when its last worker leaves, so a second handle for a held descriptor closes it under the holder. Looked up by number: `udp4` and `udp6` are two keys for one descriptor. +function errnoOfFdQuery(message) { + const fd = message.fd; + const isUdp = message.addressType === "udp4" || message.addressType === "udp6"; + let held = isUdp && (isFdOfDgramSocket ??= $newRustFunction("udp_socket.rs", "jsDgramIsFdAdopted", 1))(fd); + if (!held) { + handles.$forEach(handle => { + if (handle.fd === fd) held = true; + }); + } + if (!held) return 0; + // The kind comes before the holder, as in node. Not as in node v26.3.0: EEXIST also under SCHED_NONE and for udp4 then udp6, where its worker dies on an assertion (nodejs/node#64869) or node serves and closes the number two times. + return (validateFd ??= $newRustFunction("node_cluster_binding.rs", "clusterValidateFd", 2))(fd, isUdp) || UV_EEXIST; +} + function queryServer(worker, message) { // Stop processing if worker already disconnecting if (worker.exitedAfterDisconnect) return; @@ -294,6 +311,13 @@ function queryServer(worker, message) { worker.emit("error", error); return; } + if (typeof message.fd === "number" && message.fd >= 0) { + const errno = errnoOfFdQuery(message); + if (errno !== 0) { + send(worker, { errno, key, ack: message.seq, data: cachedHandle ? cachedHandle.data : message.data }, null); + return; + } + } if ( schedulingPolicy !== SCHED_RR || message.sharedOnly === true || diff --git a/src/runtime/node/node_cluster_binding.rs b/src/runtime/node/node_cluster_binding.rs index 8eb044595d0e..3c5825db6c86 100644 --- a/src/runtime/node/node_cluster_binding.rs +++ b/src/runtime/node/node_cluster_binding.rs @@ -731,6 +731,7 @@ pub(crate) fn cluster_raw_bind(global: &JSGlobalObject, frame: &CallFrame) -> Js } } +/// `(fd, wantDgram)` → 0 when `fd` can serve a query of that kind (datagram, else stream), or a negative errno. #[bun_jsc::host_fn] pub(crate) fn cluster_validate_fd(global: &JSGlobalObject, frame: &CallFrame) -> JsResult { let _ = global; @@ -741,12 +742,21 @@ pub(crate) fn cluster_validate_fd(global: &JSGlobalObject, frame: &CallFrame) -> #[cfg(not(windows))] { let fd = value.to_int32(); + // `fd: 3.5` is not descriptor 3: https://github.com/nodejs/node/blob/v26.3.0/lib/net.js#L1904-L1911 + if f64::from(fd) != value.as_number() { + return Ok(JSValue::js_number_from_int32(-bun_sys::UV_E::INVAL)); + } if fd < 0 { return Ok(JSValue::js_number_from_int32(-bun_sys::UV_E::BADF)); } + let wanted = if frame.argument(1).to_boolean() { + libc::SOCK_DGRAM + } else { + libc::SOCK_STREAM + }; let mut ty: libc::c_int = 0; let mut len = core::mem::size_of::() as libc::socklen_t; - // SAFETY: plain getsockopt on a caller-supplied fd; out-params are + // SAFETY: plain getsockopt on a caller-supplied fd; out-params are live locals, and `len` is the size of `ty`. let rc = unsafe { libc::getsockopt( fd, @@ -756,8 +766,8 @@ pub(crate) fn cluster_validate_fd(global: &JSGlobalObject, frame: &CallFrame) -> &raw mut len, ) }; - // node's createServerHandle: EINVAL for anything that cannot listen (e.g. a connected stdio socketpair), fd left untouched. - if rc != 0 || (ty != libc::SOCK_STREAM && ty != libc::SOCK_DGRAM) { + // A descriptor of another kind is EINVAL before node opens it: https://github.com/nodejs/node/blob/v26.3.0/lib/net.js#L1904-L1911 and https://github.com/nodejs/node/blob/v26.3.0/lib/internal/dgram.js#L67-L70 + if rc != 0 || ty != wanted { return Ok(JSValue::js_number_from_int32(-bun_sys::UV_E::INVAL)); } if ty == libc::SOCK_STREAM { @@ -771,6 +781,21 @@ pub(crate) fn cluster_validate_fd(global: &JSGlobalObject, frame: &CallFrame) -> if connected { return Ok(JSValue::js_number_from_int32(-bun_sys::UV_E::INVAL)); } + } else { + // SAFETY: sockaddr_storage is plain data; getsockname only writes within `name_len`. + let family = unsafe { + let mut name: libc::sockaddr_storage = bun_core::ffi::zeroed_unchecked(); + let mut name_len = + core::mem::size_of::() as libc::socklen_t; + if libc::getsockname(fd, (&raw mut name).cast(), &raw mut name_len) != 0 { + return Ok(JSValue::js_number_from_int32(-bun_sys::UV_E::INVAL)); + } + libc::c_int::from(name.ss_family) + }; + // https://github.com/nodejs/node/blob/v26.3.0/deps/uv/src/unix/tty.c#L458-L460 + if family != libc::AF_INET && family != libc::AF_INET6 { + return Ok(JSValue::js_number_from_int32(-bun_sys::UV_E::INVAL)); + } } Ok(JSValue::js_number_from_int32(0)) } diff --git a/test/js/node/cluster.test.ts b/test/js/node/cluster.test.ts index 3248646bd3b0..0d3fd677d2ad 100644 --- a/test/js/node/cluster.test.ts +++ b/test/js/node/cluster.test.ts @@ -1,4 +1,5 @@ -import { expect, test } from "bun:test"; +import { dlopen } from "bun:ffi"; +import { beforeAll, describe, expect, test } from "bun:test"; import { bunEnv, bunExe, @@ -7,10 +8,12 @@ import { isLinux, isWindows, joinP, + libcPathForDlopen, tempDir, tempDirWithFiles, tls as tlsCerts, } from "harness"; +import { closeSync } from "node:fs"; import net from "node:net"; test.concurrent("cloneable and transferable equals", async () => { @@ -932,21 +935,21 @@ if (cluster.isPrimary) { }, ); -test.skipIf(isWindows)("dgram worker releases a shared fd it failed to adopt", async () => { - using dir = tempDir("cluster-dgram-adopt-fail", { - "main.ts": ` +test.skipIf(isWindows)( + "dgram bind({ fd }) on a stream socket fails EINVAL like node and leaves the primary's server listening", + async () => { + using dir = tempDir("cluster-dgram-stream-fd", { + "main.ts": ` const cluster = require("node:cluster"); const dgram = require("node:dgram"); const net = require("node:net"); if (cluster.isPrimary) { - // A stream socket passes the primary's fd check but cannot be adopted as a dgram socket in the worker. const tcp = net.createServer().listen(0, "127.0.0.1", () => { const { port } = tcp.address(); const worker = cluster.fork(); worker.on("message", m => { - console.log("worker error code:", m.code); - // Refused once both processes closed their copy; a leaked copy in either keeps the socket accepting. + console.log("worker error:", JSON.stringify(m)); const probe = net.connect(port, "127.0.0.1"); probe.on("connect", () => { console.log("probe: connected"); probe.destroy(); finish(); }); probe.on("error", err => { console.log("probe:", err.code); finish(); }); @@ -961,22 +964,575 @@ if (cluster.isPrimary) { process.on("message", ({ fd }) => { const socket = dgram.createSocket("udp4"); socket.on("listening", () => process.send({ code: "listening" })); - socket.on("error", err => process.send({ code: err.code })); + socket.on("error", err => process.send({ code: err.code, errno: err.errno, syscall: err.syscall })); socket.bind({ fd }); }); } +`, + }); + await using proc = Bun.spawn({ + cmd: [bunExe(), "main.ts"], + env: bunEnv, + cwd: String(dir), + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect({ stdout: stdout.trim(), stderr }).toEqual({ + stdout: 'worker error: {"code":"EINVAL","errno":-22,"syscall":"open"}\nprobe: connected', + stderr: "", + }); + expect(exitCode).toBe(0); + }, +); + +// A socket of this process that a primary gets as a descriptor. The types of these sockets do not have `fd`. +async function socketForPrimary(sockets: DisposableStack, kind: "tcp" | "udp" | "unix dgram") { + if (kind === "tcp") { + const listener = Bun.listen({ hostname: "127.0.0.1", port: 0, socket: { data() {} } }); + sockets.defer(() => listener.stop(true)); + return listener as typeof listener & { fd: number }; + } + if (kind === "udp") { + const udp = await Bun.udpSocket({ hostname: "127.0.0.1", port: 0 }); + sockets.defer(() => udp.close()); + return udp as typeof udp & { fd: number }; + } + const libc = dlopen(libcPathForDlopen(), { socket: { args: ["int", "int", "int"], returns: "int" } }); + sockets.defer(() => libc.close()); + const AF_UNIX = 1; + const SOCK_DGRAM = 2; + const fd = libc.symbols.socket(AF_UNIX, SOCK_DGRAM, 0); + if (fd < 0) throw new Error("socket(AF_UNIX, SOCK_DGRAM, 0) failed"); + sockets.defer(() => closeSync(fd)); + return { fd, port: 0 }; +} + +test.skipIf(isWindows)("dgram worker releases a shared fd it failed to adopt", async () => { + using dir = tempDir("cluster-dgram-adopt-fail", { + "main.ts": ` +const cluster = require("node:cluster"); +const dgram = require("node:dgram"); +const fs = require("node:fs"); + +if (cluster.isPrimary) { + const worker = cluster.fork(); + worker.on("message", m => { + console.log("worker error code:", m.code); + try { + fs.fstatSync(3); + console.log("descriptor of the primary: open"); + } catch (err) { + console.log("descriptor of the primary:", err.code); + } + // Free once both processes closed their copy; a leaked copy in either keeps the port bound. + const probe = dgram.createSocket("udp4"); + probe.on("listening", () => { console.log("probe: listening"); probe.close(finish); }); + probe.on("error", err => { console.log("probe:", err.code); finish(); }); + probe.bind(Number(process.env.PORT), "127.0.0.1"); + }); + function finish() { + worker.kill(); + worker.on("exit", () => process.exit(0)); + } +} else { + // The primary shares a datagram socket only, and a worker adopts every one. So the handle names a file here. + const getServer = cluster._getServer; + cluster._getServer = (socket, options, callback) => + getServer(socket, options, (err, handle) => { + if (handle) handle.sharedFd = fs.openSync(__filename, "r"); + callback(err, handle); + }); + const socket = dgram.createSocket("udp4"); + socket.on("listening", () => process.send({ code: "listening" })); + socket.on("error", err => process.send({ code: err.code })); + socket.bind({ fd: 3 }); +} `, }); + using sockets = new DisposableStack(); + const udp = await socketForPrimary(sockets, "udp"); await using proc = Bun.spawn({ cmd: [bunExe(), "main.ts"], - env: bunEnv, + env: { ...bunEnv, PORT: String(udp.port) }, cwd: String(dir), - stdout: "pipe", - stderr: "pipe", + // Descriptor 3 of the primary. + stdio: ["ignore", "pipe", "pipe", udp.fd], }); + // The primary has its copy. A copy in this process keeps the port bound. + sockets.dispose(); const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); expect({ stdout: stdout.trim(), stderr }).toEqual({ - stdout: "worker error code: EINVAL\nprobe: ECONNREFUSED", + stdout: "worker error code: ENOTSOCK\ndescriptor of the primary: EBADF\nprobe: listening", + stderr: "", + }); + expect(exitCode).toBe(0); +}); + +// Runs as the primary and as its workers. SCENARIO is { policy, cases: [{ fd, port, steps }] }. The test made one +// socket for each case, and the primary has it as descriptor `fd`. A step is one of: +// { worker, kind, fd } the worker calls listen({ fd }) or bind({ fd }) on a new server or socket +// { worker, kind, fd, twice } one server calls listen({ fd }) two times before the first answer +// { worker, kind, port: 0 } the worker listens on a port, so the primary makes a socket (`fd: "created"`) +// { primary: kind } the primary binds a socket of its own to the descriptor +// { connect: true } the primary connects a client to the socket of the case +const fdQueryFixture = ` +const cluster = require("node:cluster"); +const fs = require("node:fs"); +const net = require("node:net"); +const path = require("node:path"); + +const scenario = JSON.parse(process.env.SCENARIO); +cluster.schedulingPolicy = scenario.policy === "SCHED_NONE" ? cluster.SCHED_NONE : cluster.SCHED_RR; + +function stateOf(fd) { + try { + fs.fstatSync(fd); + return "open"; + } catch (error) { + return error.code; + } +} + +function socketsOfPrimary() { + const fds = []; + for (let fd = 3; fd < 256; fd++) { + if (stateOf(fd) === "open" && fs.fstatSync(fd).isSocket()) fds.push(fd); + } + return fds; +} + +function bindInPrimary(kind, fd) { + const { promise, resolve, reject } = Promise.withResolvers(); + const socket = require("node:dgram").createSocket(kind); + socket.once("error", reject); + socket.bind({ fd }, () => resolve(socket)); + return promise; +} + +// Takes each free number up to the descriptor. If the primary closed the descriptor too early, a sentinel has its +// number, and the close by the holder then closes that sentinel. +function openSentinels(fd) { + const sentinels = []; + do sentinels.push(fs.openSync(__filename, "r")); + while (sentinels.at(-1) < fd); + return sentinels; +} + +function connectTo(port) { + const { promise, resolve } = Promise.withResolvers(); + let data = ""; + const client = net.connect(port, "127.0.0.1", () => client.end("ping")); + client.setEncoding("utf8"); + client.on("data", chunk => (data += chunk)); + client.on("error", error => resolve(error.code)); + client.on("close", () => resolve(data)); + return promise; +} + +async function primary() { + const count = 1 + Math.max(...scenario.cases.flatMap(({ steps }) => steps.map(step => step.worker ?? 0))); + const workers = []; + const pending = []; + let cases = true; + const start = i => { + workers[i] = cluster.fork(); + workers[i].on("message", message => pending[i].resolve(message.result)); + // A worker that left fails its case only. The next case has a new worker. + workers[i].on("exit", (code, signal) => { + if (!cases) return; + start(i); + pending[i]?.reject(new Error("worker " + i + " left: " + code + " " + signal)); + }); + }; + for (let i = 0; i < count; i++) start(i); + const tell = (i, command) => { + pending[i] = Promise.withResolvers(); + workers[i].send(command); + return pending[i].promise; + }; + const closeAll = async () => { + for (const i of workers.keys()) await tell(i, { close: true }); + }; + + for (const { fd, port, steps } of scenario.cases) { + const answers = []; + const own = []; + // The descriptor that the case is about: the one of the test, or the socket that the primary made. + let held = fd; + try { + for (const step of steps) { + if (step.connect) { + answers.push({ client: await connectTo(port) }); + } else if (step.primary) { + own.push(await bindInPrimary(step.primary, fd)); + answers.push({ primary: "listening", fd: stateOf(fd) }); + } else if (step.port === 0) { + const before = socketsOfPrimary(); + answers.push({ result: await tell(step.worker, step) }); + const created = socketsOfPrimary().filter(n => !before.includes(n)); + if (created.length !== 1) throw new Error("expected one new socket in the primary, found [" + created + "]"); + held = created[0]; + } else { + const result = await tell(step.worker, { ...step, fd: step.fd === "created" ? held : step.fd }); + answers.push({ result, fd: stateOf(held) }); + } + } + const sentinels = openSentinels(held); + await closeAll(); + const closed = sentinels.map(stateOf).find(state => state !== "open"); + console.log(JSON.stringify({ answers, left: { fd: stateOf(held), sentinels: closed ?? "open" } })); + } catch (error) { + process.exitCode = 1; + console.log(JSON.stringify({ answers, error: error.message })); + await closeAll(); + } + for (const socket of own) socket.close(); + } + + cases = false; + for (const worker of workers) { + const { promise, resolve } = Promise.withResolvers(); + worker.once("exit", resolve); + worker.disconnect(); + await promise; + } +} + +function make(kind) { + const tlsOptions = () => ({ + key: fs.readFileSync(path.join(__dirname, "key.pem")), + cert: fs.readFileSync(path.join(__dirname, "cert.pem")), + }); + switch (kind) { + case "net": + return net.createServer(socket => socket.resume().end("served")); + case "tls": + return require("node:tls").createServer(tlsOptions()); + default: + return require("node:dgram").createSocket(kind); + } +} + +function worker() { + const live = []; + process.on("message", async step => { + if (step.close) { + for (const target of live.splice(0)) { + const { promise, resolve } = Promise.withResolvers(); + target.close(resolve); + await promise; + } + process.send({ result: "closed" }); + return; + } + const target = make(step.kind); + const done = result => process.send({ result }); + target.once("error", error => done({ code: error.code, syscall: error.syscall, errno: error.errno })); + const listening = () => { + live.push(target); + done("listening"); + }; + if (step.kind === "udp4" || step.kind === "udp6") target.bind({ fd: step.fd }, listening); + else if (step.port === 0) target.listen(0, "127.0.0.1", listening); + else { + if (step.twice) target.listen({ fd: step.fd }); + target.listen({ fd: step.fd }, listening); + } + }); +} + +if (cluster.isPrimary) { + primary().catch(error => { + console.error(error); + process.exit(1); + }); +} else { + worker(); +} +`; + +type FdQueryStep = + | { worker: number; kind: string; fd: number | "created"; twice?: true } + | { worker: number; kind: string; port: 0 } + | { primary: string } + | { connect: true }; +type FdQueryCase = { + name: string; + socket: "tcp" | "udp" | "unix dgram"; + // `fd` is the descriptor of the case in the primary. + steps: (fd: number) => FdQueryStep[]; + answers: object[]; + // After the workers closed what they listened on. The holder of the descriptor closed it one time. + left?: { fd: string; sentinels: string }; +}; + +const served = { result: "listening", fd: "open" }; +const refused = (code: "EEXIST" | "EINVAL", syscall: "bind" | "open") => ({ + result: { code, syscall, errno: code === "EEXIST" ? -17 : -22 }, + fd: "open", +}); +const ask = (kind: string, fd: number | "created", worker = 0): FdQueryStep => ({ worker, kind, fd }); +const twoAsks = (first: string, second: string) => (fd: number) => [ask(first, fd), ask(second, fd)]; + +// Each answer is the answer of node v26.3.0 for the same fixture. A row with another answer says what node does. +const fdQueryCases: Record<"SCHED_NONE" | "SCHED_RR", FdQueryCase[]> = { + SCHED_NONE: [ + { + name: "net, then udp4 on the stream socket", + socket: "tcp", + steps: fd => [ask("net", fd), ask("udp4", fd), { connect: true }], + answers: [served, refused("EINVAL", "open"), { client: "served" }], + }, + { + name: "udp4, then net on the datagram socket", + socket: "udp", + steps: twoAsks("udp4", "net"), + answers: [served, refused("EINVAL", "bind")], + }, + { + name: "net on a datagram socket, then udp4", + socket: "udp", + steps: twoAsks("net", "udp4"), + answers: [refused("EINVAL", "bind"), served], + }, + { + name: "udp4 on a stream socket, then net", + socket: "tcp", + steps: twoAsks("udp4", "net"), + answers: [refused("EINVAL", "open"), served], + }, + { + name: "net on the descriptor plus 0.5, then net on the descriptor", + socket: "tcp", + steps: fd => [ask("net", fd + 0.5), ask("net", fd)], + answers: [refused("EINVAL", "bind"), served], + }, + { + name: "udp4 on a unix datagram socket", + socket: "unix dgram", + steps: fd => [ask("udp4", fd)], + answers: [refused("EINVAL", "open")], + left: { fd: "open", sentinels: "open" }, + }, + { + // node: the worker dies on ERR_INTERNAL_ASSERTION at the second answer (nodejs/node#64869). + name: "net, net, then net in a second worker", + socket: "tcp", + steps: fd => [ask("net", fd), ask("net", fd), { connect: true }, ask("net", fd, 1)], + answers: [served, refused("EEXIST", "bind"), { client: "served" }, served], + }, + { + // node: the worker dies on ERR_INTERNAL_ASSERTION at the second answer. + name: "udp4, udp4", + socket: "udp", + steps: twoAsks("udp4", "udp4"), + answers: [served, refused("EEXIST", "open")], + }, + { + // The refused net ask leaves the descriptor open, so the third ask is the second ask of "udp4, udp4". + // node: the worker dies on ERR_INTERNAL_ASSERTION at the third answer. + name: "udp4, net, then udp4 on the datagram socket", + socket: "udp", + steps: fd => [ask("udp4", fd), ask("net", fd), ask("udp4", fd)], + answers: [served, refused("EINVAL", "bind"), refused("EEXIST", "open")], + }, + { + // node: 'listening'. Its primary then closes the descriptor two times. + name: "udp4, then udp6 in a second worker", + socket: "udp", + steps: fd => [ask("udp4", fd), ask("udp6", fd, 1)], + answers: [served, refused("EEXIST", "open")], + }, + { + // node: 'listening'. Its primary then closes the socket two times. + name: "net on a port, then net on the socket that the primary made for it", + socket: "tcp", + steps: () => [{ worker: 0, kind: "net", port: 0 }, ask("net", "created")], + answers: [{ result: "listening" }, refused("EEXIST", "bind")], + }, + { + // node: 'listening' under this policy, EEXIST under SCHED_RR. The worker gave the first handle back, so the + // holder closed the descriptor, and a sentinel has its number. + name: "one net server that listens two times before the first answer", + socket: "tcp", + steps: fd => [{ ...ask("net", fd), twice: true }], + answers: [{ ...refused("EEXIST", "bind"), fd: "EBADF" }], + left: { fd: "open", sentinels: "open" }, + }, + { + name: "a udp4 socket of the primary, then udp4", + socket: "udp", + steps: fd => [{ primary: "udp4" }, ask("udp4", fd)], + answers: [{ primary: "listening", fd: "open" }, refused("EEXIST", "open")], + left: { fd: "open", sentinels: "open" }, + }, + ], + SCHED_RR: [ + { + name: "net, then udp4 on the stream socket", + socket: "tcp", + steps: fd => [ask("net", fd), ask("udp4", fd), { connect: true }], + answers: [served, refused("EINVAL", "open"), { client: "served" }], + }, + { + // A TLS server has a shared handle under this policy, so this row reaches the check that a net server skips. + name: "tls on a datagram socket, then udp4", + socket: "udp", + steps: twoAsks("tls", "udp4"), + answers: [refused("EINVAL", "bind"), served], + }, + { + // A net server skips that check. The listener of the primary refuses the number: Bun.listen takes an integer. + name: "net on the descriptor plus 0.5, then net on the descriptor", + socket: "tcp", + steps: fd => [ask("net", fd + 0.5), ask("net", fd), { connect: true }], + answers: [refused("EINVAL", "bind"), served, { client: "served" }], + }, + { + // With epoll the second listener of the primary fails, so this row passes without the lookup. Not with kqueue. + name: "net, net, then net in a second worker", + socket: "tcp", + steps: fd => [ask("net", fd), ask("net", fd), { connect: true }, ask("net", fd, 1)], + answers: [served, refused("EEXIST", "bind"), { client: "served" }, served], + }, + { name: "tls, tls", socket: "tcp", steps: twoAsks("tls", "tls"), answers: [served, refused("EEXIST", "bind")] }, + { + // node: its primary answers EEXIST, and its worker dies on a TypeError in tls.Server._setServerData(null). + name: "net, tls", + socket: "tcp", + steps: fd => [ask("net", fd), ask("tls", fd), { connect: true }], + answers: [served, refused("EEXIST", "bind"), { client: "served" }], + }, + { name: "tls, net", socket: "tcp", steps: twoAsks("tls", "net"), answers: [served, refused("EEXIST", "bind")] }, + { + // The kind comes before the holder. This row passes without the lookup. + name: "udp4, then net on the datagram socket", + socket: "udp", + steps: twoAsks("udp4", "net"), + answers: [served, refused("EINVAL", "bind")], + }, + { + name: "net on a port, then tls on the socket that the primary made for it", + socket: "tcp", + steps: () => [{ worker: 0, kind: "net", port: 0 }, ask("tls", "created")], + answers: [{ result: "listening" }, refused("EEXIST", "bind")], + }, + ], +}; + +async function runFdQueryCases(policy: "SCHED_NONE" | "SCHED_RR") { + using dir = tempDir("cluster-fd-query", { + "cert.pem": tlsCerts.cert, + "key.pem": tlsCerts.key, + "fixture.cjs": fdQueryFixture, + }); + using sockets = new DisposableStack(); + const cases: { fd: number; port: number; steps: FdQueryStep[] }[] = []; + const inherited: number[] = []; + for (const { socket, steps } of fdQueryCases[policy]) { + const { fd, port } = await socketForPrimary(sockets, socket); + inherited.push(fd); + // The primary has the sockets as descriptors 3, 4, 5 and so on. + cases.push({ fd: 3 + cases.length, port, steps: steps(3 + cases.length) }); + } + await using proc = Bun.spawn({ + cmd: [bunExe(), "fixture.cjs"], + env: { ...bunEnv, SCENARIO: JSON.stringify({ policy, cases }) }, + cwd: String(dir), + stdio: ["ignore", "pipe", "pipe", ...inherited], + }); + // The primary has its copies. A listener in this process takes the clients of the workers. + sockets.dispose(); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + return { lines: stdout.split("\n").filter(Boolean), stderr, exitCode }; +} + +describe.skipIf(isWindows).each(["SCHED_NONE", "SCHED_RR"] as const)( + "%s: a worker names a descriptor of the primary", + policy => { + // One primary and two workers run all the cases of a policy. + let run: Awaited>; + beforeAll(async () => { + run = await runFdQueryCases(policy); + }); + + test("the primary and its workers leave with no error", () => { + expect({ lines: run.lines.length, stderr: run.stderr, exitCode: run.exitCode }).toEqual({ + lines: fdQueryCases[policy].length, + stderr: "", + exitCode: 0, + }); + }); + + test.each(fdQueryCases[policy].map((row, i) => [row.name, row, i] as const))( + "%s", + (_, { answers, left = { fd: "EBADF", sentinels: "open" } }, i) => { + // The primary prints one line for each case. A line is missing when the primary stopped before the case. + expect(run.lines[i] === undefined ? undefined : JSON.parse(run.lines[i])).toEqual({ answers, left }); + }, + ); + }, +); + +// One worker asks for descriptor 3 in the order of ASKS and closes nothing. Then the primary disconnects it. +const fdDisconnectFixture = ` +const cluster = require("node:cluster"); +cluster.schedulingPolicy = cluster.SCHED_NONE; + +if (cluster.isPrimary) { + const answers = []; + const worker = cluster.fork(); + worker.on("message", line => { + if (line !== "asked") { + answers.push(line); + console.log(line); + return; + } + // After a wrong answer the worker can have two handles under one key. Then it never leaves on disconnect. + if (answers.join() === process.env.ANSWERS) worker.disconnect(); + else worker.kill(); + }); + worker.on("exit", (code, signal) => console.log("worker left:", code, signal)); +} else { + const asks = process.env.ASKS.split(","); + const ask = i => { + if (i === asks.length) return process.send("asked"); + const kind = asks[i]; + const target = kind === "net" ? require("node:net").createServer() : require("node:dgram").createSocket(kind); + const answer = result => { + process.send(kind + ": " + result); + ask(i + 1); + }; + target.once("error", error => answer(error.syscall + " " + error.code)); + if (kind === "net") target.listen({ fd: 3 }, () => answer("listening")); + else target.bind({ fd: 3 }, () => answer("listening")); + }; + ask(0); +} +`; + +// A worker has one handle for each key, and a disconnect closes the handles that it has. With two handles under one +// key the first one stays open, and the worker never leaves. +test.concurrent.skipIf(isWindows).each([ + ["udp", "udp4,net,udp4", ["udp4: listening", "net: bind EINVAL", "udp4: open EEXIST"]], + ["tcp", "net,udp4,net", ["net: listening", "udp4: open EINVAL", "net: bind EEXIST"]], +] as const)("a worker leaves on disconnect after it asked for a %s descriptor: %s", async (socket, asks, answers) => { + using dir = tempDir("cluster-fd-disconnect", { "fixture.cjs": fdDisconnectFixture }); + using sockets = new DisposableStack(); + const { fd } = await socketForPrimary(sockets, socket); + await using proc = Bun.spawn({ + cmd: [bunExe(), "fixture.cjs"], + env: { ...bunEnv, ASKS: asks, ANSWERS: answers.join() }, + cwd: String(dir), + // Descriptor 3 of the primary. + stdio: ["ignore", "pipe", "pipe", fd], + }); + // The primary has its copy. + sockets.dispose(); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect({ lines: stdout.trim().split("\n"), stderr }).toEqual({ + lines: [...answers, "worker left: 0 null"], stderr: "", }); expect(exitCode).toBe(0);