Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
71 changes: 53 additions & 18 deletions src/js/node/worker_threads.ts
Original file line number Diff line number Diff line change
Expand Up @@ -76,6 +76,8 @@ const {
14: MessageChannel,
15: BroadcastChannel,
16: WebWorker,
17: _createJSTransferableClaim,
18: _claimJSTransferable,
} = $cpp("Worker.cpp", "createNodeWorkerThreadsBinding") as [
unknown,
number,
Expand All @@ -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<typeof globalThis.Worker>, 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;
Expand Down Expand Up @@ -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 {
Expand All @@ -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<string, any>): 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:
Expand All @@ -394,8 +420,10 @@ function unpackJSTransferables(value: unknown, memo?: Map<object, unknown>): unk
if (cached !== undefined) return cached;
if (isJSTransferableMarker(value)) {
const instance = deserializeJSTransferable(value as Record<string, any>);
memo.set(value, instance);
return instance;
if (instance !== undefined) {
memo.set(value, instance);
return instance;
}
}
memo.set(value, value);
if ($isArray(value)) {
Expand Down Expand Up @@ -434,6 +462,7 @@ function unpackJSTransferables(value: unknown, memo?: Map<object, unknown>): unk

const kRestoreJSTransferables = Symbol("kRestoreJSTransferables");
const kFinalizeJSTransferables = Symbol("kFinalizeJSTransferables");
const kReclaimJSTransferables = Symbol("kReclaimJSTransferables");

function packJSTransferables(options: NodeWorkerOptions): NodeWorkerOptions {
const transferList = options?.transferList;
Expand All @@ -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 {
Expand All @@ -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 {
Expand Down Expand Up @@ -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;
}

Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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] <name>` thread-name metadata event. No-op when tracing is
// off — the agent module is a tiny one-time load.
Expand Down Expand Up @@ -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) {
Expand Down
42 changes: 41 additions & 1 deletion src/jsc/bindings/webcore/Worker.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,10 @@
#include <JavaScriptCore/ScriptCallStack.h>
#include <wtf/TZoneMallocInlines.h>
#include <wtf/Scope.h>
#include <wtf/CryptographicallyRandomNumber.h>
#include <wtf/HashSet.h>
#include <wtf/Lock.h>
#include <wtf/NeverDestroyed.h>
#include "SerializedScriptValue.h"
#include "ScriptExecutionContext.h"
#include <JavaScriptCore/JSMap.h>
Expand Down Expand Up @@ -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<uint64_t>& transferClaims() WTF_REQUIRES_LOCK(s_transferClaimsLock)
{
static NeverDestroyed<HashSet<uint64_t>> 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<uint64_t>() >> 11;
} while (!id || !transferClaims().add(id).isNewEntry);
return JSValue::encode(jsNumber(static_cast<double>(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<uint64_t>(number))));
}

JSValue createNodeWorkerThreadsBinding(Zig::GlobalObject* globalObject)
{
VM& vm = globalObject->vm();
Expand Down Expand Up @@ -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, {});
Expand Down Expand Up @@ -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;
}

Expand Down
Loading
Loading