Skip to content

Commit f484e35

Browse files
committed
child_process: watch child_process stdin pipe peer close event
Watch the child end of the stdin pipe for peer-close (EOF) events, only supported on non-Windows platforms. When the child closes its end of the pipe, subprocess.stdin is destroyed so that its 'close' event is emitted, instead of only being reported as an EPIPE error on the next write. Fixes: #25131
1 parent 6f41e41 commit f484e35

6 files changed

Lines changed: 204 additions & 5 deletions

File tree

‎doc/api/child_process.md‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2191,6 +2191,11 @@ until this stream has been closed via `end()`.
21912191
If the child process was spawned with `stdio[0]` set to anything other than `'pipe'`,
21922192
then this will be `null`.
21932193
2194+
On non-Windows platforms, when `stdio[0]` is `'pipe'`, Node.js watches for the
2195+
child process closing its end of the stdin pipe and destroys `subprocess.stdin`
2196+
when that happens. This helps surface pipe peer-close semantics consistently for
2197+
the writable side of the stream.
2198+
21942199
`subprocess.stdin` is an alias for `subprocess.stdio[0]`. Both properties will
21952200
refer to the same value.
21962201

‎lib/internal/child_process.js‎

Lines changed: 10 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -332,8 +332,15 @@ function flushStdio(subprocess) {
332332
}
333333

334334

335-
function createSocket(pipe, readable) {
336-
return net.Socket({ handle: pipe, readable });
335+
function createSocket(pipe, readable, watchPeerClose) {
336+
const sock = net.Socket({ handle: pipe, readable });
337+
if (watchPeerClose &&
338+
process.platform !== 'win32' &&
339+
typeof pipe?.watchPeerClose === 'function') {
340+
pipe.watchPeerClose(true, () => sock.destroy());
341+
sock.once('close', () => pipe.watchPeerClose(false));
342+
}
343+
return sock;
337344
}
338345

339346

@@ -489,7 +496,7 @@ ChildProcess.prototype.spawn = function spawn(options) {
489496

490497
if (stream.handle) {
491498
stream.socket = createSocket(this.pid !== 0 ?
492-
stream.handle : null, i > 0);
499+
stream.handle : null, i > 0, i === 0);
493500

494501
if (i > 0 && this.pid !== 0) {
495502
this._closesNeeded++;

‎src/pipe_wrap.cc‎

Lines changed: 100 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@
2828
#include "handle_wrap.h"
2929
#include "node.h"
3030
#include "node_buffer.h"
31+
#include "node_errors.h"
3132
#include "node_external_reference.h"
3233
#include "stream_base-inl.h"
3334
#include "stream_wrap.h"
@@ -80,6 +81,7 @@ void PipeWrap::Initialize(Local<Object> target,
8081
SetProtoMethod(isolate, t, "listen", Listen);
8182
SetProtoMethod(isolate, t, "connect", Connect);
8283
SetProtoMethod(isolate, t, "open", Open);
84+
SetProtoMethod(isolate, t, "watchPeerClose", WatchPeerClose);
8385

8486
#ifdef _WIN32
8587
SetProtoMethod(isolate, t, "setPendingInstances", SetPendingInstances);
@@ -110,6 +112,7 @@ void PipeWrap::RegisterExternalReferences(ExternalReferenceRegistry* registry) {
110112
registry->Register(Listen);
111113
registry->Register(Connect);
112114
registry->Register(Open);
115+
registry->Register(WatchPeerClose);
113116
#ifdef _WIN32
114117
registry->Register(SetPendingInstances);
115118
#endif
@@ -219,6 +222,103 @@ void PipeWrap::Open(const FunctionCallbackInfo<Value>& args) {
219222
args.GetReturnValue().Set(err);
220223
}
221224

225+
void PipeWrap::WatchPeerClose(const FunctionCallbackInfo<Value>& args) {
226+
PipeWrap* wrap;
227+
ASSIGN_OR_RETURN_UNWRAP(&wrap, args.This());
228+
229+
CHECK_GT(args.Length(), 0);
230+
CHECK(args[0]->IsBoolean());
231+
const bool enable = args[0].As<v8::Boolean>()->Value();
232+
Environment* env = wrap->env();
233+
Isolate* isolate = env->isolate();
234+
v8::HandleScope handle_scope(isolate);
235+
v8::Context::Scope context_scope(env->context());
236+
Local<Object> obj = wrap->object();
237+
238+
// UnwatchPeerClose
239+
if (!enable) {
240+
if (obj->GetInternalField(kPeerCloseCallbackField)
241+
.As<Value>()
242+
->IsUndefined()) {
243+
return;
244+
}
245+
246+
obj->SetInternalField(kPeerCloseCallbackField, v8::Undefined(isolate));
247+
uv_read_stop(wrap->stream());
248+
return;
249+
}
250+
251+
if (!wrap->IsAlive()) {
252+
return;
253+
}
254+
if (!obj->GetInternalField(kPeerCloseCallbackField)
255+
.As<Value>()
256+
->IsUndefined()) {
257+
return;
258+
}
259+
260+
CHECK_GT(args.Length(), 1);
261+
CHECK(args[1]->IsFunction());
262+
263+
// Store the JS callback in an internal field.
264+
obj->SetInternalField(kPeerCloseCallbackField, args[1]);
265+
266+
// Start reading to detect EOF/ECONNRESET from the peer.
267+
// We use our custom allocator and reader, ignoring actual data.
268+
int err = uv_read_start(wrap->stream(), PeerCloseAlloc, PeerCloseRead);
269+
if (err != 0) {
270+
obj->SetInternalField(kPeerCloseCallbackField, v8::Undefined(isolate));
271+
}
272+
}
273+
274+
void PipeWrap::PeerCloseAlloc(uv_handle_t* handle,
275+
size_t suggested_size,
276+
uv_buf_t* buf) {
277+
// We only care about EOF, not the actual data.
278+
// Using a static 1-byte buffer avoids dynamic memory allocation overhead.
279+
static char scratch;
280+
*buf = uv_buf_init(&scratch, 1);
281+
}
282+
283+
void PipeWrap::PeerCloseRead(uv_stream_t* stream,
284+
ssize_t nread,
285+
const uv_buf_t* buf) {
286+
PipeWrap* wrap = static_cast<PipeWrap*>(stream->data);
287+
if (wrap == nullptr) return;
288+
289+
// Ignore actual data reads or EAGAIN (0). We only watch for disconnects.
290+
if (nread > 0 || nread == 0) return;
291+
292+
// Wait specifically for EOF or connection reset (peer closed).
293+
if (nread != UV_EOF && nread != UV_ECONNRESET) return;
294+
295+
// Peer has closed the connection. Stop reading immediately.
296+
uv_read_stop(stream);
297+
298+
Environment* env = wrap->env();
299+
Isolate* isolate = env->isolate();
300+
301+
// Set up V8 context and handles to safely execute the JS callback.
302+
v8::HandleScope handle_scope(isolate);
303+
v8::Context::Scope context_scope(env->context());
304+
Local<Object> obj = wrap->object();
305+
306+
// Check if callback is set
307+
if (obj->GetInternalField(kPeerCloseCallbackField)
308+
.As<Value>()
309+
->IsUndefined()) {
310+
return;
311+
}
312+
Local<Value> cb_value =
313+
obj->GetInternalField(kPeerCloseCallbackField).As<Value>();
314+
Local<Function> cb = cb_value.As<Function>();
315+
// Reset before calling to prevent re-entrancy issues
316+
obj->SetInternalField(kPeerCloseCallbackField, v8::Undefined(isolate));
317+
318+
// MakeCallback properly tracks AsyncHooks context and flushes microtasks.
319+
wrap->MakeCallback(cb, 0, nullptr);
320+
}
321+
222322
void PipeWrap::Connect(const FunctionCallbackInfo<Value>& args) {
223323
Environment* env = Environment::GetCurrent(args);
224324

@@ -258,7 +358,6 @@ void PipeWrap::Connect(const FunctionCallbackInfo<Value>& args) {
258358

259359
args.GetReturnValue().Set(err);
260360
}
261-
262361
} // namespace node
263362

264363
NODE_BINDING_CONTEXT_AWARE_INTERNAL(pipe_wrap, node::PipeWrap::Initialize)

‎src/pipe_wrap.h‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,11 @@ class PipeWrap : public ConnectionWrap<PipeWrap, uv_pipe_t> {
4040
IPC
4141
};
4242

43+
enum InternalFields {
44+
kPeerCloseCallbackField = LibuvStreamWrap::kInternalFieldCount,
45+
kInternalFieldCount
46+
};
47+
4348
static v8::MaybeLocal<v8::Object> Instantiate(Environment* env,
4449
AsyncWrap* parent,
4550
SocketType type);
@@ -64,6 +69,13 @@ class PipeWrap : public ConnectionWrap<PipeWrap, uv_pipe_t> {
6469
static void Listen(const v8::FunctionCallbackInfo<v8::Value>& args);
6570
static void Connect(const v8::FunctionCallbackInfo<v8::Value>& args);
6671
static void Open(const v8::FunctionCallbackInfo<v8::Value>& args);
72+
static void WatchPeerClose(const v8::FunctionCallbackInfo<v8::Value>& args);
73+
static void PeerCloseAlloc(uv_handle_t* handle,
74+
size_t suggested_size,
75+
uv_buf_t* buf);
76+
static void PeerCloseRead(uv_stream_t* stream,
77+
ssize_t nread,
78+
const uv_buf_t* buf);
6779

6880
#ifdef _WIN32
6981
static void SetPendingInstances(

‎test/async-hooks/test-pipewrap.js‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,7 @@ const processwrap = processes[0];
3535
const pipe1 = pipes[0];
3636
const pipe2 = pipes[1];
3737
const pipe3 = pipes[2];
38+
const pipe1ExpectedInvocations = process.platform === 'win32' ? 1 : 2;
3839

3940
assert.strictEqual(processwrap.type, 'PROCESSWRAP');
4041
assert.strictEqual(processwrap.triggerAsyncId, 1);
@@ -83,7 +84,11 @@ function onexit() {
8384
// Usually it is just one event, but it can be more.
8485
assert.ok(ioEvents >= 3, `at least 3 stdout io events, got ${ioEvents}`);
8586

86-
checkInvocations(pipe1, { init: 1, before: 1, after: 1 },
87+
checkInvocations(pipe1, {
88+
init: 1,
89+
before: pipe1ExpectedInvocations,
90+
after: pipe1ExpectedInvocations,
91+
},
8792
'pipe wrap when sleep.spawn was called');
8893
checkInvocations(pipe2, { init: 1, before: ioEvents, after: ioEvents },
8994
'pipe wrap when sleep.spawn was called');
Lines changed: 71 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,71 @@
1+
'use strict';
2+
3+
const common = require('../common');
4+
const assert = require('assert');
5+
const { spawn } = require('child_process');
6+
7+
if (common.isWindows) {
8+
common.skip('Not applicable on Windows');
9+
}
10+
11+
function spawnChild(script) {
12+
return spawn(process.execPath, ['-e', script], {
13+
stdio: ['pipe', 'ignore', 'ignore'],
14+
});
15+
}
16+
17+
function runTest({ script, onSpawn }, done) {
18+
const child = spawnChild(script);
19+
20+
const timeout = setTimeout(() => {
21+
assert.fail('stdin close event was not emitted');
22+
}, 2000);
23+
24+
let closed = false;
25+
let exited = false;
26+
function maybeDone() {
27+
if (!closed || !exited) return;
28+
clearTimeout(timeout);
29+
done();
30+
}
31+
32+
child.stdin.once('close', common.mustCall(() => {
33+
closed = true;
34+
child.kill();
35+
maybeDone();
36+
}));
37+
38+
child.once('exit', common.mustCall(() => {
39+
exited = true;
40+
maybeDone();
41+
}));
42+
43+
onSpawn?.(child);
44+
}
45+
46+
runTest({
47+
script: 'setTimeout(() => require("fs").closeSync(0), 50); setTimeout(() => {}, 2000)',
48+
}, common.mustCall(() => {
49+
runTest({
50+
script: 'setTimeout(() => require("fs").closeSync(0), 200); setTimeout(() => {}, 2000)',
51+
onSpawn: common.mustCall((child) => {
52+
const handle = child.stdin?._handle;
53+
assert.strictEqual(typeof handle?.watchPeerClose, 'function');
54+
handle.watchPeerClose(true, common.mustNotCall());
55+
}),
56+
}, common.mustCall(() => {
57+
runTest({
58+
script: 'setTimeout(() => {}, 2000)',
59+
onSpawn: common.mustCall((child) => {
60+
const handle = child.stdin?._handle;
61+
assert.strictEqual(typeof handle?.watchPeerClose, 'function');
62+
63+
child.stdin.once('close', common.mustCall(() => {
64+
// Calling watchPeerClose again after close must not throw.
65+
handle.watchPeerClose(true, common.mustNotCall());
66+
}));
67+
child.stdin.destroy();
68+
}),
69+
}, common.mustCall());
70+
}));
71+
}));

0 commit comments

Comments
 (0)