Skip to content

Commit de4d49d

Browse files
committed
stream: narrow pipeTo guard to writer release check
The shuttingDown guard added by the previous commit is too broad: it also skips an already-read chunk while the destination is still writable. The WHATWG Streams spec requires pipeTo to write such chunks during shutdown. Check writer[kState].stream instead. It becomes undefined only after finalize() releases the writer, which is precisely the assertion case. Add a regression test that aborts after enqueue() and verifies that the already-read chunk is written. Signed-off-by: Qingyu Wang <wangqingyu.c0l1n@bytedance.com> Assisted-by: Codex
1 parent 456484b commit de4d49d

2 files changed

Lines changed: 31 additions & 1 deletion

File tree

‎lib/internal/webstreams/readablestream.js‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1760,7 +1760,7 @@ function readableStreamPipeTo(
17601760
function forwardChunk() {
17611761
const chunk = pendingChunk;
17621762
pendingChunk = undefined;
1763-
if (shuttingDown) return;
1763+
if (writer[kState].stream === undefined) return;
17641764
writableStreamDefaultWriterWriteWithRequest(writer, chunk, writeTracker);
17651765
pump();
17661766
}

‎test/parallel/test-webstreams-pipeto-writer-released-race.js‎

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -31,3 +31,33 @@ const { ReadableStream, WritableStream } = require('stream/web');
3131
sourceController.enqueue('chunk');
3232
}));
3333
}
34+
35+
{
36+
const ac = new AbortController();
37+
let sourceController;
38+
const chunks = [];
39+
40+
const source = new ReadableStream({
41+
start(controller) {
42+
sourceController = controller;
43+
},
44+
}, { highWaterMark: 0 });
45+
46+
const dest = new WritableStream({
47+
write: common.mustCall((chunk) => {
48+
chunks.push(chunk);
49+
}),
50+
}, { highWaterMark: 1 });
51+
52+
assert.rejects(
53+
source.pipeTo(dest, { signal: ac.signal }),
54+
{ name: 'AbortError' },
55+
).then(common.mustCall(() => {
56+
assert.deepStrictEqual(chunks, ['chunk']);
57+
}));
58+
59+
setImmediate(common.mustCall(() => {
60+
sourceController.enqueue('chunk');
61+
ac.abort();
62+
}));
63+
}

0 commit comments

Comments
 (0)