Skip to content

Commit 7ff6267

Browse files
authored
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> PR-URL: #66266 Reviewed-By: Mattias Buelens <mattias@buelens.com> Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent 147ade5 commit 7ff6267

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,
@@ -185,8 +186,7 @@ class BroadcastImpl {
185186
}
186187

187188
#createRawConsumer() {
188-
const state = {
189-
__proto__: null,
189+
const state = ObjectSetPrototypeOf({
190190
// Start at the oldest buffered entry so late-joining consumers
191191
// can read data already in the buffer.
192192
cursor: this.#bufferStart,
@@ -195,7 +195,7 @@ class BroadcastImpl {
195195
pending: [],
196196
detached: false,
197197
error: kNoBroadcastError,
198-
};
198+
}, null);
199199

200200
this.#consumers.add(state);
201201
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,
@@ -143,15 +144,14 @@ class ShareImpl {
143144
}
144145

145146
#createRawConsumer() {
146-
const state = {
147-
__proto__: null,
147+
const state = ObjectSetPrototypeOf({
148148
cursor: this.#bufferStart,
149149
resolve: null,
150150
reject: null,
151151
detached: false,
152152
error: kNoShareError,
153153
pendingNext: PromiseResolve(),
154-
};
154+
}, null);
155155

156156
this.#consumers.add(state);
157157
if (this.#consumers.size === 1) {
@@ -549,12 +549,11 @@ class SyncShareImpl {
549549
}
550550

551551
#createRawConsumer() {
552-
const state = {
553-
__proto__: null,
552+
const state = ObjectSetPrototypeOf({
554553
cursor: this.#bufferStart,
555554
detached: false,
556555
error: kNoShareError,
557-
};
556+
}, null);
558557

559558
this.#consumers.add(state);
560559
if (this.#consumers.size === 1) {

0 commit comments

Comments
 (0)