Skip to content

Commit a7f8de6

Browse files
guybedfordaduh95
authored andcommitted
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> PR-URL: #64725 Reviewed-By: Matteo Collina <matteo.collina@gmail.com> Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent f4665f1 commit a7f8de6

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
@@ -1581,7 +1581,7 @@ added: v0.5.9
15811581
15821582
* `message` {Object} A parsed JSON object or primitive value.
15831583
* `sendHandle` {Handle|undefined} `undefined` or a [`net.Socket`][],
1584-
[`net.Server`][], or [`dgram.Socket`][] object.
1584+
[`net.Server`][], [`net.BoundSocket`][], or [`dgram.Socket`][] object.
15851585
15861586
The `'message'` event is triggered when a child process uses
15871587
[`process.send()`][] to send messages.
@@ -1891,6 +1891,9 @@ subprocess.ref();
18911891
<!-- YAML
18921892
added: v0.5.9
18931893
changes:
1894+
- version: REPLACEME
1895+
pr-url: https://github.com/nodejs/node/pull/64725
1896+
description: '`net.BoundSocket` instances can now be sent.'
18941897
- version: v5.8.0
18951898
pr-url: https://github.com/nodejs/node/pull/5283
18961899
description: The `options` parameter, and the `keepOpen` option
@@ -1905,7 +1908,7 @@ changes:
19051908
19061909
* `message` {Object}
19071910
* `sendHandle` {Handle|undefined} `undefined`, or a [`net.Socket`][],
1908-
[`net.Server`][], or [`dgram.Socket`][] object.
1911+
[`net.Server`][], [`net.BoundSocket`][], or [`dgram.Socket`][] object.
19091912
* `options` {Object} The `options` argument, if present, is an object used to
19101913
parameterize the sending of certain types of handles. `options` supports
19111914
the following properties:
@@ -1972,12 +1975,18 @@ Applications should avoid using such messages or listening for
19721975
`'internalMessage'` events as it is subject to change without notice.
19731976
19741977
The optional `sendHandle` argument that may be passed to `subprocess.send()` is
1975-
for passing a TCP server or socket object to the child process. The child process will
1978+
for passing a TCP server, socket or [`net.BoundSocket`][] object to the child
1979+
process. The child process will
19761980
receive the object as the second argument passed to the callback function
19771981
registered on the [`'message'`][] event. Any data that is received
19781982
and buffered in the socket will not be sent to the child. Sending IPC sockets is
19791983
not supported on Windows.
19801984
1985+
Sending a `net.BoundSocket` moves its underlying TCP handle to the child
1986+
process, leaving the source instance in the adopted state as if it had been
1987+
adopted by a server or socket. The bound socket must not have been adopted or
1988+
closed, and pipe (`path`) binds cannot be sent.
1989+
19811990
The optional `callback` is a function that is invoked after the message is
19821991
sent but before the child process may have received it. The function is called with a
19831992
single argument: `null` on success, or an [`Error`][] object on failure.
@@ -2405,6 +2414,7 @@ or [`child_process.fork()`][].
24052414
[`child_process.spawnSync()`]: #child_processspawnsynccommand-args-options
24062415
[`dgram.Socket`]: dgram.md#class-dgramsocket
24072416
[`maxBuffer` and Unicode]: #maxbuffer-and-unicode
2417+
[`net.BoundSocket`]: net.md#class-netboundsocket
24082418
[`net.Server`]: net.md#class-netserver
24092419
[`net.Socket`]: net.md#class-netsocket
24102420
[`options.detached`]: #optionsdetached

‎doc/api/errors.md‎

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

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

30043005
<a id="ERR_SOURCE_MAP_CORRUPT"></a>
30053006

@@ -3567,9 +3568,10 @@ The `Response` that has been passed to `WebAssembly.compileStreaming` or to
35673568

35683569
### `ERR_WORKER_HANDLE_NOT_TRANSFERABLE`
35693570

3570-
An attempt was made to transfer a `net.Socket` or `net.Server` to another thread
3571-
via a `worker_threads` `postMessage()` call while it was not in a transferable
3572-
state, for example because it had already started reading or had buffered data.
3571+
An attempt was made to transfer a `net.Socket`, `net.Server` or
3572+
`net.BoundSocket` to another thread via a `worker_threads` `postMessage()` call
3573+
while it was not in a transferable state, for example because it had already
3574+
started reading, had buffered data, or had already been adopted.
35733575

35743576
<a id="ERR_WORKER_INIT_FAILED"></a>
35753577

‎doc/api/net.md‎

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

943+
An un-adopted TCP [`BoundSocket`][] can also be transferred, which moves the
944+
bound (but not yet listening or connected) socket. This allows a port to be
945+
reserved synchronously on one thread and adopted by a server or outgoing
946+
connection on another. Pipe binds are not transferable. After the transfer, the
947+
source `BoundSocket` behaves as if it had been adopted: `address()`, `fd()` and
948+
`close()` throw [`ERR_SOCKET_HANDLE_ADOPTED`][].
949+
943950
### `new net.Socket([options])`
944951

945952
<!-- YAML
@@ -1863,6 +1870,13 @@ file system entry; abstract and TCP binds have none to remove.
18631870
When a pipe `BoundSocket` bound to a source `path` is adopted as a client, that
18641871
path is reported as the socket's `localAddress` once it connects.
18651872

1873+
An un-adopted TCP `BoundSocket` can be moved to another thread by listing it in
1874+
the `transferList` of a [`worker_threads`][] `postMessage()` call, see
1875+
[Transferring TCP handles to other threads][]. It can likewise be sent to a
1876+
child process as the `sendHandle` argument of [`subprocess.send()`][]. In both
1877+
cases the source is left in the adopted state. Pipe binds cannot be moved
1878+
either way.
1879+
18661880
When an adopted `BoundSocket` connects to a numeric IP literal, `connect(2)` is
18671881
issued synchronously, so [`socket.localAddress`][] is resolved once
18681882
[`socket.connect()`][] returns. Connection failures are still reported via a
@@ -2493,6 +2507,7 @@ net.isIPv6('fhqwhgads'); // returns false
24932507
[`socket.setTimeout()`]: #socketsettimeouttimeout-callback
24942508
[`socket.setTimeout(timeout)`]: #socketsettimeouttimeout-callback
24952509
[`stream.getDefaultHighWaterMark()`]: stream.md#streamgetdefaulthighwatermarkobjectmode
2510+
[`subprocess.send()`]: child_process.md#subprocesssendmessage-sendhandle-options-callback
24962511
[`worker_threads`]: worker_threads.md
24972512
[`writable.destroy()`]: stream.md#writabledestroyerror
24982513
[`writable.destroyed`]: stream.md#writabledestroyed

‎doc/api/worker_threads.md‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1236,7 +1236,7 @@ port2.postMessage(circularData);
12361236
```
12371237
12381238
`transferList` may be a list of {ArrayBuffer}, [`MessagePort`][],
1239-
[`FileHandle`][], {net.Server}, and {net.Socket} objects.
1239+
[`FileHandle`][], {net.Server}, {net.Socket}, and {net.BoundSocket} objects.
12401240
After transferring, they are not usable on the sending side of the channel
12411241
anymore (even if they are not contained in `value`).
12421242
@@ -1247,6 +1247,8 @@ freshly accepted or created TCP connection that has not yet started reading and
12471247
has no buffered data, otherwise `postMessage()` throws
12481248
`ERR_WORKER_HANDLE_NOT_TRANSFERABLE`. This makes it possible to accept
12491249
connections on one thread and distribute them across a pool of worker threads.
1250+
Transferring a {net.BoundSocket} moves an un-adopted pre-bound socket, so a
1251+
port can be reserved synchronously on one thread and adopted on another.
12501252
Only TCP handles are supported.
12511253
12521254
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
@@ -55,6 +55,10 @@ const { TCP } = internalBinding('tcp_wrap');
5555
const { TTY } = internalBinding('tty_wrap');
5656
const { UDP } = internalBinding('udp_wrap');
5757
const SocketList = require('internal/socket_list');
58+
const {
59+
kDeserialize,
60+
kTransfer,
61+
} = require('internal/worker/js_transferable');
5862
const { owner_symbol } = require('internal/async_hooks').symbols;
5963
const { convertToValidSignal, deprecate } = require('internal/util');
6064
const { isArrayBufferView } = require('internal/util/types');
@@ -87,6 +91,24 @@ const kChannelHandle = Symbol('kChannelHandle');
8791
const kIsUsedAsStdio = Symbol('kIsUsedAsStdio');
8892
const kPendingMessages = Symbol('kPendingMessages');
8993

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

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

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

187216
got(message, handle, emit) {
188217
const socket = new net.Socket({
@@ -838,6 +867,8 @@ function setupChannel(target, channel, serializationMode) {
838867
message.type = 'net.Socket';
839868
} else if (handle instanceof net.Server) {
840869
message.type = 'net.Server';
870+
} else if (handle instanceof net.BoundSocket) {
871+
message.type = 'net.BoundSocket';
841872
} else if (handle instanceof TCP || handle instanceof Pipe) {
842873
message.type = 'net.Native';
843874
} 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);
@@ -2937,6 +2999,7 @@ Server.prototype.unref = function() {
29372999
module.exports = {
29383000
_createServerHandle: createServerHandle,
29393001
_normalizeArgs: normalizeArgs,
3002+
_TransferredBoundSocket,
29403003
get BlockList() {
29413004
BlockList ??= require('internal/blocklist').BlockList;
29423005
return BlockList;

0 commit comments

Comments
 (0)