Skip to content

Commit 1c7c0f5

Browse files
jasnelladuh95
authored andcommitted
stream: make broadcast waiters use SafeSet
Signed-off-by: James M Snell <jasnell@gmail.com> Assisted-by: Opencode PR-URL: #66030 Reviewed-By: Filip Skokan <panva.ip@gmail.com> Reviewed-By: Trivikram Kamat <trivikr.dev@gmail.com>
1 parent 9b50b62 commit 1c7c0f5

2 files changed

Lines changed: 45 additions & 8 deletions

File tree

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

Lines changed: 11 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -110,7 +110,7 @@ class BroadcastImpl {
110110
#buffer = new RingBuffer();
111111
#bufferStart = 0;
112112
#consumers = new SafeSet();
113-
#waiters = []; // Consumers with pending resolve (subset of #consumers)
113+
#waiters = new SafeSet(); // Consumers with pending resolve
114114
#ended = false;
115115
#error;
116116
#errored = false;
@@ -203,6 +203,7 @@ class BroadcastImpl {
203203

204204
function detach() {
205205
state.detached = true;
206+
self.#waiters.delete(state);
206207
if (state.resolve) {
207208
state.resolve({ __proto__: null, done: true, value: undefined });
208209
}
@@ -262,7 +263,7 @@ class BroadcastImpl {
262263
const { promise, resolve, reject } = PromiseWithResolvers();
263264
state.resolve = resolve;
264265
state.reject = reject;
265-
ArrayPrototypePush(self.#waiters, state);
266+
self.#waiters.add(state);
266267
return promise;
267268
},
268269

@@ -313,6 +314,7 @@ class BroadcastImpl {
313314
consumer.detached = true;
314315
}
315316
this.#consumers.clear();
317+
this.#waiters.clear();
316318
this.#cachedMinCursorConsumers = 0;
317319
const onCancel = this[kOnCancel];
318320
this[kOnCancel] = null;
@@ -401,6 +403,7 @@ class BroadcastImpl {
401403
}
402404
}
403405
}
406+
this.#waiters.clear();
404407
this.#notifyEndDrained();
405408
}
406409

@@ -422,6 +425,7 @@ class BroadcastImpl {
422425
consumer.detached = true;
423426
}
424427
this.#consumers.clear();
428+
this.#waiters.clear();
425429
this.#cachedMinCursorConsumers = 0;
426430
}
427431

@@ -489,12 +493,11 @@ class BroadcastImpl {
489493

490494
#notifyConsumers() {
491495
const waiters = this.#waiters;
492-
if (waiters.length === 0) return;
496+
if (waiters.size === 0) return;
493497
// Swap out the waiters list so consumers that re-wait during
494498
// resolve don't get processed twice in this cycle.
495-
this.#waiters = [];
496-
for (let i = 0; i < waiters.length; i++) {
497-
const consumer = waiters[i];
499+
this.#waiters = new SafeSet();
500+
for (const consumer of waiters) {
498501
if (consumer.resolve) {
499502
const bufferIndex = consumer.cursor - this.#bufferStart;
500503
if (bufferIndex < this.#buffer.length) {
@@ -513,11 +516,11 @@ class BroadcastImpl {
513516
if (consumer.detached && this.#deleteConsumer(consumer)) {
514517
this.#tryTrimBuffer();
515518
} else if (this.#promotePending(consumer)) {
516-
ArrayPrototypePush(this.#waiters, consumer);
519+
this.#waiters.add(consumer);
517520
}
518521
} else {
519522
// Still waiting -- put back
520-
ArrayPrototypePush(this.#waiters, consumer);
523+
this.#waiters.add(consumer);
521524
}
522525
}
523526
}
Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,34 @@
1+
// Flags: --experimental-stream-iter --expose-gc
2+
'use strict';
3+
4+
const common = require('../common');
5+
const assert = require('assert');
6+
const { broadcast } = require('stream/iter');
7+
8+
async function detachConsumers(shared, count) {
9+
for (let i = 0; i < count; i++) {
10+
const iterator = shared.push()[Symbol.asyncIterator]();
11+
const pending = iterator.next();
12+
await iterator.return();
13+
await pending;
14+
}
15+
}
16+
17+
async function testDetachedWaitersAreReleased() {
18+
const { broadcast: shared } = broadcast();
19+
20+
await detachConsumers(shared, 100);
21+
global.gc();
22+
const before = process.memoryUsage().heapUsed;
23+
24+
await detachConsumers(shared, 50_000);
25+
global.gc();
26+
const retained = process.memoryUsage().heapUsed - before;
27+
28+
assert.ok(retained < 8 * 1024 * 1024,
29+
`Detached Broadcast waiters retained ${retained} bytes`);
30+
assert.strictEqual(shared.consumerCount, 0);
31+
shared.cancel();
32+
}
33+
34+
testDetachedWaitersAreReleased().then(common.mustCall());

0 commit comments

Comments
 (0)