Skip to content

Commit b1e4e88

Browse files
committed
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
1 parent 0032189 commit b1e4e88

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() {
@@ -281,7 +287,11 @@ async function testConsumerThrowRejectsWrites() {
281287
writer.writeSync('a');
282288

283289
const iter = readable[Symbol.asyncIterator]();
284-
await iter.throw(new Error('consumer boom'));
290+
const err = new Error('consumer boom');
291+
await assert.rejects(
292+
() => iter.throw(err),
293+
(e) => e === err,
294+
);
285295

286296
// Subsequent async writes should reject with the consumer's error
287297
await assert.rejects(
@@ -290,6 +300,18 @@ async function testConsumerThrowRejectsWrites() {
290300
);
291301
}
292302

303+
async function testConsumerThrowRejectsWithThrownError() {
304+
const { readable } = push();
305+
306+
const iter = readable[Symbol.asyncIterator]();
307+
const err = new Error('boom');
308+
309+
await assert.rejects(
310+
() => iter.throw(err),
311+
(e) => e === err,
312+
);
313+
}
314+
293315
// end() resolves a pending read as done:true
294316
async function testEndResolvesPendingRead() {
295317
const { writer, readable } = push();
@@ -351,14 +373,16 @@ async function testConsumerThrowRejectsPendingRead() {
351373
await new Promise(setImmediate);
352374

353375
const err = new Error('consumer read boom');
354-
const throwResult = await iter.throw(err);
355-
assert.strictEqual(throwResult.value, undefined);
356-
assert.strictEqual(throwResult.done, true);
357-
358-
await assert.rejects(
376+
const readRejects = assert.rejects(
359377
() => readPromise,
360378
(e) => e === err,
361379
);
380+
await assert.rejects(
381+
() => iter.throw(err),
382+
(e) => e === err,
383+
);
384+
385+
await readRejects;
362386
}
363387

364388
// end() while writes are pending rejects those writes
@@ -504,6 +528,7 @@ Promise.all([
504528
testWriteUint8Array(),
505529
testOndrainWaitsForDrain(),
506530
testConsumerThrowRejectsWrites(),
531+
testConsumerThrowRejectsWithThrownError(),
507532
testEndResolvesPendingRead(),
508533
testFailRejectsPendingRead(),
509534
testFailRejectsFutureReadWithFalsyReason(),

0 commit comments

Comments
 (0)