Skip to content

Commit a546fba

Browse files
committed
stream: fix pipeTo assert when writer is released during microtask
The deferred write introduced by commit 65aa8f6 ("stream: fix pipeTo to defer writes per WHATWG spec") wraps the write operation in queueMicrotask(). A race condition exists where the pipe can shutdown and release the writer between when [kChunk] schedules the microtask and when it executes, causing an ERR_INTERNAL_ASSERTION because writer[[stream]] is undefined. Fix by moving the shuttingDown flag into the shared pipe state object so that PipeToReadableStreamReadRequest can check it. The microtask now skips the write if the pipe has already begun shutting down. This aligns with the WHATWG Streams spec step 15 which states: "if shuttingDown becomes true, the user agent must not initiate further reads from reader, and must only perform writes of already-read chunks". Fixes: #63732 Signed-off-by: Qingyu Wang <wangqingyu.c0l1n@bytedance.com>
1 parent b345a17 commit a546fba

2 files changed

Lines changed: 47 additions & 9 deletions

File tree

‎lib/internal/webstreams/readablestream.js‎

Lines changed: 12 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -1424,8 +1424,6 @@ function readableStreamPipeTo(
14241424

14251425
source[kState].disturbed = true;
14261426

1427-
let shuttingDown = false;
1428-
14291427
if (signal !== undefined) {
14301428
try {
14311429
validateAbortSignal(signal, 'options.signal');
@@ -1438,6 +1436,7 @@ function readableStreamPipeTo(
14381436

14391437
const state = {
14401438
currentWrite: PromiseResolve(),
1439+
shuttingDown: false,
14411440
};
14421441

14431442
// The error here can be undefined. The rejected arg
@@ -1462,8 +1461,8 @@ function readableStreamPipeTo(
14621461
}
14631462

14641463
function shutdownWithAnAction(action, rejected, originalError) {
1465-
if (shuttingDown) return;
1466-
shuttingDown = true;
1464+
if (state.shuttingDown) return;
1465+
state.shuttingDown = true;
14671466
if (dest[kState].state === 'writable' &&
14681467
!writableStreamCloseQueuedOrInFlight(dest)) {
14691468
PromisePrototypeThen(
@@ -1483,8 +1482,8 @@ function readableStreamPipeTo(
14831482
}
14841483

14851484
function shutdown(rejected, error) {
1486-
if (shuttingDown) return;
1487-
shuttingDown = true;
1485+
if (state.shuttingDown) return;
1486+
state.shuttingDown = true;
14881487
if (dest[kState].state === 'writable' &&
14891488
!writableStreamCloseQueuedOrInFlight(dest)) {
14901489
PromisePrototypeThen(
@@ -1546,11 +1545,11 @@ function readableStreamPipeTo(
15461545
}
15471546

15481547
async function step() {
1549-
if (shuttingDown) return true;
1548+
if (state.shuttingDown) return true;
15501549

15511550
if (dest[kState].backpressure) {
15521551
await writer[kState].ready.promise;
1553-
if (shuttingDown) return true;
1552+
if (state.shuttingDown) return true;
15541553
}
15551554

15561555
const controller = source[kState].controller;
@@ -1563,7 +1562,7 @@ function readableStreamPipeTo(
15631562
controller[kState].queue.length > 0) {
15641563

15651564
while (controller[kState].queue.length > 0) {
1566-
if (shuttingDown) return true;
1565+
if (state.shuttingDown) return true;
15671566

15681567
const chunk = dequeueValue(controller);
15691568

@@ -1678,6 +1677,10 @@ class PipeToReadableStreamReadRequest {
16781677
// synchronous write during enqueue(). See WHATWG Streams spec
16791678
// "ReadableStreamPipeTo" step 15's "chunk steps".
16801679
queueMicrotask(() => {
1680+
if (this.state.shuttingDown) {
1681+
this.promise.resolve(false);
1682+
return;
1683+
}
16811684
this.state.currentWrite = writableStreamDefaultWriterWrite(this.writer, chunk);
16821685
markPromiseAsHandled(this.state.currentWrite);
16831686
this.promise.resolve(false);
Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
1+
'use strict';
2+
3+
const common = require('../common');
4+
const assert = require('assert');
5+
const { ReadableStream, WritableStream } = require('stream/web');
6+
7+
{
8+
let sourceController;
9+
let destController;
10+
11+
const source = new ReadableStream({
12+
start(controller) {
13+
sourceController = controller;
14+
},
15+
});
16+
17+
const dest = new WritableStream({
18+
start(controller) {
19+
destController = controller;
20+
},
21+
write() {},
22+
});
23+
24+
source.pipeTo(dest, { preventCancel: true }).then(
25+
common.mustNotCall('pipeTo should not resolve'),
26+
common.mustCall((err) => {
27+
assert.strictEqual(err.message, 'destination errored');
28+
})
29+
);
30+
31+
setImmediate(common.mustCall(() => {
32+
destController.error(new Error('destination errored'));
33+
sourceController.enqueue('chunk');
34+
}));
35+
}

0 commit comments

Comments
 (0)