Skip to content

Commit 3cfa60c

Browse files
stream: fix unhandled rejection for errored pipeTo writes
markPromiseAsHandled does not cancel V8's unhandled-rejection notification when the write promise is already rejected. Fall back to setPromiseHandled when the destination is no longer writable. Fixes: #64561 Refs: #63572 Signed-off-by: dushyant <dushyanthada90@gmail.com>
1 parent fb5d01f commit 3cfa60c

2 files changed

Lines changed: 66 additions & 2 deletions

File tree

lib/internal/webstreams/readablestream.js

Lines changed: 17 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1495,6 +1495,19 @@ function readableStreamFromIterable(iterable) {
14951495
return stream;
14961496
}
14971497

1498+
// markPromiseAsHandled only silences unhandled-rejection if the promise is
1499+
// still pending. writableStreamDefaultWriterWrite can return an already-
1500+
// rejected promise when the destination has errored or closed; in that case
1501+
// attach a real handler via setPromiseHandled.
1502+
function markPipeToWritePromiseHandled(promise, dest) {
1503+
if (dest[kState].state === 'writable' &&
1504+
!writableStreamCloseQueuedOrInFlight(dest)) {
1505+
markPromiseAsHandled(promise);
1506+
} else {
1507+
setPromiseHandled(promise);
1508+
}
1509+
}
1510+
14981511
function readableStreamPipeTo(
14991512
source,
15001513
dest,
@@ -1668,7 +1681,7 @@ function readableStreamPipeTo(
16681681
// Write the chunk - we're already in a separate microtask from enqueue
16691682
// because we awaited the writer ready promise above.
16701683
state.currentWrite = writableStreamDefaultWriterWrite(writer, chunk);
1671-
markPromiseAsHandled(state.currentWrite);
1684+
markPipeToWritePromiseHandled(state.currentWrite, dest);
16721685

16731686
// Check backpressure after each write
16741687
if (dest[kState].backpressure) {
@@ -1772,7 +1785,9 @@ class PipeToReadableStreamReadRequest {
17721785
// "ReadableStreamPipeTo" step 15's "chunk steps".
17731786
queueMicrotask(() => {
17741787
this.state.currentWrite = writableStreamDefaultWriterWrite(this.writer, chunk);
1775-
markPromiseAsHandled(this.state.currentWrite);
1788+
markPipeToWritePromiseHandled(
1789+
this.state.currentWrite,
1790+
this.writer[kState].stream);
17761791
this.promise.resolve(false);
17771792
});
17781793
}
Lines changed: 49 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,49 @@
1+
'use strict';
2+
3+
const common = require('../common');
4+
const assert = require('assert');
5+
const {
6+
ReadableStream,
7+
TransformStream,
8+
WritableStream,
9+
} = require('stream/web');
10+
11+
// A late write after the destination has errored returns an already-rejected
12+
// promise. That must not surface as an unhandled rejection when the pipeTo()
13+
// rejection is already handled.
14+
// Refs: https://github.com/nodejs/node/issues/64561
15+
16+
process.on('unhandledRejection', common.mustNotCall());
17+
18+
(async () => {
19+
const source = new ReadableStream({
20+
start(controller) {
21+
controller.enqueue('a');
22+
controller.enqueue('b');
23+
controller.close();
24+
},
25+
});
26+
27+
const asyncPassThrough = new TransformStream({
28+
async transform(chunk, controller) {
29+
controller.enqueue(chunk);
30+
},
31+
});
32+
33+
const identity = new TransformStream();
34+
35+
const failingTransform = new TransformStream({
36+
transform() {
37+
throw new Error('boom');
38+
},
39+
});
40+
41+
await assert.rejects(
42+
source
43+
.pipeThrough(asyncPassThrough)
44+
.pipeThrough(identity)
45+
.pipeThrough(failingTransform)
46+
.pipeTo(new WritableStream()),
47+
{ message: 'boom' },
48+
);
49+
})().then(common.mustCall());

0 commit comments

Comments
 (0)