Skip to content

Commit 2b31b0d

Browse files
jasnelladuh95
authored andcommitted
stream: ensure factory signals remain active through closing
Signed-off-by: James M Snell <jasnell@gmail.com> Assisted-by: Opencode PR-URL: #65658 Reviewed-By: Trivikram Kamat <trivikr.dev@gmail.com>
1 parent 6e97890 commit 2b31b0d

3 files changed

Lines changed: 21 additions & 2 deletions

File tree

‎doc/api/stream_iter.md‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -765,7 +765,9 @@ added: v24.20.0
765765
**Default:** `16384`.
766766
* `backpressure` {string} Backpressure policy: `'strict'`, `'unbounded'`,
767767
`'drop-oldest'`, or `'drop-newest'`. **Default:** `'strict'`.
768-
* `signal` {AbortSignal} Abort the stream.
768+
* `signal` {AbortSignal} Abort the stream. The signal remains active while
769+
buffered data drains after `writer.end()`; aborting during that time fails
770+
the writer and rejects the pending `end()` promise.
769771
* Returns: {Object}
770772
* `writer` {Writable} The writer side.
771773
* `readable` {AsyncIterable} whose chunks fulfill with {Uint8Array\[]}

‎lib/internal/streams/iter/push.js‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -331,14 +331,14 @@ class PushQueue {
331331
return this.#bytesWritten; // Idempotent
332332
}
333333

334-
this.#cleanup();
335334
this.#rejectPendingWrites(
336335
new ERR_INVALID_STATE.TypeError('Writer closed'));
337336
this.#resolvePendingDrains(false);
338337

339338
// If buffer is empty, close immediately
340339
if (this.#slots.length === 0) {
341340
this.#writerState = 'closed';
341+
this.#cleanup();
342342
this.#resolvePendingReads();
343343
return this.#bytesWritten;
344344
}
@@ -356,6 +356,7 @@ class PushQueue {
356356
endDrained() {
357357
if (this.#writerState !== 'closing') return;
358358
this.#writerState = 'closed';
359+
this.#cleanup();
359360
if (this.#pendingEnd) {
360361
this.#pendingEnd.resolve(this.#bytesWritten);
361362
this.#pendingEnd = null;

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

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -294,6 +294,21 @@ async function testEndSignalAbortWhileDraining() {
294294
assert.strictEqual(await completedEnd, 5);
295295
}
296296

297+
async function testFactorySignalAbortWhileDraining() {
298+
const controller = new AbortController();
299+
const reason = new Error('stream aborted while draining');
300+
const { writer, readable } = push({ signal: controller.signal });
301+
302+
writer.writeSync('hello');
303+
const end = writer.end();
304+
const endRejected = assert.rejects(end, (error) => error === reason);
305+
controller.abort(reason);
306+
307+
await endRejected;
308+
await assert.rejects(text(readable), (error) => error === reason);
309+
await assert.rejects(writer.end(), (error) => error === reason);
310+
}
311+
297312
async function testEndAfterEndSyncWaitsForDrain() {
298313
const { writer, readable } = push();
299314
writer.writeSync('hello');
@@ -603,6 +618,7 @@ Promise.all([
603618
testEndAsyncReturnValue(),
604619
testEndWithPreAbortedSignal(),
605620
testEndSignalAbortWhileDraining(),
621+
testFactorySignalAbortWhileDraining(),
606622
testEndAfterEndSyncWaitsForDrain(),
607623
testWriteUint8Array(),
608624
testOndrainWaitsForDrain(),

0 commit comments

Comments
 (0)