diff --git a/src/js/node/worker_threads.ts b/src/js/node/worker_threads.ts index 1658bd7cb4eb..b0b40107e19a 100644 --- a/src/js/node/worker_threads.ts +++ b/src/js/node/worker_threads.ts @@ -76,6 +76,8 @@ const { 14: MessageChannel, 15: BroadcastChannel, 16: WebWorker, + 17: _createJSTransferableClaim, + 18: _claimJSTransferable, } = $cpp("Worker.cpp", "createNodeWorkerThreadsBinding") as [ unknown, number, @@ -97,6 +99,9 @@ const { // instance. This is so that it can emit the `worker` event on the process with the // node:worker_threads instance instead of the Web Worker instance. new (...args: [...ConstructorParameters, nodeWorker: Worker]) => WebWorker, + () => number, + // true: the id was live, and the caller now owns the resource in transit. + (claim: unknown) => boolean, ]; type NodeWorkerOptions = import("node:worker_threads").WorkerOptions; @@ -355,10 +360,7 @@ function setupWorkerStdio(stdio) { // on receive, markers are swapped back for reconstructed instances. // A plain string key on purpose: Symbols don't survive structured clone, and // Bun has no native HostObject hook, so the marker must ride along inside the -// cloned graph (including Map/Set entries). This is in-band signaling: a user -// object that fabricates the key in workerData will deserialize on the worker -// side where node would deliver it unchanged. That's accepted - it is not a -// privilege boundary (worker threads share the parent's fd table anyway). +// cloned graph (including Map/Set entries). const kJSTransferableMarker = "__bunNodeWorkerJSTransferable"; function isJSTransferableMarker(value: object): boolean { @@ -368,13 +370,37 @@ function isJSTransferableMarker(value: object): boolean { ); } +// The first thread to take a marker's claim id owns its fd, so one fd is never closed or restored twice. +function closeIfUnclaimed(data: unknown, claim: number) { + if (!_claimJSTransferable(claim)) return; + const fd = (data as any)?.fd; + if (typeof fd !== "number" || fd < 0) return; + try { + require("node:fs").closeSync(fd); + } catch { + // already closed + } +} + +// Module scope: a closure made in packJSTransferables() would pin the user's workerData graph until the worker exits. +function makeReclaimJSTransferables(pending: Array<[data: unknown, claim: number]>) { + return function reclaimJSTransferables() { + for (const { 0: data, 1: claim } of pending) closeIfUnclaimed(data, claim); + }; +} + +// undefined: user data that only looks like a marker. It stays plain data and must not throw in the worker bootstrap. function deserializeJSTransferable(marker: Record): unknown { const deserializeInfo = marker[kJSTransferableMarker]; switch (deserializeInfo) { case "internal/fs/promises:FileHandle": { + const data = marker.data; + if (data === null || typeof data !== "object" || typeof data.fd !== "number") return undefined; const { FileHandle, kDeserialize } = require("node:fs").promises.$data; const handle = new FileHandle(-1); - handle[kDeserialize](marker.data); + // Claim last: a terminate() between the claim and kDeserialize() leaks the fd, so the node:fs load stays above. + if (!_claimJSTransferable(marker.claim)) return undefined; + handle[kDeserialize](data); return handle; } default: @@ -394,8 +420,10 @@ function unpackJSTransferables(value: unknown, memo?: Map): unk if (cached !== undefined) return cached; if (isJSTransferableMarker(value)) { const instance = deserializeJSTransferable(value as Record); - memo.set(value, instance); - return instance; + if (instance !== undefined) { + memo.set(value, instance); + return instance; + } } memo.set(value, value); if ($isArray(value)) { @@ -434,6 +462,7 @@ function unpackJSTransferables(value: unknown, memo?: Map): unk const kRestoreJSTransferables = Symbol("kRestoreJSTransferables"); const kFinalizeJSTransferables = Symbol("kFinalizeJSTransferables"); +const kReclaimJSTransferables = Symbol("kReclaimJSTransferables"); function packJSTransferables(options: NodeWorkerOptions): NodeWorkerOptions { const transferList = options?.transferList; @@ -460,9 +489,10 @@ function packJSTransferables(options: NodeWorkerOptions): NodeWorkerOptions { // kTransfer() neuters the handle (extracts the bare fd); if anything later // in the pack/construct sequence throws, restore the already-neutered // handles so their fds aren't orphaned. - const neutered: Array<[item: any, data: unknown]> = []; + const neutered: Array<[item: any, data: unknown, claim: number]> = []; function restoreNeutered() { - for (const { 0: item, 1: data } of neutered) { + for (const { 0: item, 1: data, 2: claim } of neutered) { + if (!_claimJSTransferable(claim)) continue; try { item[kDeserialize](data); } catch { @@ -485,10 +515,12 @@ function packJSTransferables(options: NodeWorkerOptions): NodeWorkerOptions { const extraTransfers = item[kTransferList]?.(); // May throw DataCloneError (e.g. FileHandle in use); propagate synchronously like Node. const { data, deserializeInfo } = item[kTransfer](); - neutered.push([item, data]); + const claim = _createJSTransferableClaim(); + neutered.push([item, data, claim]); (replacements ??= new Map()).set(item, { [kJSTransferableMarker]: deserializeInfo, data, + claim, }); if ($isArray(extraTransfers)) nativeTransferList.push(...extraTransfers); } else { @@ -570,16 +602,14 @@ function packJSTransferables(options: NodeWorkerOptions): NodeWorkerOptions { // rollback above must still find the fd open to restore the handle (node // leaves the handle fully usable in that case). packed[kFinalizeJSTransferables] = function finalizeJSTransferables() { - for (const { 0: item, 1: data } of neutered) { - if (!usedMarkers.has(item) && typeof (data as any)?.fd === "number" && (data as any).fd >= 0) { - try { - require("node:fs").closeSync((data as any).fd); - } catch { - // already closed - } - } + for (const { 0: item, 1: data, 2: claim } of neutered) { + if (!usedMarkers.has(item)) closeIfUnclaimed(data, claim); } }; + // Runs on 'close', which the worker thread posts after its VM is destroyed, so every claim is final by then. + const pending: Array<[data: unknown, claim: number]> = []; + for (const { 1: data, 2: claim } of neutered) pending.push([data, claim]); + packed[kReclaimJSTransferables] = makeReclaimJSTransferables(pending); return packed; } @@ -820,6 +850,7 @@ class Worker extends EventEmitter { #urlToRevoke = ""; // threadId captured for cleaning up the messaging control port on close. #messagingThreadId: number | undefined = undefined; + #reclaimJSTransferables: (() => void) | undefined = undefined; constructor(filename: string, options: NodeWorkerOptions = {}) { super(); @@ -952,6 +983,7 @@ class Worker extends EventEmitter { // The transfer is committed - release fds that were transferred but are // not referenced from workerData (nothing will deserialize them). options[kFinalizeJSTransferables]?.(); + this.#reclaimJSTransferables = options[kReclaimJSTransferables]; // Tracing active (CLI flag or dynamic enable): record the Node-style // `[worker N] ` thread-name metadata event. No-op when tracing is // off — the agent module is a tiny one-time load. @@ -1170,6 +1202,9 @@ class Worker extends EventEmitter { messaging.destroyMainThreadPort(this.#messagingThreadId); this.#messagingThreadId = undefined; } + // node closes an undelivered fd too: https://github.com/nodejs/node/blob/v26.3.0/src/node_file.cc#L318-L327 + this.#reclaimJSTransferables?.(); + this.#reclaimJSTransferables = undefined; // End captured stdio readables when the worker exits, even if it was // terminated before its own streams finished. if (this.#stdout) { diff --git a/src/jsc/bindings/webcore/Worker.cpp b/src/jsc/bindings/webcore/Worker.cpp index 15b0b9ff83e8..90e33d6f572f 100644 --- a/src/jsc/bindings/webcore/Worker.cpp +++ b/src/jsc/bindings/webcore/Worker.cpp @@ -37,6 +37,10 @@ #include #include #include +#include +#include +#include +#include #include "SerializedScriptValue.h" #include "ScriptExecutionContext.h" #include @@ -286,6 +290,38 @@ JSC_DEFINE_HOST_FUNCTION(jsFunctionSetEntryEvaluatedHook, (JSC::JSGlobalObject * return JSC::JSValue::encode(jsUndefined()); } +// Ids of resources in transit to a worker (worker_threads.ts). Process-wide: the parent VM makes one, the worker VM takes it. +static Lock s_transferClaimsLock; +static HashSet& transferClaims() WTF_REQUIRES_LOCK(s_transferClaimsLock) +{ + static NeverDestroyed> claims; + return claims.get(); +} + +JSC_DEFINE_HOST_FUNCTION(jsFunctionCreateJSTransferableClaim, (JSGlobalObject*, CallFrame*)) +{ + Locker locker { s_transferClaimsLock }; + uint64_t id; + do { + // 53 bits, so the id is exact as a JS number. 0 is the empty value of the set. + id = cryptographicallyRandomNumber() >> 11; + } while (!id || !transferClaims().add(id).isNewEntry); + return JSValue::encode(jsNumber(static_cast(id))); +} + +// true: the id was live and the caller now owns the resource. Anything else that user data can hold is false. +JSC_DEFINE_HOST_FUNCTION(jsFunctionClaimJSTransferable, (JSGlobalObject*, CallFrame* callFrame)) +{ + JSValue value = callFrame->argument(0); + if (!value.isNumber()) + return JSValue::encode(jsBoolean(false)); + double number = value.asNumber(); + if (!(number >= 1 && number < 9007199254740992.0) || number != std::trunc(number)) + return JSValue::encode(jsBoolean(false)); + Locker locker { s_transferClaimsLock }; + return JSValue::encode(jsBoolean(transferClaims().remove(static_cast(number)))); +} + JSValue createNodeWorkerThreadsBinding(Zig::GlobalObject* globalObject) { VM& vm = globalObject->vm(); @@ -342,7 +378,7 @@ JSValue createNodeWorkerThreadsBinding(Zig::GlobalObject* globalObject) bool isNodeWorker = proxy && proxy->options().kind == WorkerOptions::Kind::Node; - JSObject* array = constructEmptyArray(globalObject, nullptr, 17); + JSObject* array = constructEmptyArray(globalObject, nullptr, 19); RETURN_IF_EXCEPTION(scope, {}); array->putDirectIndex(globalObject, 0, workerData); RETURN_IF_EXCEPTION(scope, {}); @@ -380,6 +416,10 @@ JSValue createNodeWorkerThreadsBinding(Zig::GlobalObject* globalObject) RETURN_IF_EXCEPTION(scope, {}); array->putDirectIndex(globalObject, 16, JSWorker::getConstructor(vm, globalObject)); RETURN_IF_EXCEPTION(scope, {}); + array->putDirectIndex(globalObject, 17, JSFunction::create(vm, globalObject, 0, "createJSTransferableClaim"_s, jsFunctionCreateJSTransferableClaim, ImplementationVisibility::Public, NoIntrinsic)); + RETURN_IF_EXCEPTION(scope, {}); + array->putDirectIndex(globalObject, 18, JSFunction::create(vm, globalObject, 1, "claimJSTransferable"_s, jsFunctionClaimJSTransferable, ImplementationVisibility::Public, NoIntrinsic)); + RETURN_IF_EXCEPTION(scope, {}); return array; } diff --git a/test/js/node/worker_threads/worker_threads.test.ts b/test/js/node/worker_threads/worker_threads.test.ts index 809f6376c75c..92b59e33817f 100644 --- a/test/js/node/worker_threads/worker_threads.test.ts +++ b/test/js/node/worker_threads/worker_threads.test.ts @@ -1037,6 +1037,180 @@ test("FileHandles nested in Map and Set workerData are transferred", async () => expect(message).toEqual({ sameInstance: true, text: "hello" }); }); +// These tests watch descriptor numbers, so they stay serial. +describe("the fd of a FileHandle transferred through workerData", () => { + // False once the descriptor is closed, or once its number belongs to another file. + function refersTo(fd: number, file: fs.Stats) { + try { + const now = fs.fstatSync(fd); + return now.ino === file.ino && now.dev === file.dev; + } catch (e: any) { + if (e.code !== "EBADF") throw e; + return false; + } + } + + // Disposal closes what a failed expectation leaves open: the handle if it still owns the fd, else the bare fd. + async function openToTransfer(path: string) { + const fh = await fs.promises.open(path, "r"); + const fd = fh.fd; + const file = fs.fstatSync(fd); + return { + fh, + fd, + isOpen: () => refersTo(fd, file), + async [Symbol.asyncDispose]() { + if (fh.fd !== -1) await fh.close(); + else if (refersTo(fd, file)) fs.closeSync(fd); + }, + }; + } + + // open() returns the lowest free descriptor, so these take the closed number `fd` again. A second close hits one. + function reopenUpTo(fd: number, path: string) { + const held = [fs.openSync(path, "r")]; + const file = fs.fstatSync(held[0]); + while (held.at(-1)! < fd) held.push(fs.openSync(path, "r")); + return { + closedByOthers: () => held.filter(descriptor => !refersTo(descriptor, file)), + [Symbol.dispose]() { + for (const descriptor of held) if (refersTo(descriptor, file)) fs.closeSync(descriptor); + }, + }; + } + + test("is closed when the worker entry does not resolve", async () => { + using dir = tempDir("worker-fh-undelivered", { "x.txt": "hello" }); + await using transferred = await openToTransfer(join(String(dir), "x.txt")); + const { fh } = transferred; + const worker = new Worker(join(String(dir), "missing.js"), { workerData: { fh }, transferList: [fh as any] }); + const errors: string[] = []; + worker.on("error", error => errors.push(error.code)); + const code = await new Promise(resolve => worker.on("exit", resolve)); + expect({ errors, code, parentFd: fh.fd, open: transferred.isOpen() }).toEqual({ + errors: ["MODULE_NOT_FOUND"], + code: 1, + parentFd: -1, + open: false, + }); + }); + + // A thread that wins the race against terminate() receives the handle and owns the fd, so that attempt is repeated. + test("is closed when terminate() stops the worker before it starts", async () => { + using dir = tempDir("worker-fh-undelivered", { "x.txt": "hello" }); + let closed = false; + for (let attempt = 0; attempt < 10 && !closed; attempt++) { + // The worker is gone at disposal, so nothing else can close the fd of an attempt it won. + await using transferred = await openToTransfer(join(String(dir), "x.txt")); + const { fh } = transferred; + const worker = new Worker("setInterval(() => {}, 1000)", { + eval: true, + workerData: { fh }, + transferList: [fh as any], + }); + await worker.terminate(); + expect(fh.fd).toBe(-1); + closed = !transferred.isOpen(); + } + expect(closed).toBe(true); + }); + + // The constructor closes a handle that workerData does not reference. The worker's exit must not close it again. + test("is closed only once when workerData does not reference the handle", async () => { + using dir = tempDir("worker-fh-unreferenced", { "x.txt": "hello", "y.txt": "world" }); + const other = join(String(dir), "y.txt"); + // A starting worker thread opens descriptors of its own. It takes these lower numbers, not the handle's. + const parked = Array.from({ length: 16 }, () => fs.openSync(other, "r")); + await using transferred = await openToTransfer(join(String(dir), "x.txt")); + for (const descriptor of parked) fs.closeSync(descriptor); + const { fh } = transferred; + const worker = new Worker(join(String(dir), "missing.js"), { workerData: {}, transferList: [fh as any] }); + const closedByConstructor = !transferred.isOpen(); + using reopened = reopenUpTo(transferred.fd, other); + const errors: string[] = []; + worker.on("error", error => errors.push(error.code)); + const code = await new Promise(resolve => worker.on("exit", resolve)); + expect({ errors, code, closedByConstructor, closedAgain: reopened.closedByOthers() }).toEqual({ + errors: ["MODULE_NOT_FOUND"], + code: 1, + closedByConstructor: true, + closedAgain: [], + }); + }); + + // Fails when the worker's claim does not reach the parent: the parent then closes a number that is another file's. + test("is not closed by the parent when the worker received the handle", async () => { + using dir = tempDir("worker-fh-delivered", { "x.txt": "hello", "y.txt": "world" }); + await using transferred = await openToTransfer(join(String(dir), "x.txt")); + const { fh } = transferred; + await using worker = new Worker( + `const { workerData, parentPort } = require("node:worker_threads"); + parentPort.once("message", () => {}); + workerData.fh.close().then(() => parentPort.postMessage("closed"));`, + { eval: true, workerData: { fh }, transferList: [fh as any] }, + ); + const [message] = await once(worker, "message"); + using reopened = reopenUpTo(transferred.fd, join(String(dir), "y.txt")); + worker.postMessage("exit"); + const [code] = await once(worker, "exit"); + expect({ message, code, closedByParent: reopened.closedByOthers() }).toEqual({ + message: "closed", + code: 0, + closedByParent: [], + }); + }); + + // The claim is a plain number, so the transfer needs neither the SharedArrayBuffer global nor the JSC option. + test("is transferred when SharedArrayBuffer is not available", async () => { + using dir = tempDir("worker-fh-no-sab", { "x.txt": "hello" }); + await using proc = Bun.spawn({ + cmd: [ + bunExe(), + "-e", + `const { Worker } = require("node:worker_threads"); + const fs = require("node:fs"); + globalThis.SharedArrayBuffer = globalThis.Atomics = undefined; + fs.promises.open("x.txt", "r").then(fh => { + const worker = new Worker( + \`const { workerData, parentPort } = require("node:worker_threads"); + workerData.fh.readFile("utf8").then(text => workerData.fh.close().then(() => parentPort.postMessage(text)));\`, + { eval: true, workerData: { fh }, transferList: [fh] }, + ); + worker.on("message", text => console.log(JSON.stringify({ text, parentFd: fh.fd }))); + });`, + ], + env: { ...bunEnv, BUN_JSC_useSharedArrayBuffer: "0" }, + 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, exitCode }).toEqual({ + stdout: JSON.stringify({ text: "hello", parentFd: -1 }), + stderr: "", + exitCode: 0, + }); + }); + + // The worker unpacks workerData before user code runs, so a lookalike must not throw there. + test.each([ + ["no claim id", { data: { fd: 1 << 20 } }], + ["a claim id that is not live", { data: { fd: 1 << 20 }, claim: 12345 }], + ["a claim id of the wrong type", { data: { fd: 1 << 20 }, claim: "12345" }], + ["no fd", { data: {}, claim: 12345 }], + ["no data", { claim: 12345 }], + ])("workerData that imitates a marker with %s stays plain data", async (_, rest) => { + const imitation = { __bunNodeWorkerJSTransferable: "internal/fs/promises:FileHandle", ...rest }; + await using worker = new Worker( + `const { workerData, parentPort } = require("node:worker_threads"); + parentPort.postMessage(workerData);`, + { eval: true, workerData: { imitation } }, + ); + const [message] = await once(worker, "message"); + expect(message).toEqual({ imitation }); + }); +}); + test("MessagePort.hasRef() reports actual loop-ref state", () => { const { port1 } = new MessageChannel(); expect(port1.hasRef()).toBe(false);