Skip to content

Commit 977c20e

Browse files
authored
stream: avoid leaking consumers on signal failure
Validate Broadcast.push() and Share.pull() signals before registering raw consumers. Return a rejecting iterable for pre-aborted signals without adding a cursor. This prevents failed subscriptions from leaving unreachable cursors that inflate consumerCount and can permanently impose backpressure. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: codex:gpt-5.6-sol PR-URL: #65299 Fixes: #65298 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Jason Zhang <xzha4350@gmail.com>
1 parent 3abf65f commit 977c20e

5 files changed

Lines changed: 81 additions & 6 deletions

File tree

lib/internal/streams/iter/broadcast.js

Lines changed: 18 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -116,17 +116,32 @@ class BroadcastImpl {
116116

117117
push(...args) {
118118
const { transforms, options } = parsePullArgs(args);
119+
const signal = options?.signal;
120+
validateAbortSignal(signal, 'options.signal');
121+
122+
// Avoid registering a consumer that the pre-aborted pipeline will never
123+
// read or detach.
124+
if (signal?.aborted) {
125+
return {
126+
__proto__: null,
127+
// eslint-disable-next-line require-yield
128+
async *[SymbolAsyncIterator]() {
129+
throw signal.reason;
130+
},
131+
};
132+
}
133+
119134
const rawConsumer = this.#createRawConsumer();
120135

121136
// When transforms are present, delegate to pull() which creates its
122137
// own internal AbortController that follows the external signal.
123138
// When no transforms, return rawConsumer directly (controller elided
124139
// per PULL-02 optimization -- no transforms means no signal recipient).
125-
if (transforms.length > 0 || options?.signal) {
140+
if (transforms.length > 0 || signal) {
126141
const pullArgs = [...transforms];
127-
if (options?.signal) {
142+
if (signal) {
128143
ArrayPrototypePush(pullArgs,
129-
{ __proto__: null, signal: options.signal });
144+
{ __proto__: null, signal });
130145
}
131146
return pullWithTransforms(rawConsumer, ...pullArgs);
132147
}

lib/internal/streams/iter/share.js

Lines changed: 21 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -93,11 +93,29 @@ class ShareImpl {
9393

9494
pull(...args) {
9595
const { transforms, options } = parsePullArgs(args);
96+
const signal = options?.signal;
97+
validateAbortSignal(signal, 'options.signal');
98+
99+
// Avoid registering a consumer that the pre-aborted pipeline will never
100+
// read or detach.
101+
if (signal?.aborted) {
102+
return {
103+
__proto__: null,
104+
// eslint-disable-next-line require-yield
105+
async *[SymbolAsyncIterator]() {
106+
throw signal.reason;
107+
},
108+
};
109+
}
110+
96111
const rawConsumer = this.#createRawConsumer();
97112

98-
if (transforms.length > 0 || options?.signal) {
99-
if (options) {
100-
return pullWithTransforms(rawConsumer, ...transforms, options);
113+
if (transforms.length > 0 || signal) {
114+
if (signal) {
115+
return pullWithTransforms(
116+
rawConsumer,
117+
...transforms,
118+
{ __proto__: null, signal });
101119
}
102120
return pullWithTransforms(rawConsumer, ...transforms);
103121
}

test/parallel/test-stream-iter-broadcast-basic.js

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -203,6 +203,17 @@ async function testPushAbortSignalRejectsPendingNext() {
203203
await rejected;
204204
}
205205

206+
async function testPushPreAbortedSignalDoesNotAddConsumer() {
207+
const reason = new Error('already aborted');
208+
const signal = AbortSignal.abort(reason);
209+
const { broadcast: bc } = broadcast();
210+
const iter = bc.push({ signal })[Symbol.asyncIterator]();
211+
212+
assert.strictEqual(bc.consumerCount, 0);
213+
await assert.rejects(iter.next(), (error) => error === reason);
214+
assert.strictEqual(bc.consumerCount, 0);
215+
}
216+
206217
// =============================================================================
207218
// Writer fail detaches consumers
208219
// =============================================================================
@@ -331,6 +342,7 @@ Promise.all([
331342
testCancelWithFalsyReason(),
332343
testPendingNextSettlesAfterReturn(),
333344
testPushAbortSignalRejectsPendingNext(),
345+
testPushPreAbortedSignalDoesNotAddConsumer(),
334346
testFailDetachesConsumers(),
335347
testWriterFailIdempotent(),
336348
testLateJoinerSeesBufferedData(),

test/parallel/test-stream-iter-share-async.js

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -226,6 +226,17 @@ async function testSharePullAbortSignalRejectsPendingNext() {
226226
shared.cancel();
227227
}
228228

229+
async function testSharePullPreAbortedSignalDoesNotAddConsumer() {
230+
const reason = new Error('already aborted');
231+
const signal = AbortSignal.abort(reason);
232+
const shared = share(from('data'));
233+
const iter = shared.pull({ signal })[Symbol.asyncIterator]();
234+
235+
assert.strictEqual(shared.consumerCount, 0);
236+
await assert.rejects(iter.next(), (error) => error === reason);
237+
assert.strictEqual(shared.consumerCount, 0);
238+
}
239+
229240
async function testShareAlreadyAborted() {
230241
const shared = share(from('data'), { signal: AbortSignal.abort() });
231242
const consumer = shared.pull();
@@ -372,6 +383,7 @@ Promise.all([
372383
testShareAbortSignal(),
373384
testShareAbortSignalWhileSourcePullPending(),
374385
testSharePullAbortSignalRejectsPendingNext(),
386+
testSharePullPreAbortedSignalDoesNotAddConsumer(),
375387
testShareAlreadyAborted(),
376388
testShareSourceError(),
377389
testShareLateJoiningConsumer(),

test/parallel/test-stream-iter-validation.js

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -156,6 +156,15 @@ assert.throws(() => broadcast({ budget: 16383 }), { code: 'ERR_OUT_OF_RANGE' });
156156
assert.throws(() => broadcast({ signal: {} }), { code: 'ERR_INVALID_ARG_TYPE' });
157157
assert.throws(() => broadcast({ backpressure: 'bad' }), { code: 'ERR_INVALID_ARG_VALUE' });
158158

159+
// Broadcast consumer options.signal must be AbortSignal and validation must
160+
// not leave a consumer registered.
161+
{
162+
const { broadcast: bc } = broadcast();
163+
assert.throws(() => bc.push({ signal: {} }),
164+
{ code: 'ERR_INVALID_ARG_TYPE' });
165+
assert.strictEqual(bc.consumerCount, 0);
166+
}
167+
159168
// BroadcastWriter options.signal must be AbortSignal
160169
{
161170
const { writer } = broadcast();
@@ -212,6 +221,15 @@ assert.throws(() => share(from('a'), { budget: Number.MAX_SAFE_INTEGER + 1 }),
212221
assert.throws(() => share(from('a'), { signal: {} }), { code: 'ERR_INVALID_ARG_TYPE' });
213222
assert.throws(() => share(from('a'), { backpressure: 'bad' }), { code: 'ERR_INVALID_ARG_VALUE' });
214223

224+
// Share consumer options.signal must be AbortSignal and validation must not
225+
// leave a consumer registered.
226+
{
227+
const shared = share(from('a'));
228+
assert.throws(() => shared.pull({ signal: {} }),
229+
{ code: 'ERR_INVALID_ARG_TYPE' });
230+
assert.strictEqual(shared.consumerCount, 0);
231+
}
232+
215233
// share() values < 16384 are rejected
216234
assert.throws(() => share(from('a'), { budget: 0 }), { code: 'ERR_OUT_OF_RANGE' });
217235
assert.throws(() => share(from('a'), { budget: -1 }), { code: 'ERR_OUT_OF_RANGE' });

0 commit comments

Comments
 (0)