node/lib/internal/webstreams/util.js
Matteo Collina fa2eaa54fb
stream: reduce allocations on WHATWG streams hot paths
Pure-JavaScript optimizations to lib/internal/webstreams/* that reduce
per-chunk and per-construction allocations on hot paths without
observable behavior change.

Per-chunk: reuse promise-reaction closures per controller, add buffered
fast path for async iterator, specialize callback wrappers by arity,
and share immutable nil records for writable stream resets.

Per-construction: use queueMicrotask for non-object start results,
materialize reader/writer .closed and .ready records lazily, and remove
dead allocations.

Assisted-by: Claude Fable 5
Signed-off-by: Matteo Collina <hello@matteocollina.com>
PR-URL: https://github.com/nodejs/node/pull/63876
Reviewed-By: Yagiz Nizipli <yagiz@nizipli.com>
Reviewed-By: Mattias Buelens <mattias@buelens.com>
Reviewed-By: Gürgün Dayıoğlu <hey@gurgun.day>
Reviewed-By: James M Snell <jasnell@gmail.com>
Reviewed-By: Antoine du Hamel <duhamelantoine1995@gmail.com>
2026-06-20 14:19:08 +00:00

283 lines
7.2 KiB
JavaScript

'use strict';
const {
ArrayBufferPrototypeGetByteLength,
ArrayBufferPrototypeGetDetached,
ArrayBufferPrototypeSlice,
ArrayPrototypePush,
ArrayPrototypeShift,
AsyncIteratorPrototype,
FunctionPrototypeCall,
MathMax,
NumberIsNaN,
PromisePrototypeThen,
PromiseReject,
PromiseResolve,
ReflectGet,
Symbol,
Uint8Array,
} = primordials;
const {
codes: {
ERR_INVALID_ARG_VALUE,
},
} = require('internal/errors');
const {
copyArrayBuffer,
} = internalBinding('buffer');
const {
inspect,
} = require('util');
const {
constants: {
kPending,
},
getPromiseDetails,
} = internalBinding('util');
const assert = require('internal/assert');
const {
validateFunction,
} = require('internal/validators');
const kState = Symbol('kState');
const kType = Symbol('kType');
const AsyncIterator = {
__proto__: AsyncIteratorPrototype,
next: undefined,
return: undefined,
};
const getNonWritablePropertyDescriptor = (value) => {
return {
__proto__: null,
configurable: true,
value,
};
};
function extractHighWaterMark(value, defaultHWM) {
if (value === undefined) return defaultHWM;
const coercedValue = +value;
if (NumberIsNaN(coercedValue) ||
coercedValue < 0)
throw new ERR_INVALID_ARG_VALUE.RangeError('strategy.highWaterMark', value);
return coercedValue;
}
// The default size algorithm is never exposed to user code, so a single
// shared function avoids one closure allocation per stream.
const defaultSizeAlgorithm = () => 1;
function extractSizeAlgorithm(size) {
if (size === undefined) return defaultSizeAlgorithm;
validateFunction(size, 'strategy.size');
return size;
}
function customInspect(depth, options, name, data) {
if (depth < 0)
return this;
const opts = {
...options,
depth: options.depth == null ? null : options.depth - 1,
};
return `${name} ${inspect(data, opts)}`;
}
// These are defensive to work around the possibility that
// the buffer, byteLength, and byteOffset properties on
// ArrayBuffer and ArrayBufferView's may have been tampered with.
function ArrayBufferViewGetBuffer(view) {
return ReflectGet(view.constructor.prototype, 'buffer', view);
}
function ArrayBufferViewGetByteLength(view) {
return ReflectGet(view.constructor.prototype, 'byteLength', view);
}
function ArrayBufferViewGetByteOffset(view) {
return ReflectGet(view.constructor.prototype, 'byteOffset', view);
}
function cloneAsUint8Array(view) {
const buffer = ArrayBufferViewGetBuffer(view);
const byteOffset = ArrayBufferViewGetByteOffset(view);
const byteLength = ArrayBufferViewGetByteLength(view);
return new Uint8Array(
ArrayBufferPrototypeSlice(buffer, byteOffset, byteOffset + byteLength),
);
}
function canCopyArrayBuffer(toBuffer, toIndex, fromBuffer, fromIndex, count) {
return toBuffer !== fromBuffer &&
!ArrayBufferPrototypeGetDetached(toBuffer) &&
!ArrayBufferPrototypeGetDetached(fromBuffer) &&
toIndex + count <= ArrayBufferPrototypeGetByteLength(toBuffer) &&
fromIndex + count <= ArrayBufferPrototypeGetByteLength(fromBuffer);
}
function isBrandCheck(brand) {
return (value) => {
return value != null &&
value[kState] !== undefined &&
value[kType] === brand;
};
}
function dequeueValue(controller) {
assert(controller[kState].queue !== undefined);
assert(controller[kState].queueTotalSize !== undefined);
assert(controller[kState].queue.length);
const {
value,
size,
} = ArrayPrototypeShift(controller[kState].queue);
controller[kState].queueTotalSize =
MathMax(0, controller[kState].queueTotalSize - size);
return value;
}
function resetQueue(controller) {
assert(controller[kState].queue !== undefined);
assert(controller[kState].queueTotalSize !== undefined);
controller[kState].queue = [];
controller[kState].queueTotalSize = 0;
}
function peekQueueValue(controller) {
assert(controller[kState].queue !== undefined);
assert(controller[kState].queueTotalSize !== undefined);
assert(controller[kState].queue.length);
return controller[kState].queue[0].value;
}
function enqueueValueWithSize(controller, value, size) {
assert(controller[kState].queue !== undefined);
assert(controller[kState].queueTotalSize !== undefined);
const coercedSize = +size;
if (NumberIsNaN(coercedSize) ||
coercedSize < 0 ||
coercedSize === Infinity) {
throw new ERR_INVALID_ARG_VALUE.RangeError('size', size);
}
size = coercedSize;
ArrayPrototypePush(controller[kState].queue, { value, size });
controller[kState].queueTotalSize += size;
}
// Arity-specialized variants of the promise-callback wrapper. The generic
// rest-parameter + ReflectApply form allocated an arguments array on every
// invocation; these run on per-chunk hot paths (pull/write/transform), so
// each known call-site arity gets its own wrapper. The exact number of
// arguments passed through to the user callback is observable and must be
// preserved.
function createPromiseCallbackNoParams(name, fn, thisArg) {
validateFunction(fn, name);
return async () => FunctionPrototypeCall(fn, thisArg);
}
function createPromiseCallback1Param(name, fn, thisArg) {
validateFunction(fn, name);
return async (arg) => FunctionPrototypeCall(fn, thisArg, arg);
}
function createPromiseCallback2Params(name, fn, thisArg) {
validateFunction(fn, name);
return async (arg1, arg2) => FunctionPrototypeCall(fn, thisArg, arg1, arg2);
}
function isPromisePending(promise) {
if (promise === undefined) return false;
const details = getPromiseDetails(promise);
return details?.[0] === kPending;
}
// Shared shapes for lazily-materialized { promise, resolve, reject }
// records whose settlement is already known.
function resolvedRecord() {
return {
promise: PromiseResolve(),
resolve: undefined,
reject: undefined,
};
}
function rejectedHandledRecord(error) {
const record = {
promise: PromiseReject(error),
resolve: undefined,
reject: undefined,
};
setPromiseHandled(record.promise);
return record;
}
function setPromiseHandled(promise) {
// Alternatively, we could use the native API
// MarkAsHandled, but this avoids the extra boundary cross
// and is hopefully faster at the cost of an extra Promise
// allocation.
PromisePrototypeThen(promise, undefined, () => {});
}
async function nonOpFlush() {}
function nonOpStart() {}
async function nonOpPull() {}
async function nonOpCancel() {}
async function nonOpWrite() {}
let transfer;
function lazyTransfer() {
if (transfer === undefined)
transfer = require('internal/webstreams/transfer');
return transfer;
}
module.exports = {
ArrayBufferViewGetBuffer,
ArrayBufferViewGetByteLength,
ArrayBufferViewGetByteOffset,
AsyncIterator,
canCopyArrayBuffer,
cloneAsUint8Array,
copyArrayBuffer,
createPromiseCallbackNoParams,
createPromiseCallback1Param,
createPromiseCallback2Params,
customInspect,
defaultSizeAlgorithm,
dequeueValue,
enqueueValueWithSize,
extractHighWaterMark,
extractSizeAlgorithm,
getNonWritablePropertyDescriptor,
isBrandCheck,
isPromisePending,
kState,
kType,
lazyTransfer,
nonOpCancel,
nonOpFlush,
nonOpPull,
nonOpStart,
nonOpWrite,
peekQueueValue,
rejectedHandledRecord,
resetQueue,
resolvedRecord,
setPromiseHandled,
};