From c31b9edea286fc38ee84669e2ca48ea0566d4bc0 Mon Sep 17 00:00:00 2001 From: Matteo Collina Date: Sat, 26 Sep 2026 12:35:49 +0200 Subject: [PATCH 1/2] worker: preload streams into startup snapshot Worker stdio loads stream after restoring the main-thread context, initializing its constructors and prototype methods on every startup. Preload stream during snapshot construction so workers can restore that initialized state instead. Worker-specific stdio instances and message ports are still created at runtime. Add a concurrent startup benchmark and verify snapshot membership, stdio methods, and isolation from parent-thread stream mutations. Assisted-by: pi Signed-off-by: Matteo Collina --- benchmark/worker/startup.js | 28 +++++++++++++ .../bootstrap/switches/is_main_thread.js | 3 ++ test/parallel/test-bootstrap-modules.js | 19 +++++++++ test/parallel/test-worker-snapshot-streams.js | 41 +++++++++++++++++++ 4 files changed, 91 insertions(+) create mode 100644 benchmark/worker/startup.js create mode 100644 test/parallel/test-worker-snapshot-streams.js 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/test/parallel/test-bootstrap-modules.js b/test/parallel/test-bootstrap-modules.js index e6a40d45e26f..9a99fa33d4a6 100644 --- a/test/parallel/test-bootstrap-modules.js +++ b/test/parallel/test-bootstrap-modules.js @@ -146,6 +146,25 @@ if (contextFromSnapshot) { 'NativeModule internal/mime', 'NativeModule internal/modules/typescript', 'NativeModule internal/net', + 'NativeModule internal/abort_controller', + '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/writable', + 'NativeModule stream', + 'NativeModule stream/promises', + 'NativeModule string_decoder', ].forEach(expected.beforePreExec.add.bind(expected.beforePreExec)); } else if (isMainThread) { expected.beforePreExec.delete(getFormatNativeModule); 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()); +} From 1a9691a476af05de819f6634a291414d56c95be9 Mon Sep 17 00:00:00 2001 From: Matteo Collina Date: Sat, 26 Sep 2026 14:47:24 +0200 Subject: [PATCH 2/2] stream: defer unused stream initialization Keep Readable and Writable eagerly initialized for worker stdio, but load operators, pipeline, compose, promises, and the other stream constructors on demand. This keeps their implementation state out of the startup snapshot when it is not needed. Preserve export identity, prototype descriptors, and promisification. Install subclass methods before linking prototype chains so lazy loading also works when a base constructor or prototype has been frozen. Update the bootstrap inventory and add lazy-loading compatibility tests. Assisted-by: pi Signed-off-by: Matteo Collina --- lib/internal/streams/duplex.js | 12 +- lib/internal/streams/lazy.js | 22 ++++ lib/internal/streams/operators.js | 39 +++++- lib/internal/streams/passthrough.js | 6 +- lib/internal/streams/pipeline.js | 14 ++- lib/internal/streams/transform.js | 6 +- lib/stream.js | 83 ++----------- lib/stream/promises.js | 3 +- test/fixtures/stream-lazy-frozen.js | 37 ++++++ test/fixtures/stream-lazy-load.js | 152 ++++++++++++++++++++++++ test/parallel/test-bootstrap-modules.js | 18 --- test/parallel/test-stream-lazy-load.js | 27 +++++ 12 files changed, 319 insertions(+), 100 deletions(-) create mode 100644 lib/internal/streams/lazy.js create mode 100644 test/fixtures/stream-lazy-frozen.js create mode 100644 test/fixtures/stream-lazy-load.js create mode 100644 test/parallel/test-stream-lazy-load.js 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 9a99fa33d4a6..1c5db787a8bc 100644 --- a/test/parallel/test-bootstrap-modules.js +++ b/test/parallel/test-bootstrap-modules.js @@ -146,24 +146,15 @@ if (contextFromSnapshot) { 'NativeModule internal/mime', 'NativeModule internal/modules/typescript', 'NativeModule internal/net', - 'NativeModule internal/abort_controller', '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/writable', 'NativeModule stream', - 'NativeModule stream/promises', 'NativeModule string_decoder', ].forEach(expected.beforePreExec.add.bind(expected.beforePreExec)); } else if (isMainThread) { @@ -173,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}`); +}