diff --git a/lib/internal/streams/iter/broadcast.js b/lib/internal/streams/iter/broadcast.js index a1384c9e4d51..4c63774bbb32 100644 --- a/lib/internal/streams/iter/broadcast.js +++ b/lib/internal/streams/iter/broadcast.js @@ -116,17 +116,32 @@ class BroadcastImpl { push(...args) { const { transforms, options } = parsePullArgs(args); + const signal = options?.signal; + validateAbortSignal(signal, 'options.signal'); + + // Avoid registering a consumer that the pre-aborted pipeline will never + // read or detach. + if (signal?.aborted) { + return { + __proto__: null, + // eslint-disable-next-line require-yield + async *[SymbolAsyncIterator]() { + throw signal.reason; + }, + }; + } + const rawConsumer = this.#createRawConsumer(); // When transforms are present, delegate to pull() which creates its // own internal AbortController that follows the external signal. // When no transforms, return rawConsumer directly (controller elided // per PULL-02 optimization -- no transforms means no signal recipient). - if (transforms.length > 0 || options?.signal) { + if (transforms.length > 0 || signal) { const pullArgs = [...transforms]; - if (options?.signal) { + if (signal) { ArrayPrototypePush(pullArgs, - { __proto__: null, signal: options.signal }); + { __proto__: null, signal }); } return pullWithTransforms(rawConsumer, ...pullArgs); } diff --git a/lib/internal/streams/iter/share.js b/lib/internal/streams/iter/share.js index 97efe7f3065b..662e57a7df55 100644 --- a/lib/internal/streams/iter/share.js +++ b/lib/internal/streams/iter/share.js @@ -93,11 +93,29 @@ class ShareImpl { pull(...args) { const { transforms, options } = parsePullArgs(args); + const signal = options?.signal; + validateAbortSignal(signal, 'options.signal'); + + // Avoid registering a consumer that the pre-aborted pipeline will never + // read or detach. + if (signal?.aborted) { + return { + __proto__: null, + // eslint-disable-next-line require-yield + async *[SymbolAsyncIterator]() { + throw signal.reason; + }, + }; + } + const rawConsumer = this.#createRawConsumer(); - if (transforms.length > 0 || options?.signal) { - if (options) { - return pullWithTransforms(rawConsumer, ...transforms, options); + if (transforms.length > 0 || signal) { + if (signal) { + return pullWithTransforms( + rawConsumer, + ...transforms, + { __proto__: null, signal }); } return pullWithTransforms(rawConsumer, ...transforms); } diff --git a/test/parallel/test-stream-iter-broadcast-basic.js b/test/parallel/test-stream-iter-broadcast-basic.js index 125f386210c9..4438e82a7817 100644 --- a/test/parallel/test-stream-iter-broadcast-basic.js +++ b/test/parallel/test-stream-iter-broadcast-basic.js @@ -203,6 +203,17 @@ async function testPushAbortSignalRejectsPendingNext() { await rejected; } +async function testPushPreAbortedSignalDoesNotAddConsumer() { + const reason = new Error('already aborted'); + const signal = AbortSignal.abort(reason); + const { broadcast: bc } = broadcast(); + const iter = bc.push({ signal })[Symbol.asyncIterator](); + + assert.strictEqual(bc.consumerCount, 0); + await assert.rejects(iter.next(), (error) => error === reason); + assert.strictEqual(bc.consumerCount, 0); +} + // ============================================================================= // Writer fail detaches consumers // ============================================================================= @@ -331,6 +342,7 @@ Promise.all([ testCancelWithFalsyReason(), testPendingNextSettlesAfterReturn(), testPushAbortSignalRejectsPendingNext(), + testPushPreAbortedSignalDoesNotAddConsumer(), testFailDetachesConsumers(), testWriterFailIdempotent(), testLateJoinerSeesBufferedData(), diff --git a/test/parallel/test-stream-iter-share-async.js b/test/parallel/test-stream-iter-share-async.js index 314e7bfcf01c..c96a0cb0f3c3 100644 --- a/test/parallel/test-stream-iter-share-async.js +++ b/test/parallel/test-stream-iter-share-async.js @@ -226,6 +226,17 @@ async function testSharePullAbortSignalRejectsPendingNext() { shared.cancel(); } +async function testSharePullPreAbortedSignalDoesNotAddConsumer() { + const reason = new Error('already aborted'); + const signal = AbortSignal.abort(reason); + const shared = share(from('data')); + const iter = shared.pull({ signal })[Symbol.asyncIterator](); + + assert.strictEqual(shared.consumerCount, 0); + await assert.rejects(iter.next(), (error) => error === reason); + assert.strictEqual(shared.consumerCount, 0); +} + async function testShareAlreadyAborted() { const shared = share(from('data'), { signal: AbortSignal.abort() }); const consumer = shared.pull(); @@ -372,6 +383,7 @@ Promise.all([ testShareAbortSignal(), testShareAbortSignalWhileSourcePullPending(), testSharePullAbortSignalRejectsPendingNext(), + testSharePullPreAbortedSignalDoesNotAddConsumer(), testShareAlreadyAborted(), testShareSourceError(), testShareLateJoiningConsumer(), diff --git a/test/parallel/test-stream-iter-validation.js b/test/parallel/test-stream-iter-validation.js index 9f77fe330478..8dcfb46f173f 100644 --- a/test/parallel/test-stream-iter-validation.js +++ b/test/parallel/test-stream-iter-validation.js @@ -156,6 +156,15 @@ assert.throws(() => broadcast({ budget: 16383 }), { code: 'ERR_OUT_OF_RANGE' }); assert.throws(() => broadcast({ signal: {} }), { code: 'ERR_INVALID_ARG_TYPE' }); assert.throws(() => broadcast({ backpressure: 'bad' }), { code: 'ERR_INVALID_ARG_VALUE' }); +// Broadcast consumer options.signal must be AbortSignal and validation must +// not leave a consumer registered. +{ + const { broadcast: bc } = broadcast(); + assert.throws(() => bc.push({ signal: {} }), + { code: 'ERR_INVALID_ARG_TYPE' }); + assert.strictEqual(bc.consumerCount, 0); +} + // BroadcastWriter options.signal must be AbortSignal { const { writer } = broadcast(); @@ -212,6 +221,15 @@ assert.throws(() => share(from('a'), { budget: Number.MAX_SAFE_INTEGER + 1 }), assert.throws(() => share(from('a'), { signal: {} }), { code: 'ERR_INVALID_ARG_TYPE' }); assert.throws(() => share(from('a'), { backpressure: 'bad' }), { code: 'ERR_INVALID_ARG_VALUE' }); +// Share consumer options.signal must be AbortSignal and validation must not +// leave a consumer registered. +{ + const shared = share(from('a')); + assert.throws(() => shared.pull({ signal: {} }), + { code: 'ERR_INVALID_ARG_TYPE' }); + assert.strictEqual(shared.consumerCount, 0); +} + // share() values < 16384 are rejected assert.throws(() => share(from('a'), { budget: 0 }), { code: 'ERR_OUT_OF_RANGE' }); assert.throws(() => share(from('a'), { budget: -1 }), { code: 'ERR_OUT_OF_RANGE' });