Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 28 additions & 0 deletions benchmark/worker/startup.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
'use strict';

const common = require('../common.js');
const { Worker } = require('worker_threads');

const bench = common.createBenchmark(main, {
n: [1, 4, 8, 16],
});

function main({ n }) {
const workers = [];
let online = 0;

bench.start();
for (let i = 0; i < n; i++) {
// Keep workers alive so their teardown does not compete with startup.
const worker = new Worker('setInterval(() => {}, 1e6)', { eval: true });
workers.push(worker);
worker.on('online', () => {
if (++online === n) {
bench.end(n);
for (const worker of workers) {
worker.terminate();
}
}
});
}
}
3 changes: 3 additions & 0 deletions lib/internal/bootstrap/switches/is_main_thread.js
Original file line number Diff line number Diff line change
Expand Up @@ -312,6 +312,9 @@ if (isBuildingSnapshot()) {
require('internal/modules/typescript');
require('internal/blob');
require('internal/dns/utils');
// Workers restore this context before loading internal/worker/io, which
// needs the initialized stream constructors for its stdio subclasses.
require('stream');
}

// Needed to refresh the time origin.
Expand Down
12 changes: 8 additions & 4 deletions lib/internal/streams/duplex.js
Original file line number Diff line number Diff line change
Expand Up @@ -46,15 +46,14 @@ const {
const destroyImpl = require('internal/streams/destroy');
const { kOnConstructed } = require('internal/streams/utils');

ObjectSetPrototypeOf(Duplex.prototype, Readable.prototype);
ObjectSetPrototypeOf(Duplex, Readable);

{
const keys = ObjectKeys(Writable.prototype);
// Allow the keys array to be GC'ed.
for (let i = 0; i < keys.length; i++) {
const method = keys[i];
Duplex.prototype[method] ||= Writable.prototype[method];
if (!Readable.prototype[method]) {
Duplex.prototype[method] = Writable.prototype[method];
}
}
}

Expand Down Expand Up @@ -201,3 +200,8 @@ Duplex.from = function(body) {
duplexify ??= require('internal/streams/duplexify');
return duplexify(body, 'body');
};

// Define own methods before inheriting from constructors or prototypes that
// user code may have frozen before this module was loaded.
ObjectSetPrototypeOf(Duplex.prototype, Readable.prototype);
ObjectSetPrototypeOf(Duplex, Readable);
22 changes: 22 additions & 0 deletions lib/internal/streams/lazy.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
'use strict';

// These modules export a single function rather than named properties.
// Expose them by name for defineLazyProperties() without loading them together.
module.exports = {
__proto__: null,
get Duplex() {
return require('internal/streams/duplex');
},
get Transform() {
return require('internal/streams/transform');
},
get PassThrough() {
return require('internal/streams/passthrough');
},
get duplexPair() {
return require('internal/streams/duplexpair');
},
get compose() {
return require('internal/streams/compose');
},
};
39 changes: 37 additions & 2 deletions lib/internal/streams/operators.js
Original file line number Diff line number Diff line change
Expand Up @@ -6,10 +6,13 @@ const {
MathFloor,
Number,
NumberIsNaN,
ObjectDefineProperty,
ObjectKeys,
Promise,
PromisePrototypeThen,
PromiseReject,
PromiseResolve,
ReflectApply,
Symbol,
} = primordials;

Expand All @@ -18,6 +21,7 @@ const { AbortController, AbortSignal } = require('internal/abort_controller');
const {
AbortError,
codes: {
ERR_ILLEGAL_CONSTRUCTOR,
ERR_MISSING_ARGS,
ERR_OUT_OF_RANGE,
},
Expand All @@ -31,6 +35,7 @@ const {
const { kWeakHandler, kResistStopPropagation } = require('internal/event_target');
const destroyImpl = require('internal/streams/destroy');
const { finished } = require('internal/streams/end-of-stream');
const Stream = require('stream');

const kEmpty = Symbol('kEmpty');
const kEof = Symbol('kEof');
Expand Down Expand Up @@ -373,19 +378,49 @@ function take(number, options = undefined) {
}.call(this);
}

module.exports.streamReturningOperators = {
const streamReturningOperators = {
drop,
filter,
flatMap,
map,
take,
};

module.exports.promiseReturningOperators = {
const promiseReturningOperators = {
every,
forEach,
reduce,
toArray,
some,
find,
};

const streamKeys = ObjectKeys(streamReturningOperators);
for (let i = 0; i < streamKeys.length; i++) {
const key = streamKeys[i];
const op = streamReturningOperators[key];
function fn(...args) {
if (new.target) {
throw new ERR_ILLEGAL_CONSTRUCTOR();
}
return Stream.Readable.from(ReflectApply(op, this, args));
}
ObjectDefineProperty(fn, 'name', { __proto__: null, value: op.name });
ObjectDefineProperty(fn, 'length', { __proto__: null, value: op.length });
module.exports[key] = fn;
}

const promiseKeys = ObjectKeys(promiseReturningOperators);
for (let i = 0; i < promiseKeys.length; i++) {
const key = promiseKeys[i];
const op = promiseReturningOperators[key];
function fn(...args) {
if (new.target) {
throw new ERR_ILLEGAL_CONSTRUCTOR();
}
return ReflectApply(op, this, args);
}
ObjectDefineProperty(fn, 'name', { __proto__: null, value: op.name });
ObjectDefineProperty(fn, 'length', { __proto__: null, value: op.length });
module.exports[key] = fn;
}
6 changes: 4 additions & 2 deletions lib/internal/streams/passthrough.js
Original file line number Diff line number Diff line change
Expand Up @@ -32,8 +32,6 @@ const {
module.exports = PassThrough;

const Transform = require('internal/streams/transform');
ObjectSetPrototypeOf(PassThrough.prototype, Transform.prototype);
ObjectSetPrototypeOf(PassThrough, Transform);

function PassThrough(options) {
if (!(this instanceof PassThrough))
Expand All @@ -45,3 +43,7 @@ function PassThrough(options) {
PassThrough.prototype._transform = function(chunk, encoding, cb) {
cb(null, chunk);
};

// Install the override before inheriting a potentially frozen prototype.
ObjectSetPrototypeOf(PassThrough.prototype, Transform.prototype);
ObjectSetPrototypeOf(PassThrough, Transform);
14 changes: 13 additions & 1 deletion lib/internal/streams/pipeline.js
Original file line number Diff line number Diff line change
Expand Up @@ -5,13 +5,17 @@

const {
ArrayIsArray,
ObjectDefineProperty,
Promise,
SymbolAsyncIterator,
SymbolDispose,
} = primordials;

const { eos } = require('internal/streams/end-of-stream');
const { once } = require('internal/util');
const {
once,
promisify: { custom: customPromisify },
} = require('internal/util');
const destroyImpl = require('internal/streams/destroy');
const Duplex = require('internal/streams/duplex');
const {
Expand Down Expand Up @@ -469,4 +473,12 @@ function pipe(src, dst, finish, finishOnlyHandleError, { end }) {
return eos(dst, { readable: false, writable: true }, finish);
}

ObjectDefineProperty(pipeline, customPromisify, {
__proto__: null,
enumerable: true,
get() {
return require('stream/promises').pipeline;
},
});

module.exports = { pipelineImpl, pipeline };
6 changes: 4 additions & 2 deletions lib/internal/streams/transform.js
Original file line number Diff line number Diff line change
Expand Up @@ -74,8 +74,6 @@ const {
} = require('internal/errors').codes;
const Duplex = require('internal/streams/duplex');
const { getHighWaterMark } = require('internal/streams/state');
ObjectSetPrototypeOf(Transform.prototype, Duplex.prototype);
ObjectSetPrototypeOf(Transform, Duplex);

const kCallback = Symbol('kCallback');

Expand Down Expand Up @@ -202,3 +200,7 @@ Transform.prototype._read = function() {
callback();
}
};

// Install overrides before inheriting potentially frozen prototype methods.
ObjectSetPrototypeOf(Transform.prototype, Duplex.prototype);
ObjectSetPrototypeOf(Transform, Duplex);
83 changes: 13 additions & 70 deletions lib/stream.js
Original file line number Diff line number Diff line change
Expand Up @@ -23,32 +23,19 @@

const {
ObjectDefineProperty,
ObjectKeys,
ReflectApply,
} = primordials;

const {
defineLazyProperties,
promisify: { custom: customPromisify },
} = require('internal/util');

const {
streamReturningOperators,
promiseReturningOperators,
} = require('internal/streams/operators');

const {
codes: {
ERR_ILLEGAL_CONSTRUCTOR,
},
} = require('internal/errors');
const compose = require('internal/streams/compose');
const { setDefaultHighWaterMark, getDefaultHighWaterMark } = require('internal/streams/state');
const { pipeline } = require('internal/streams/pipeline');
const { destroyer } = require('internal/streams/destroy');
const { eos } = require('internal/streams/end-of-stream');
const internalBuffer = require('internal/buffer');

const promises = require('stream/promises');
let promises;
const utils = require('internal/streams/utils');
const { isArrayBufferView, isUint8Array } = require('internal/util/types');

Expand All @@ -61,57 +48,20 @@ Stream.isReadable = utils.isReadable;
Stream.isWritable = utils.isWritable;

Stream.Readable = require('internal/streams/readable');
const streamKeys = ObjectKeys(streamReturningOperators);
for (let i = 0; i < streamKeys.length; i++) {
const key = streamKeys[i];
const op = streamReturningOperators[key];
function fn(...args) {
if (new.target) {
throw new ERR_ILLEGAL_CONSTRUCTOR();
}
return Stream.Readable.from(ReflectApply(op, this, args));
}
ObjectDefineProperty(fn, 'name', { __proto__: null, value: op.name });
ObjectDefineProperty(fn, 'length', { __proto__: null, value: op.length });
ObjectDefineProperty(Stream.Readable.prototype, key, {
__proto__: null,
value: fn,
enumerable: false,
configurable: true,
writable: true,
});
}
const promiseKeys = ObjectKeys(promiseReturningOperators);
for (let i = 0; i < promiseKeys.length; i++) {
const key = promiseKeys[i];
const op = promiseReturningOperators[key];
function fn(...args) {
if (new.target) {
throw new ERR_ILLEGAL_CONSTRUCTOR();
}
return ReflectApply(op, this, args);
}
ObjectDefineProperty(fn, 'name', { __proto__: null, value: op.name });
ObjectDefineProperty(fn, 'length', { __proto__: null, value: op.length });
ObjectDefineProperty(Stream.Readable.prototype, key, {
__proto__: null,
value: fn,
enumerable: false,
configurable: true,
writable: true,
});
}
defineLazyProperties(Stream.Readable.prototype, 'internal/streams/operators', [
'drop', 'filter', 'flatMap', 'map', 'take',
'every', 'forEach', 'reduce', 'toArray', 'some', 'find',
], false);
Stream.Writable = require('internal/streams/writable');
Stream.Duplex = require('internal/streams/duplex');
Stream.Transform = require('internal/streams/transform');
Stream.PassThrough = require('internal/streams/passthrough');
Stream.duplexPair = require('internal/streams/duplexpair');
Stream.pipeline = pipeline;
defineLazyProperties(Stream, 'internal/streams/lazy', [
'Duplex', 'Transform', 'PassThrough', 'duplexPair',
]);
defineLazyProperties(Stream, 'internal/streams/pipeline', ['pipeline']);
const { addAbortSignal } = require('internal/streams/add-abort-signal');
Stream.addAbortSignal = addAbortSignal;
Stream.finished = eos;
Stream.destroy = destroyer;
Stream.compose = compose;
defineLazyProperties(Stream, 'internal/streams/lazy', ['compose']);
Stream.setDefaultHighWaterMark = setDefaultHighWaterMark;
Stream.getDefaultHighWaterMark = getDefaultHighWaterMark;

Expand All @@ -120,22 +70,15 @@ ObjectDefineProperty(Stream, 'promises', {
configurable: true,
enumerable: true,
get() {
return promises;
},
});

ObjectDefineProperty(pipeline, customPromisify, {
__proto__: null,
enumerable: true,
get() {
return promises.pipeline;
return promises ??= require('stream/promises');
},
});

ObjectDefineProperty(eos, customPromisify, {
__proto__: null,
enumerable: true,
get() {
promises ??= require('stream/promises');
return promises.finished;
},
});
Expand Down
3 changes: 2 additions & 1 deletion lib/stream/promises.js
Original file line number Diff line number Diff line change
Expand Up @@ -11,12 +11,13 @@ const {
isWebStream,
} = require('internal/streams/utils');

const { pipelineImpl: pl } = require('internal/streams/pipeline');
let pl;
const { finished } = require('internal/streams/end-of-stream');

require('stream');

function pipeline(...streams) {
pl ??= require('internal/streams/pipeline').pipelineImpl;
return new Promise((resolve, reject) => {
let signal;
let end;
Expand Down
Loading
Loading