Skip to content

Commit 350436a

Browse files
committed
stream: keep consumer state in fast mode
Create null-prototype share and broadcast consumer state with fast properties instead of V8 dictionary properties. Assisted-by: Pi Signed-off-by: Matteo Collina <hello@matteocollina.com>
1 parent e7d8ab5 commit 350436a

3 files changed

Lines changed: 45 additions & 9 deletions

File tree

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,37 @@
1+
'use strict';
2+
3+
const common = require('../common.js');
4+
5+
const bench = common.createBenchmark(main, {
6+
consumers: [2, 8, 32],
7+
batches: [1e4],
8+
n: [5],
9+
}, {
10+
flags: ['--experimental-stream-iter'],
11+
});
12+
13+
function main({ consumers, batches, n }) {
14+
const { shareSync } = require('stream/iter');
15+
const chunk = Buffer.alloc(1024);
16+
let bytes = 0;
17+
18+
function* source() {
19+
for (let i = 0; i < batches; i++) yield [chunk];
20+
}
21+
22+
bench.start();
23+
for (let run = 0; run < n; run++) {
24+
const shared = shareSync(source(), { budget: 65536 });
25+
const readers = Array.from({ length: consumers }, () =>
26+
shared.pull()[Symbol.iterator]());
27+
for (let i = 0; i < batches; i++) {
28+
for (let j = 0; j < consumers; j++) {
29+
bytes += readers[j].next().value[0].byteLength;
30+
}
31+
}
32+
}
33+
if (bytes !== batches * consumers * n * chunk.byteLength) {
34+
throw new Error('Incorrect byte count');
35+
}
36+
bench.end(batches * consumers * n);
37+
}

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

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@ const {
1111
ArrayPrototypePush,
1212
ArrayPrototypeShift,
1313
FunctionPrototypeCall,
14+
ObjectSetPrototypeOf,
1415
PromisePrototypeThen,
1516
PromiseReject,
1617
PromiseResolve,
@@ -183,8 +184,7 @@ class BroadcastImpl {
183184
}
184185

185186
#createRawConsumer() {
186-
const state = {
187-
__proto__: null,
187+
const state = ObjectSetPrototypeOf({
188188
// Start at the oldest buffered entry so late-joining consumers
189189
// can read data already in the buffer.
190190
cursor: this.#bufferStart,
@@ -193,7 +193,7 @@ class BroadcastImpl {
193193
pending: [],
194194
detached: false,
195195
error: kNoBroadcastError,
196-
};
196+
}, null);
197197

198198
this.#consumers.add(state);
199199
if (this.#consumers.size === 1) {

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

Lines changed: 5 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@
88
const {
99
ArrayPrototypePush,
1010
FunctionPrototypeCall,
11+
ObjectSetPrototypeOf,
1112
PromisePrototypeThen,
1213
PromiseResolve,
1314
PromiseWithResolvers,
@@ -141,15 +142,14 @@ class ShareImpl {
141142
}
142143

143144
#createRawConsumer() {
144-
const state = {
145-
__proto__: null,
145+
const state = ObjectSetPrototypeOf({
146146
cursor: this.#bufferStart,
147147
resolve: null,
148148
reject: null,
149149
detached: false,
150150
error: kNoShareError,
151151
pendingNext: PromiseResolve(),
152-
};
152+
}, null);
153153

154154
this.#consumers.add(state);
155155
if (this.#consumers.size === 1) {
@@ -546,12 +546,11 @@ class SyncShareImpl {
546546
}
547547

548548
#createRawConsumer() {
549-
const state = {
550-
__proto__: null,
549+
const state = ObjectSetPrototypeOf({
551550
cursor: this.#bufferStart,
552551
detached: false,
553552
error: kNoShareError,
554-
};
553+
}, null);
555554

556555
this.#consumers.add(state);
557556
if (this.#consumers.size === 1) {

0 commit comments

Comments
 (0)