Skip to content

Commit 608112a

Browse files
authored
stream: reject push iterator.throw() with error
Keep the existing consumer cancellation side effects, but reject the iterator.throw() call with the supplied error instead of resolving with done: true. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: openai:gpt-5.5 PR-URL: #64380 Fixes: #64378 Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent f2ba101 commit 608112a

2 files changed

Lines changed: 35 additions & 10 deletions

File tree

lib/internal/streams/iter/push.js

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -668,7 +668,7 @@ function createReadable(queue) {
668668
},
669669
async throw(error) {
670670
queue.consumerThrow(error);
671-
return { __proto__: null, done: true, value: undefined };
671+
throw error;
672672
},
673673
};
674674
},

test/parallel/test-stream-iter-push-writer.js

Lines changed: 34 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -147,10 +147,16 @@ async function testOndrainRejectsOnConsumerThrow() {
147147
// Consumer throws via iterator.throw() before draining enough
148148
// to clear backpressure. The drain should reject.
149149
const iter = readable[Symbol.asyncIterator]();
150-
await iter.throw(new Error('consumer error'));
150+
const err = new Error('consumer error');
151+
const drainRejects = assert.rejects(drainPromise, (e) => e === err);
152+
const pendingWriteRejects = pendingWrite.catch(() => {});
153+
await assert.rejects(
154+
() => iter.throw(err),
155+
(e) => e === err,
156+
);
151157

152-
await assert.rejects(drainPromise, /consumer error/);
153-
await pendingWrite.catch(() => {}); // Ignore write rejection
158+
await drainRejects;
159+
await pendingWriteRejects; // Ignore write rejection
154160
}
155161

156162
async function testWritev() {
@@ -303,7 +309,11 @@ async function testConsumerThrowRejectsWrites() {
303309
writer.writeSync('a');
304310

305311
const iter = readable[Symbol.asyncIterator]();
306-
await iter.throw(new Error('consumer boom'));
312+
const err = new Error('consumer boom');
313+
await assert.rejects(
314+
() => iter.throw(err),
315+
(e) => e === err,
316+
);
307317

308318
// Subsequent async writes should reject with the consumer's error
309319
await assert.rejects(
@@ -312,6 +322,18 @@ async function testConsumerThrowRejectsWrites() {
312322
);
313323
}
314324

325+
async function testConsumerThrowRejectsWithThrownError() {
326+
const { readable } = push();
327+
328+
const iter = readable[Symbol.asyncIterator]();
329+
const err = new Error('boom');
330+
331+
await assert.rejects(
332+
() => iter.throw(err),
333+
(e) => e === err,
334+
);
335+
}
336+
315337
// end() resolves a pending read as done:true
316338
async function testEndResolvesPendingRead() {
317339
const { writer, readable } = push();
@@ -373,14 +395,16 @@ async function testConsumerThrowRejectsPendingRead() {
373395
await new Promise(setImmediate);
374396

375397
const err = new Error('consumer read boom');
376-
const throwResult = await iter.throw(err);
377-
assert.strictEqual(throwResult.value, undefined);
378-
assert.strictEqual(throwResult.done, true);
379-
380-
await assert.rejects(
398+
const readRejects = assert.rejects(
381399
() => readPromise,
382400
(e) => e === err,
383401
);
402+
await assert.rejects(
403+
() => iter.throw(err),
404+
(e) => e === err,
405+
);
406+
407+
await readRejects;
384408
}
385409

386410
// end() while writes are pending rejects those writes
@@ -527,6 +551,7 @@ Promise.all([
527551
testWriteUint8Array(),
528552
testOndrainWaitsForDrain(),
529553
testConsumerThrowRejectsWrites(),
554+
testConsumerThrowRejectsWithThrownError(),
530555
testEndResolvesPendingRead(),
531556
testFailRejectsPendingRead(),
532557
testFailRejectsFutureReadWithFalsyReason(),

0 commit comments

Comments
 (0)