diff --git a/benchmark/worker/startup.js b/benchmark/worker/startup.js new file mode 100644 index 000000000000..c114ab7e9efb --- /dev/null +++ b/benchmark/worker/startup.js @@ -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(); + } + } + }); + } +} diff --git a/lib/internal/bootstrap/switches/is_main_thread.js b/lib/internal/bootstrap/switches/is_main_thread.js index e9e37275b17f..3e18fe6e3228 100644 --- a/lib/internal/bootstrap/switches/is_main_thread.js +++ b/lib/internal/bootstrap/switches/is_main_thread.js @@ -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. diff --git a/lib/internal/streams/duplex.js b/lib/internal/streams/duplex.js index dce8bd6e0fdf..6d31f62d0998 100644 --- a/lib/internal/streams/duplex.js +++ b/lib/internal/streams/duplex.js @@ -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]; + } } } @@ -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); diff --git a/lib/internal/streams/lazy.js b/lib/internal/streams/lazy.js new file mode 100644 index 000000000000..8d80b5ef6a7b --- /dev/null +++ b/lib/internal/streams/lazy.js @@ -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'); + }, +}; diff --git a/lib/internal/streams/operators.js b/lib/internal/streams/operators.js index a5cbe1af428f..87b82b65b56b 100644 --- a/lib/internal/streams/operators.js +++ b/lib/internal/streams/operators.js @@ -6,10 +6,13 @@ const { MathFloor, Number, NumberIsNaN, + ObjectDefineProperty, + ObjectKeys, Promise, PromisePrototypeThen, PromiseReject, PromiseResolve, + ReflectApply, Symbol, } = primordials; @@ -18,6 +21,7 @@ const { AbortController, AbortSignal } = require('internal/abort_controller'); const { AbortError, codes: { + ERR_ILLEGAL_CONSTRUCTOR, ERR_MISSING_ARGS, ERR_OUT_OF_RANGE, }, @@ -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'); @@ -373,7 +378,7 @@ function take(number, options = undefined) { }.call(this); } -module.exports.streamReturningOperators = { +const streamReturningOperators = { drop, filter, flatMap, @@ -381,7 +386,7 @@ module.exports.streamReturningOperators = { take, }; -module.exports.promiseReturningOperators = { +const promiseReturningOperators = { every, forEach, reduce, @@ -389,3 +394,33 @@ module.exports.promiseReturningOperators = { 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; +} diff --git a/lib/internal/streams/passthrough.js b/lib/internal/streams/passthrough.js index d37f9caf0116..587361295111 100644 --- a/lib/internal/streams/passthrough.js +++ b/lib/internal/streams/passthrough.js @@ -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)) @@ -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); diff --git a/lib/internal/streams/pipeline.js b/lib/internal/streams/pipeline.js index 546ff579d8eb..6937d4354340 100644 --- a/lib/internal/streams/pipeline.js +++ b/lib/internal/streams/pipeline.js @@ -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 { @@ -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 }; diff --git a/lib/internal/streams/transform.js b/lib/internal/streams/transform.js index 2f4d498bc780..2d16f2568edc 100644 --- a/lib/internal/streams/transform.js +++ b/lib/internal/streams/transform.js @@ -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'); @@ -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); diff --git a/lib/stream.js b/lib/stream.js index 078fa092f483..1eb255dd4865 100644 --- a/lib/stream.js +++ b/lib/stream.js @@ -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'); @@ -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; @@ -120,15 +70,7 @@ 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'); }, }); @@ -136,6 +78,7 @@ ObjectDefineProperty(eos, customPromisify, { __proto__: null, enumerable: true, get() { + promises ??= require('stream/promises'); return promises.finished; }, }); diff --git a/lib/stream/promises.js b/lib/stream/promises.js index a8b65d62b096..9d36d4da4f15 100644 --- a/lib/stream/promises.js +++ b/lib/stream/promises.js @@ -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; diff --git a/test/fixtures/stream-lazy-frozen.js b/test/fixtures/stream-lazy-frozen.js new file mode 100644 index 000000000000..051529c204fc --- /dev/null +++ b/test/fixtures/stream-lazy-frozen.js @@ -0,0 +1,37 @@ +'use strict'; + +// Public assert initializes stdio, which can load Duplex before it is tested. +const assert = require('internal/assert'); +const stream = require('stream'); + +assert(!process.moduleLoadList.includes('NativeModule internal/streams/duplex')); +Object.freeze(stream); +Object.freeze(stream.Readable); +Object.freeze(stream.Readable.prototype); +Object.freeze(stream.Writable); +Object.freeze(stream.Writable.prototype); + +const { Duplex } = stream; +assert(Object.getPrototypeOf(Duplex) === stream.Readable); +assert(Object.getPrototypeOf(Duplex.prototype) === stream.Readable.prototype); +assert(Duplex.prototype.destroy === stream.Writable.prototype.destroy); +Object.freeze(Duplex); +Object.freeze(Duplex.prototype); + +const { Transform } = stream; +assert(Object.getPrototypeOf(Transform) === Duplex); +assert(Object.getPrototypeOf(Transform.prototype) === Duplex.prototype); +assert(Object.hasOwn(Transform.prototype, '_read')); +assert(Object.hasOwn(Transform.prototype, '_write')); +Object.freeze(Transform); +Object.freeze(Transform.prototype); + +const { PassThrough } = stream; +assert(Object.getPrototypeOf(PassThrough) === Transform); +assert(Object.getPrototypeOf(PassThrough.prototype) === Transform.prototype); +assert(Object.hasOwn(PassThrough.prototype, '_transform')); + +const descriptor = Object.getOwnPropertyDescriptor(stream.Readable.prototype, 'map'); +assert(descriptor.value.name === 'map'); +assert(descriptor.writable === false); +assert(descriptor.configurable === false); diff --git a/test/fixtures/stream-lazy-load.js b/test/fixtures/stream-lazy-load.js new file mode 100644 index 000000000000..c76110945684 --- /dev/null +++ b/test/fixtures/stream-lazy-load.js @@ -0,0 +1,152 @@ +'use strict'; + +// test/common and assert can initialize stdio, which accesses stream.Duplex. +// Capture the load lists before importing either of them. +const stream = require('stream'); +const initiallyLoaded = process.moduleLoadList.slice(); +Object.keys(stream); +Object.getOwnPropertyNames(stream.Readable.prototype); +const { Readable, Writable } = stream; +const afterEnumeration = process.moduleLoadList.slice(); + +const assert = require('assert'); +let completed = false; +process.on('exit', () => assert(completed)); +run(process.argv[2]).then(() => { + completed = true; +}); + +async function run(scenario) { + const optional = [ + 'internal/abort_controller', + 'internal/streams/compose', + 'internal/streams/duplex', + 'internal/streams/duplexpair', + 'internal/streams/operators', + 'internal/streams/passthrough', + 'internal/streams/pipeline', + 'internal/streams/transform', + 'stream/promises', + ]; + const loaded = (id) => process.moduleLoadList.includes(`NativeModule ${id}`); + for (const id of optional) { + assert.strictEqual(initiallyLoaded.includes(`NativeModule ${id}`), false, id); + assert.strictEqual(afterEnumeration.includes(`NativeModule ${id}`), false, id); + } + assert.strictEqual(Readable, require('internal/streams/readable')); + assert.strictEqual(Writable, require('internal/streams/writable')); + + switch (scenario) { + case 'constructors': { + for (const [name, id] of [ + ['Duplex', 'duplex'], + ['Transform', 'transform'], + ['PassThrough', 'passthrough'], + ['duplexPair', 'duplexpair'], + ]) { + const descriptor = Object.getOwnPropertyDescriptor(stream, name); + const value = require(`internal/streams/${id}`); + assert.deepStrictEqual(descriptor, { + value, writable: true, enumerable: true, configurable: true, + }); + assert.strictEqual(stream[name], value); + } + assert.strictEqual(loaded('internal/streams/operators'), false); + assert.strictEqual(loaded('internal/streams/pipeline'), false); + assert.strictEqual(loaded('stream/promises'), false); + break; + } + case 'assignment': { + const replacement = () => {}; + stream.Transform = replacement; + stream.Readable.prototype.map = replacement; + assert.strictEqual(stream.Transform, replacement); + assert.strictEqual(stream.Readable.prototype.map, replacement); + assert.strictEqual(loaded('internal/streams/transform'), false); + assert.strictEqual(loaded('internal/streams/operators'), false); + break; + } + case 'operators': { + const source = stream.Readable.from([1, 2, 3]); + const map = source.map; + assert.strictEqual(Object.hasOwn(source, 'map'), false); + assert.strictEqual(map, stream.Readable.prototype.map); + const expected = { + drop: 1, filter: 2, flatMap: 2, map: 2, take: 1, + every: 1, forEach: 2, reduce: 3, toArray: 1, some: 1, find: 2, + }; + // Keep the lazy property list in sync with the implementation exports. + assert.deepStrictEqual(Object.keys(require('internal/streams/operators')), Object.keys(expected)); + for (const [name, length] of Object.entries(expected)) { + const fn = stream.Readable.prototype[name]; + assert.strictEqual(fn.name, name); + assert.strictEqual(fn.length, length); + assert.deepStrictEqual(Object.getOwnPropertyDescriptor(stream.Readable.prototype, name), { + value: fn, enumerable: false, configurable: true, writable: true, + }); + assert.throws(() => new fn(), { code: 'ERR_ILLEGAL_CONSTRUCTOR' }); + } + assert.deepStrictEqual(await source.map((value) => value * 2).toArray(), [2, 4, 6]); + assert.strictEqual(loaded('internal/streams/pipeline'), false); + break; + } + case 'finished': { + const { promisify } = require('util'); + const finished = promisify(stream.finished); + assert.strictEqual(finished, stream.promises.finished); + assert.strictEqual(stream.promises, require('stream/promises')); + assert.strictEqual(loaded('internal/streams/pipeline'), false); + await finished(stream.Readable.from([]).resume()); + assert.strictEqual(loaded('internal/streams/pipeline'), false); + break; + } + case 'pipeline': { + const callbackPipeline = stream.pipeline; + assert.strictEqual(callbackPipeline, require('internal/streams/pipeline').pipeline); + assert.strictEqual(loaded('stream/promises'), false); + const { promisify } = require('util'); + assert.strictEqual(promisify(callbackPipeline), stream.promises.pipeline); + const values = []; + await stream.promises.pipeline(stream.Readable.from([1, 2, 3]), new stream.Writable({ + objectMode: true, + write(chunk, encoding, cb) { values.push(chunk); cb(); }, + })); + assert.deepStrictEqual(values, [1, 2, 3]); + break; + } + case 'compose': { + assert.strictEqual(stream.compose, require('internal/streams/compose')); + const values = await stream.compose([1, 2], async function* (source) { + for await (const item of source) yield item * 3; + }).toArray(); + assert.deepStrictEqual(values, [3, 6]); + break; + } + case 'esm': { + const esm = await import('node:stream'); + assert.strictEqual(esm.default, stream); + for (const name of [ + 'Readable', 'Writable', 'Duplex', 'Transform', 'PassThrough', + 'pipeline', 'compose', 'promises', + ]) { + assert.strictEqual(esm[name], stream[name], name); + } + assert.strictEqual(loaded('internal/streams/operators'), false); + function replacement() {} + stream.Transform = replacement; + require('module').syncBuiltinESMExports(); + assert.strictEqual(esm.Transform, replacement); + break; + } + case 'vm': { + const vm = require('vm'); + const transform = vm.runInNewContext('stream.Transform', { stream }); + assert.strictEqual(transform, stream.Transform); + const map = vm.runInNewContext('stream.Readable.prototype.map', { stream }); + assert.strictEqual(map, stream.Readable.prototype.map); + break; + } + default: + assert.fail(`Unknown scenario: ${scenario}`); + } +} diff --git a/test/parallel/test-bootstrap-modules.js b/test/parallel/test-bootstrap-modules.js index e6a40d45e26f..1c5db787a8bc 100644 --- a/test/parallel/test-bootstrap-modules.js +++ b/test/parallel/test-bootstrap-modules.js @@ -146,6 +146,16 @@ if (contextFromSnapshot) { 'NativeModule internal/mime', 'NativeModule internal/modules/typescript', 'NativeModule internal/net', + 'NativeModule internal/streams/add-abort-signal', + 'NativeModule internal/streams/destroy', + 'NativeModule internal/streams/end-of-stream', + 'NativeModule internal/streams/from', + 'NativeModule internal/streams/legacy', + 'NativeModule internal/streams/readable', + 'NativeModule internal/streams/state', + 'NativeModule internal/streams/writable', + 'NativeModule stream', + 'NativeModule string_decoder', ].forEach(expected.beforePreExec.add.bind(expected.beforePreExec)); } else if (isMainThread) { expected.beforePreExec.delete(getFormatNativeModule); @@ -154,29 +164,20 @@ if (contextFromSnapshot) { if (!isMainThread) { [ 'NativeModule diagnostics_channel', - 'NativeModule internal/abort_controller', 'NativeModule internal/process/worker_thread_only', 'NativeModule internal/streams/add-abort-signal', - 'NativeModule internal/streams/compose', 'NativeModule internal/streams/destroy', - 'NativeModule internal/streams/duplex', - 'NativeModule internal/streams/duplexpair', 'NativeModule internal/streams/end-of-stream', 'NativeModule internal/streams/from', 'NativeModule internal/streams/legacy', - 'NativeModule internal/streams/operators', - 'NativeModule internal/streams/passthrough', - 'NativeModule internal/streams/pipeline', 'NativeModule internal/streams/readable', 'NativeModule internal/streams/state', - 'NativeModule internal/streams/transform', 'NativeModule internal/streams/utils', 'NativeModule internal/streams/writable', 'NativeModule internal/worker', 'NativeModule internal/worker/io', 'NativeModule internal/worker/messaging', 'NativeModule stream', - 'NativeModule stream/promises', 'NativeModule string_decoder', 'NativeModule worker_threads', ].forEach(expected.atRunTime.add.bind(expected.atRunTime)); diff --git a/test/parallel/test-stream-lazy-load.js b/test/parallel/test-stream-lazy-load.js new file mode 100644 index 000000000000..015547b02554 --- /dev/null +++ b/test/parallel/test-stream-lazy-load.js @@ -0,0 +1,27 @@ +// Flags: --expose-internals +'use strict'; + +require('../common'); +const assert = require('assert'); +const { spawnSync } = require('child_process'); +const fixtures = require('../common/fixtures'); + +const scenarios = [ + 'constructors', 'assignment', 'operators', 'finished', 'pipeline', + 'compose', 'esm', 'vm', +]; + +// Each scenario starts with an untouched stream module, with and without +// the embedded snapshot that preloads Readable and Writable. +for (const flags of [[], ['--no-node-snapshot']]) { + for (const scenario of scenarios) { + const child = spawnSync(process.execPath, [ + ...process.execArgv, ...flags, fixtures.path('stream-lazy-load.js'), scenario, + ], { encoding: 'utf8' }); + assert.strictEqual(child.status, 0, `${scenario}: ${child.stdout}\n${child.stderr}`); + } + const child = spawnSync(process.execPath, [ + ...process.execArgv, ...flags, fixtures.path('stream-lazy-frozen.js'), + ], { encoding: 'utf8' }); + assert.strictEqual(child.status, 0, `frozen: ${child.stdout}\n${child.stderr}`); +} diff --git a/test/parallel/test-worker-snapshot-streams.js b/test/parallel/test-worker-snapshot-streams.js new file mode 100644 index 000000000000..bda73a129c8c --- /dev/null +++ b/test/parallel/test-worker-snapshot-streams.js @@ -0,0 +1,41 @@ +// Flags: --expose-internals +'use strict'; + +const common = require('../common'); +const assert = require('assert'); +const { Worker, isMainThread, workerData } = require('worker_threads'); +const { internalBinding } = require('internal/test/binding'); + +const contextFromSnapshot = + process.config.variables.node_use_node_snapshot && + !process.execArgv.includes('--no-node-snapshot') && + (isMainThread || !process.execArgv.includes('--no-worker-snapshot')); +const { compiledInSnapshot } = internalBinding('builtins').getCacheUsage(); +for (const id of ['stream', 'internal/streams/readable', 'internal/streams/writable']) { + assert.strictEqual(compiledInSnapshot.includes(id), Boolean(contextFromSnapshot), id); +} + +const { Readable, Writable, getDefaultHighWaterMark, setDefaultHighWaterMark } = require('stream'); + +if (isMainThread) { + const highWaterMark = getDefaultHighWaterMark(false); + // Runtime changes in the parent must not become the worker's initial state. + setDefaultHighWaterMark(false, highWaterMark + 1); + Readable.prototype.parentOnly = true; + const worker = new Worker(__filename, { workerData: { highWaterMark } }); + worker.on('error', common.mustNotCall()); + worker.on('exit', common.mustCall((code) => assert.strictEqual(code, 0))); +} else { + assert.strictEqual(Readable.prototype.parentOnly, undefined); + if (workerData?.highWaterMark !== undefined) { + assert.strictEqual(getDefaultHighWaterMark(false), workerData.highWaterMark); + } + assert(process.stdin instanceof Readable); + assert(process.stdout instanceof Writable); + assert(process.stderr instanceof Writable); + assert.strictEqual(typeof process.stdin.map, 'function'); + assert.strictEqual(typeof process.stdin.toArray, 'function'); + Readable.from([1, 2, 3]).map((value) => value * 2).toArray().then((values) => { + assert.deepStrictEqual(values, [2, 4, 6]); + }).then(common.mustCall()); +}