Skip to content

Commit b21b5f1

Browse files
committed
net: support sending net.BoundSocket to threads and child processes
This adds support for transferring net.BoundSocket instances to other threads via the worker_threads postMessage() transfer list, and for sending them to child processes as the sendHandle argument of subprocess.send(), following on from the BoundSocket introduction. A BoundSocket reserves a port synchronously at construction time. Making it transferable means a port can be reserved on one thread or process and the bound (but not yet listening or connected) TCP handle handed off to another to listen or connect on, without racing on the bind. For threads, BoundSocket implements kTransfer/kTransferList/ kDeserialize, moving the underlying TCP handle with the same mechanism used for net.Socket and net.Server transfer. For child processes, the handleConversion entry reuses the same transfer protocol on the sending side and the same _TransferredBoundSocket deserialization path on the receiving side; the underlying transport is that of cluster's shared-handle scheduling: SCM_RIGHTS on Unix and WSADuplicateSocket on Windows, both of which carry bind state. In both cases the source instance is left in the adopted state: address(), fd() and close() throw ERR_SOCKET_HANDLE_ADOPTED. Transfer requires an un-adopted, open TCP handle, otherwise ERR_WORKER_HANDLE_NOT_TRANSFERABLE is thrown; pipe (path) binds are not transferable and throw ERR_INVALID_HANDLE_TYPE when sent over IPC. On the receiving side the local address is re-derived from the handle rather than trusted from serialized state. Signed-off-by: Guy Bedford <guybedford@gmail.com>
1 parent fdcb1de commit b21b5f1

9 files changed

Lines changed: 329 additions & 30 deletions

‎doc/api/child_process.md‎

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1552,7 +1552,7 @@ added: v0.5.9
15521552
15531553
* `message` {Object} A parsed JSON object or primitive value.
15541554
* `sendHandle` {Handle|undefined} `undefined` or a [`net.Socket`][],
1555-
[`net.Server`][], or [`dgram.Socket`][] object.
1555+
[`net.Server`][], [`net.BoundSocket`][], or [`dgram.Socket`][] object.
15561556
15571557
The `'message'` event is triggered when a child process uses
15581558
[`process.send()`][] to send messages.
@@ -1858,6 +1858,9 @@ subprocess.ref();
18581858
<!-- YAML
18591859
added: v0.5.9
18601860
changes:
1861+
- version: REPLACEME
1862+
pr-url: https://github.com/nodejs/node/pull/64725
1863+
description: '`net.BoundSocket` instances can now be sent.'
18611864
- version: v5.8.0
18621865
pr-url: https://github.com/nodejs/node/pull/5283
18631866
description: The `options` parameter, and the `keepOpen` option
@@ -1872,7 +1875,7 @@ changes:
18721875
18731876
* `message` {Object}
18741877
* `sendHandle` {Handle|undefined} `undefined`, or a [`net.Socket`][],
1875-
[`net.Server`][], or [`dgram.Socket`][] object.
1878+
[`net.Server`][], [`net.BoundSocket`][], or [`dgram.Socket`][] object.
18761879
* `options` {Object} The `options` argument, if present, is an object used to
18771880
parameterize the sending of certain types of handles. `options` supports
18781881
the following properties:
@@ -1939,12 +1942,18 @@ Applications should avoid using such messages or listening for
19391942
`'internalMessage'` events as it is subject to change without notice.
19401943
19411944
The optional `sendHandle` argument that may be passed to `subprocess.send()` is
1942-
for passing a TCP server or socket object to the child process. The child process will
1945+
for passing a TCP server, socket or [`net.BoundSocket`][] object to the child
1946+
process. The child process will
19431947
receive the object as the second argument passed to the callback function
19441948
registered on the [`'message'`][] event. Any data that is received
19451949
and buffered in the socket will not be sent to the child. Sending IPC sockets is
19461950
not supported on Windows.
19471951
1952+
Sending a `net.BoundSocket` moves its underlying TCP handle to the child
1953+
process, leaving the source instance in the adopted state as if it had been
1954+
adopted by a server or socket. The bound socket must not have been adopted or
1955+
closed, and pipe (`path`) binds cannot be sent.
1956+
19481957
The optional `callback` is a function that is invoked after the message is
19491958
sent but before the child process may have received it. The function is called with a
19501959
single argument: `null` on success, or an [`Error`][] object on failure.
@@ -2372,6 +2381,7 @@ or [`child_process.fork()`][].
23722381
[`child_process.spawnSync()`]: #child_processspawnsynccommand-args-options
23732382
[`dgram.Socket`]: dgram.md#class-dgramsocket
23742383
[`maxBuffer` and Unicode]: #maxbuffer-and-unicode
2384+
[`net.BoundSocket`]: net.md#class-netboundsocket
23752385
[`net.Server`]: net.md#class-netserver
23762386
[`net.Socket`]: net.md#class-netsocket
23772387
[`options.detached`]: #optionsdetached

‎doc/api/errors.md‎

Lines changed: 7 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -3004,8 +3004,9 @@ A call was made and the UDP subsystem was not running.
30043004
### `ERR_SOCKET_HANDLE_ADOPTED`
30053005

30063006
An operation was attempted on a [`BoundSocket`][] that had already been adopted
3007-
by a [`net.Server`][] or [`net.Socket`][]. Once a bound socket is adopted, its
3008-
`address()` and `close()` methods can no longer be used.
3007+
by a [`net.Server`][] or [`net.Socket`][], or transferred to another thread.
3008+
Once a bound socket is adopted or transferred, its `address()` and `close()`
3009+
methods can no longer be used.
30093010

30103011
<a id="ERR_SOURCE_MAP_CORRUPT"></a>
30113012

@@ -3573,9 +3574,10 @@ The `Response` that has been passed to `WebAssembly.compileStreaming` or to
35733574

35743575
### `ERR_WORKER_HANDLE_NOT_TRANSFERABLE`
35753576

3576-
An attempt was made to transfer a `net.Socket` or `net.Server` to another thread
3577-
via a `worker_threads` `postMessage()` call while it was not in a transferable
3578-
state, for example because it had already started reading or had buffered data.
3577+
An attempt was made to transfer a `net.Socket`, `net.Server` or
3578+
`net.BoundSocket` to another thread via a `worker_threads` `postMessage()` call
3579+
while it was not in a transferable state, for example because it had already
3580+
started reading, had buffered data, or had already been adopted.
35793581

35803582
<a id="ERR_WORKER_INIT_FAILED"></a>
35813583

‎doc/api/net.md‎

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -789,6 +789,13 @@ server.listen(8000);
789789
A listening [`net.Server`][] can be transferred the same way, which moves the
790790
listening socket itself (and its pending accept queue) to the receiving thread.
791791

792+
An un-adopted TCP [`BoundSocket`][] can also be transferred, which moves the
793+
bound (but not yet listening or connected) socket. This allows a port to be
794+
reserved synchronously on one thread and adopted by a server or outgoing
795+
connection on another. Pipe binds are not transferable. After the transfer, the
796+
source `BoundSocket` behaves as if it had been adopted: `address()`, `fd()` and
797+
`close()` throw [`ERR_SOCKET_HANDLE_ADOPTED`][].
798+
792799
### `new net.Socket([options])`
793800

794801
<!-- YAML
@@ -1718,6 +1725,13 @@ file system entry; abstract and TCP binds have none to remove.
17181725
When a pipe `BoundSocket` bound to a source `path` is adopted as a client, that
17191726
path is reported as the socket's `localAddress` once it connects.
17201727

1728+
An un-adopted TCP `BoundSocket` can be moved to another thread by listing it in
1729+
the `transferList` of a [`worker_threads`][] `postMessage()` call, see
1730+
[Transferring TCP handles to other threads][]. It can likewise be sent to a
1731+
child process as the `sendHandle` argument of [`subprocess.send()`][]. In both
1732+
cases the source is left in the adopted state. Pipe binds cannot be moved
1733+
either way.
1734+
17211735
When an adopted `BoundSocket` connects to a numeric IP literal, `connect(2)` is
17221736
issued synchronously, so [`socket.localAddress`][] is resolved once
17231737
[`socket.connect()`][] returns. Connection failures are still reported via a
@@ -2345,6 +2359,7 @@ net.isIPv6('fhqwhgads'); // returns false
23452359
[`socket.setTimeout()`]: #socketsettimeouttimeout-callback
23462360
[`socket.setTimeout(timeout)`]: #socketsettimeouttimeout-callback
23472361
[`stream.getDefaultHighWaterMark()`]: stream.md#streamgetdefaulthighwatermarkobjectmode
2362+
[`subprocess.send()`]: child_process.md#subprocesssendmessage-sendhandle-options-callback
23482363
[`worker_threads`]: worker_threads.md
23492364
[`writable.destroy()`]: stream.md#writabledestroyerror
23502365
[`writable.destroyed`]: stream.md#writabledestroyed

‎doc/api/worker_threads.md‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1238,7 +1238,7 @@ port2.postMessage(circularData);
12381238
```
12391239
12401240
`transferList` may be a list of {ArrayBuffer}, [`MessagePort`][],
1241-
[`FileHandle`][], {net.Server}, and {net.Socket} objects.
1241+
[`FileHandle`][], {net.Server}, {net.Socket}, and {net.BoundSocket} objects.
12421242
After transferring, they are not usable on the sending side of the channel
12431243
anymore (even if they are not contained in `value`).
12441244
@@ -1249,6 +1249,8 @@ freshly accepted or created TCP connection that has not yet started reading and
12491249
has no buffered data, otherwise `postMessage()` throws
12501250
`ERR_WORKER_HANDLE_NOT_TRANSFERABLE`. This makes it possible to accept
12511251
connections on one thread and distribute them across a pool of worker threads.
1252+
Transferring a {net.BoundSocket} moves an un-adopted pre-bound socket, so a
1253+
port can be reserved synchronously on one thread and adopted on another.
12521254
Only TCP handles are supported.
12531255
12541256
If `value` contains {SharedArrayBuffer} instances, those are accessible

‎lib/internal/child_process.js‎

Lines changed: 48 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -54,6 +54,10 @@ const { TCP } = internalBinding('tcp_wrap');
5454
const { TTY } = internalBinding('tty_wrap');
5555
const { UDP } = internalBinding('udp_wrap');
5656
const SocketList = require('internal/socket_list');
57+
const {
58+
kDeserialize,
59+
kTransfer,
60+
} = require('internal/worker/js_transferable');
5761
const { owner_symbol } = require('internal/async_hooks').symbols;
5862
const { convertToValidSignal } = require('internal/util');
5963
const { isArrayBufferView } = require('internal/util/types');
@@ -86,6 +90,24 @@ const kChannelHandle = Symbol('kChannelHandle');
8690
const kIsUsedAsStdio = Symbol('kIsUsedAsStdio');
8791
const kPendingMessages = Symbol('kPendingMessages');
8892

93+
// Store the handle after successfully sending it, so it can be closed when
94+
// the NODE_HANDLE_ACK is received. If the handle could not be sent, just
95+
// close it.
96+
function closeHandleAfterAck(message, handle, options, callback, target) {
97+
if (handle && !options.keepOpen) {
98+
if (target) {
99+
// There can only be one _pendingMessage as passing handles are
100+
// processed one at a time: handles are stored in _handleQueue while
101+
// waiting for the NODE_HANDLE_ACK of the current passing handle.
102+
assert(!target._pendingMessage);
103+
target._pendingMessage =
104+
{ callback, message, handle, options, retransmissions: 0 };
105+
} else {
106+
handle.close();
107+
}
108+
}
109+
}
110+
89111
// This object contain function to convert TCP objects to native handle objects
90112
// and back again.
91113
const handleConversion = {
@@ -116,6 +138,29 @@ const handleConversion = {
116138
},
117139
},
118140

141+
'net.BoundSocket': {
142+
simultaneousAccepts: true,
143+
144+
send(message, boundSocket, options) {
145+
// Pipe (path) binds cannot be sent over the IPC channel.
146+
if (boundSocket.isPipe)
147+
throw new ERR_INVALID_HANDLE_TYPE();
148+
149+
// Reuse the worker_threads transfer protocol: detaches the handle,
150+
// leaving the source in the adopted state, and throws if the bound
151+
// socket is not in a transferable state.
152+
return boundSocket[kTransfer]().data.handle;
153+
},
154+
155+
postSend: closeHandleAfterAck,
156+
157+
got(message, handle, emit) {
158+
const boundSocket = new net._TransferredBoundSocket();
159+
boundSocket[kDeserialize]({ handle });
160+
emit(boundSocket);
161+
},
162+
},
163+
119164
'net.Socket': {
120165
send(message, socket, options) {
121166
if (!socket._handle)
@@ -165,23 +210,7 @@ const handleConversion = {
165210
return handle;
166211
},
167212

168-
postSend(message, handle, options, callback, target) {
169-
// Store the handle after successfully sending it, so it can be closed
170-
// when the NODE_HANDLE_ACK is received. If the handle could not be sent,
171-
// just close it.
172-
if (handle && !options.keepOpen) {
173-
if (target) {
174-
// There can only be one _pendingMessage as passing handles are
175-
// processed one at a time: handles are stored in _handleQueue while
176-
// waiting for the NODE_HANDLE_ACK of the current passing handle.
177-
assert(!target._pendingMessage);
178-
target._pendingMessage =
179-
{ callback, message, handle, options, retransmissions: 0 };
180-
} else {
181-
handle.close();
182-
}
183-
}
184-
},
213+
postSend: closeHandleAfterAck,
185214

186215
got(message, handle, emit) {
187216
const socket = new net.Socket({
@@ -822,6 +851,8 @@ function setupChannel(target, channel, serializationMode) {
822851
message.type = 'net.Socket';
823852
} else if (handle instanceof net.Server) {
824853
message.type = 'net.Server';
854+
} else if (handle instanceof net.BoundSocket) {
855+
message.type = 'net.BoundSocket';
825856
} else if (handle instanceof TCP || handle instanceof Pipe) {
826857
message.type = 'net.Native';
827858
} else if (handle instanceof dgram.Socket) {

‎lib/net.js‎

Lines changed: 64 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -388,6 +388,10 @@ const kBoundPath = Symbol('kBoundPath');
388388

389389
const isLinux = process.platform === 'linux';
390390

391+
// Internal: construct an empty BoundSocket shell during postMessage()
392+
// deserialization; the transferred handle is installed by [kDeserialize].
393+
const kBoundSocketDeserialize = Symbol('kBoundSocketDeserialize');
394+
391395
// A role-neutral wrapper over a synchronously bound libuv handle: bound to a
392396
// local address (a numeric IP literal for TCP, or a filesystem/abstract path
393397
// for a unix-domain socket via { path }) but neither listening nor connecting
@@ -396,11 +400,20 @@ const isLinux = process.platform === 'linux';
396400
// handle must be closed by the caller. bind(2) is non-blocking, so binding
397401
// happens inline and errors throw synchronously. No DNS is performed.
398402
class BoundSocket {
399-
#handle;
403+
#handle = null;
400404
#address = {};
401405
#path;
402406

403407
constructor(options = kEmptyObject) {
408+
// An un-adopted BoundSocket can be moved to another thread by listing it
409+
// in the transferList of a worker_threads postMessage() call. See
410+
// [kTransfer]().
411+
markTransferMode(this, false, true);
412+
413+
if (options === kBoundSocketDeserialize) {
414+
return;
415+
}
416+
404417
validateObject(options, 'options');
405418

406419
if (options.path !== undefined) {
@@ -544,7 +557,56 @@ class BoundSocket {
544557
get isPipe() {
545558
return this.#path !== undefined;
546559
}
560+
561+
// A BoundSocket can be transferred only while it still owns its handle,
562+
// i.e. before it has been adopted, closed or already transferred. Only TCP
563+
// binds are transferable; pipe handles cannot move between event loops.
564+
#assertTransferable() {
565+
if (this.#handle === null || !(this.#handle instanceof TCP)) {
566+
throw new ERR_WORKER_HANDLE_NOT_TRANSFERABLE('net.BoundSocket');
567+
}
568+
}
569+
570+
[kTransferList]() {
571+
this.#assertTransferable();
572+
return [this.#handle];
573+
}
574+
575+
[kTransfer]() {
576+
this.#assertTransferable();
577+
const handle = this.#handle;
578+
// Detach the handle; the messaging layer takes ownership of it via
579+
// TCPWrap::TransferForMessaging(). Further use on the sending side throws
580+
// ERR_SOCKET_HANDLE_ADOPTED, as after adoption.
581+
this.#handle = null;
582+
return {
583+
data: { handle },
584+
deserializeInfo: 'net:_TransferredBoundSocket',
585+
};
586+
}
587+
588+
[kDeserialize](data) {
589+
const handle = data?.handle;
590+
if (handle == null || !(handle instanceof TCP)) {
591+
throw new ERR_WORKER_HANDLE_NOT_TRANSFERABLE('net.BoundSocket');
592+
}
593+
// Re-derive the bound address from the transferred handle rather than
594+
// trusting serialized state.
595+
const err = handle.getsockname(this.#address);
596+
if (err) {
597+
handle.close();
598+
throw new ERR_WORKER_HANDLE_NOT_TRANSFERABLE('net.BoundSocket');
599+
}
600+
this.#handle = handle;
601+
}
602+
}
603+
604+
// Deserialization target for a transferred BoundSocket: constructs an empty
605+
// shell without binding a new socket. Internal, not part of the public API.
606+
function _TransferredBoundSocket() {
607+
return new BoundSocket(kBoundSocketDeserialize);
547608
}
609+
_TransferredBoundSocket.prototype = BoundSocket.prototype;
548610

549611
function Socket(options) {
550612
if (!(this instanceof Socket)) return new Socket(options);
@@ -2933,6 +2995,7 @@ Server.prototype.unref = function() {
29332995
module.exports = {
29342996
_createServerHandle: createServerHandle,
29352997
_normalizeArgs: normalizeArgs,
2998+
_TransferredBoundSocket,
29362999
get BlockList() {
29373000
BlockList ??= require('internal/blocklist').BlockList;
29383001
return BlockList;

0 commit comments

Comments
 (0)