From 010d18b9baa9d814b872ffc7d1c8c4b3a5a425c2 Mon Sep 17 00:00:00 2001 From: seungwoo Date: Sat, 18 Jul 2026 23:49:24 +0900 Subject: [PATCH 01/15] stream: preserve mutable chunks in Web Stream adapters Wait for native write callbacks when they can safely represent chunk consumption. For Duplex streams, corked writes, and custom or legacy write methods, pass private BufferSource copies to preserve completion timing. Coordinate callback completion with backpressure, aborts, and stream errors so a settled Web Streams write no longer exposes mutable bytes still retained by the native stream. Preserve native HTTP validation and SharedArrayBuffer backing when fallback copies are required. Signed-off-by: seungwoo --- lib/_http_outgoing.js | 3 + lib/internal/streams/utils.js | 1 + lib/internal/streams/writable.js | 8 +- lib/internal/webstreams/adapters.js | 353 ++++++++++++++++++++-- lib/internal/webstreams/writablestream.js | 5 + 5 files changed, 342 insertions(+), 28 deletions(-) diff --git a/lib/_http_outgoing.js b/lib/_http_outgoing.js index 2ca455eb5566..817e4ddf49f2 100644 --- a/lib/_http_outgoing.js +++ b/lib/_http_outgoing.js @@ -1038,6 +1038,8 @@ OutgoingMessage.prototype.write = function write(chunk, encoding, callback) { return ret; }; +const outgoingMessagePrototypeWrite = OutgoingMessage.prototype.write; + function onError(msg, err, callback) { if (msg.destroyed) { return; @@ -1488,4 +1490,5 @@ module.exports = { validateHeaderName, validateHeaderValue, OutgoingMessage, + outgoingMessagePrototypeWrite, }; diff --git a/lib/internal/streams/utils.js b/lib/internal/streams/utils.js index 45f55316104f..06b05bab0184 100644 --- a/lib/internal/streams/utils.js +++ b/lib/internal/streams/utils.js @@ -345,6 +345,7 @@ module.exports = { isWritableEnded, isWritableFinished, isWritableErrored, + isOutgoingMessage, isServerRequest, isServerResponse, willEmitClose, diff --git a/lib/internal/streams/writable.js b/lib/internal/streams/writable.js index f8c61f819615..d39981ee05f6 100644 --- a/lib/internal/streams/writable.js +++ b/lib/internal/streams/writable.js @@ -512,6 +512,8 @@ Writable.prototype.write = function(chunk, encoding, cb) { return _write(this, chunk, encoding, cb) === true; }; +const writablePrototypeWrite = Writable.prototype.write; + Writable.prototype.cork = function() { const state = this._writableState; @@ -1187,7 +1189,11 @@ Writable.fromWeb = function(writableStream, options) { }; Writable.toWeb = function(streamWritable) { - return lazyWebStreams().newWritableStreamFromStreamWritable(streamWritable); + return lazyWebStreams().newWritableStreamFromStreamWritable( + streamWritable, + undefined, + writablePrototypeWrite, + ); }; Writable.prototype[SymbolAsyncDispose] = async function() { diff --git a/lib/internal/webstreams/adapters.js b/lib/internal/webstreams/adapters.js index a41b668f014c..0fd99fafdb98 100644 --- a/lib/internal/webstreams/adapters.js +++ b/lib/internal/webstreams/adapters.js @@ -1,25 +1,46 @@ 'use strict'; const { + ArrayBufferPrototypeGetDetached, ArrayPrototypeFilter, + BigInt64Array, + BigUint64Array, + DataView, + Float32Array, + Float64Array, + FunctionPrototypeCall, + Int16Array, + Int32Array, + Int8Array, ObjectKeys, PromisePrototypeThen, PromiseResolve, PromiseWithResolvers, SafePromiseAllReturnVoid, - SafePromisePrototypeFinally, SafeSet, StringPrototypeStartsWith, Symbol, + SymbolDispose, TypeError, TypedArrayPrototypeGetBuffer, TypedArrayPrototypeGetByteLength, TypedArrayPrototypeGetByteOffset, + TypedArrayPrototypeGetSymbolToStringTag, + TypedArrayPrototypeSet, + Uint16Array, + Uint32Array, Uint8Array, + Uint8ClampedArray, + globalThis: { + Float16Array, + }, } = primordials; const { TextEncoder } = require('internal/encoding'); +let addAbortListener; +let outgoingMessagePrototypeWrite; + const { ReadableStream, isReadableStream, @@ -31,6 +52,7 @@ const { const { WritableStream, + getWritableStreamDefaultControllerSignal, isWritableStream, writableStreamDefaultWriterWriteWithRequest, } = require('internal/webstreams/writablestream'); @@ -49,6 +71,7 @@ const { const { isDestroyed, + isOutgoingMessage, isReadable, isWritable, isWritableEnded, @@ -59,11 +82,18 @@ const { } = require('buffer'); const { + ArrayBufferViewGetBuffer, + ArrayBufferViewGetByteLength, + ArrayBufferViewGetByteOffset, kResolvedPromise, } = require('internal/webstreams/util'); const { isAnyArrayBuffer, + isArrayBufferView, + isDataView, + isSharedArrayBuffer, + isUint8Array, } = require('internal/util/types'); const { @@ -78,6 +108,7 @@ const { } = require('internal/errors'); const { + constructSharedArrayBuffer, getDeprecationWarningEmitter, kEmptyObject, normalizeEncoding, @@ -241,6 +272,62 @@ function dispatchWritev(writer, chunks, tracker) { tracker.settle(); } +function cloneArrayBufferView(view) { + const viewIsDataView = isDataView(view); + const buffer = ArrayBufferViewGetBuffer(view); + + // A detached backing buffer cannot be changed by the caller. Passing the + // original view through also preserves the native stream's validation. + if (!isSharedArrayBuffer(buffer) && + ArrayBufferPrototypeGetDetached(buffer)) { + return view; + } + + const byteLength = ArrayBufferViewGetByteLength(view); + const byteOffset = ArrayBufferViewGetByteOffset(view); + const copiedBytes = new Uint8Array( + isSharedArrayBuffer(buffer) ? + constructSharedArrayBuffer(byteLength) : + byteLength, + ); + TypedArrayPrototypeSet( + copiedBytes, + new Uint8Array(buffer, byteOffset, byteLength), + ); + const copiedBuffer = ArrayBufferViewGetBuffer(copiedBytes); + + if (viewIsDataView) { + return new DataView(copiedBuffer); + } + + switch (TypedArrayPrototypeGetSymbolToStringTag(view)) { + case 'BigInt64Array': + return new BigInt64Array(copiedBuffer); + case 'BigUint64Array': + return new BigUint64Array(copiedBuffer); + case 'Float16Array': + return new Float16Array(copiedBuffer); + case 'Float32Array': + return new Float32Array(copiedBuffer); + case 'Float64Array': + return new Float64Array(copiedBuffer); + case 'Int8Array': + return new Int8Array(copiedBuffer); + case 'Int16Array': + return new Int16Array(copiedBuffer); + case 'Int32Array': + return new Int32Array(copiedBuffer); + case 'Uint8Array': + return Buffer.isBuffer(view) ? Buffer.from(copiedBuffer) : copiedBytes; + case 'Uint8ClampedArray': + return new Uint8ClampedArray(copiedBuffer); + case 'Uint16Array': + return new Uint16Array(copiedBuffer); + case 'Uint32Array': + return new Uint32Array(copiedBuffer); + } +} + /** * @typedef {import('../../stream').Writable} Writable * @typedef {import('../../stream').Readable} Readable @@ -255,9 +342,14 @@ function dispatchWritev(writer, chunks, tracker) { /** * @param {Writable} streamWritable * @param {object} [options] + * @param {Function} [writablePrototypeWrite] * @returns {WritableStream} */ -function newWritableStreamFromStreamWritable(streamWritable, options = kEmptyObject) { +function newWritableStreamFromStreamWritable( + streamWritable, + options = kEmptyObject, + writablePrototypeWrite, +) { // Not using the internal/streams/utils isWritableNodeStream utility // here because it will return false if streamWritable is a Duplex // whose writable option is false. For a Duplex that is not writable, @@ -293,22 +385,188 @@ function newWritableStreamFromStreamWritable(streamWritable, options = kEmptyObj }; let controller; - let backpressurePromise; + let pendingWriteState; + let isTerminated = false; + let terminalReason; let closed; + let abortSignal; + let abortDisposable; + + function disposeAbortListener() { + abortDisposable?.[SymbolDispose](); + abortDisposable = undefined; + } + + function recordTerminalReason(reason) { + if (!isTerminated) { + isTerminated = true; + terminalReason = reason; + } + } + + function resolvePendingWrite(state) { + if (pendingWriteState !== state) { + return; + } + pendingWriteState = undefined; + disposeAbortListener(); + state.deferred.resolve(); + } + + function rejectPendingWrite(error) { + const state = pendingWriteState; + if (state === undefined) { + return; + } + pendingWriteState = undefined; + disposeAbortListener(); + state.deferred.reject(error); + } + + function listenForPendingAbort() { + addAbortListener ??= + require('internal/events/abort_listener').addAbortListener; + abortDisposable = addAbortListener(abortSignal, () => { + // Release an in-flight write so the sink's abort algorithm can run. + // The abort algorithm is solely responsible for destroying the stream. + recordTerminalReason(abortSignal.reason); + rejectPendingWrite(abortSignal.reason); + }); + } + + function maybeResolvePendingWrite(state) { + if (state.writeComplete && state.backpressureCleared) { + resolvePendingWrite(state); + } + } + + function markWriteComplete(state, error) { + if (pendingWriteState !== state) { + return; + } + if (error != null) { + rejectPendingWrite(error); + return; + } + state.writeComplete = true; + maybeResolvePendingWrite(state); + } + + function createPendingWriteState(waitForWriteCallback) { + const state = { + __proto__: null, + deferred: PromiseWithResolvers(), + writeComplete: !waitForWriteCallback, + backpressureCleared: false, + }; + pendingWriteState = state; + listenForPendingAbort(); + return state; + } function onDrain() { - backpressurePromise?.resolve(); + const state = pendingWriteState; + if (state === undefined) { + return; + } + state.backpressureCleared = true; + maybeResolvePendingWrite(state); + } + + function onWrite(error) { + const state = pendingWriteState; + if (state === undefined) { + return; + } + if (error != null) { + error = handleKnownInternalErrors(error); + } + markWriteComplete(state, error); + } + + function throwSyncWriteError(error) { + disposeAbortListener(); + // When the kDestroyOnSyncError flag is set (e.g. for + // CompressionStream), a sync throw must also destroy the + // stream so the readable side is errored too. Without this + // the readable side hangs forever. This replicates the + // TransformStream semantics: error both sides on any throw + // in the transform path. + if (options[kDestroyOnSyncError]) { + destroy(streamWritable, error); + } + throw error; + } + + function throwIfTerminated() { + if (abortSignal.aborted) { + recordTerminalReason(abortSignal.reason); + } + if (isTerminated) { + throw terminalReason; + } + } + + function writeChunk( + writeMethod, + chunk, + needDrainBeforeWrite, + waitForWriteCallback, + ) { + if (!waitForWriteCallback && !needDrainBeforeWrite) { + try { + const writeReturn = FunctionPrototypeCall( + writeMethod, + streamWritable, + chunk, + ); + throwIfTerminated(); + if (controller === undefined || + writeReturn || + !streamWritable.writableNeedDrain) { + return; + } + } catch (error) { + throwSyncWriteError(error); + } + + const state = createPendingWriteState(false); + return state.deferred.promise; + } + + const state = createPendingWriteState(waitForWriteCallback); + try { + // The return value only reports backpressure. The callback reports when + // the chunk has been handled and can safely be mutated by the caller. + const writeReturn = waitForWriteCallback ? + FunctionPrototypeCall(writeMethod, streamWritable, chunk, onWrite) : + FunctionPrototypeCall(writeMethod, streamWritable, chunk); + throwIfTerminated(); + const shouldWaitForDrain = + (needDrainBeforeWrite || !writeReturn) && + streamWritable.writableNeedDrain; + state.backpressureCleared = !shouldWaitForDrain; + maybeResolvePendingWrite(state); + } catch (error) { + resolvePendingWrite(state); + PromisePrototypeThen(state.deferred.promise, undefined, noop); + throwSyncWriteError(error); + } + + return state.deferred.promise; } const cleanup = eos(streamWritable, (error) => { error = handleKnownInternalErrors(error); cleanup(); + disposeAbortListener(); // This is a protection against non-standard, legacy streams // that happen to emit an error event again after finished is called. streamWritable.on('error', () => {}); if (error != null) { - backpressurePromise?.reject(error); + recordTerminalReason(error); + rejectPendingWrite(error); // If closed is not undefined, the error is happening // after the WritableStream close has already started. // We need to reject it here. @@ -326,57 +584,98 @@ function newWritableStreamFromStreamWritable(streamWritable, options = kEmptyObj closed = undefined; return; } - controller?.error(new AbortError()); + const abortError = new AbortError(); + recordTerminalReason(abortError); + rejectPendingWrite(abortError); + controller?.error(abortError); controller = undefined; }); streamWritable.on('drain', onDrain); return new WritableStream({ - start(c) { controller = c; }, + start(c) { + controller = c; + abortSignal = getWritableStreamDefaultControllerSignal(c); + }, write(chunk) { + let writeMethod; + let needDrainBeforeWrite; + let waitForWriteCallback; try { options[kValidateChunk]?.(chunk); - if (!streamWritable.writableObjectMode && isAnyArrayBuffer(chunk)) { + writeMethod = streamWritable.write; + const writableState = streamWritable._writableState; + const objectMode = streamWritable.writableObjectMode; + const isBufferSourceChunk = + isAnyArrayBuffer(chunk) || isArrayBufferView(chunk); + if (!objectMode && isAnyArrayBuffer(chunk)) { chunk = new Uint8Array(chunk); } - if (streamWritable.writableNeedDrain || !streamWritable.write(chunk)) { - backpressurePromise = PromiseWithResolvers(); - if (!streamWritable.writableNeedDrain) { - backpressurePromise.resolve(); + needDrainBeforeWrite = streamWritable.writableNeedDrain; + const isArrayBufferViewChunk = isArrayBufferView(chunk); + if (!objectMode && + isOutgoingMessage(streamWritable) && + isArrayBufferViewChunk && + !isUint8Array(chunk)) { + outgoingMessagePrototypeWrite ??= + require('_http_outgoing').outgoingMessagePrototypeWrite; + if (writeMethod === outgoingMessagePrototypeWrite) { + throw new ERR_INVALID_ARG_TYPE( + 'chunk', ['string', 'Buffer', 'Uint8Array'], chunk); } - return SafePromisePrototypeFinally( - backpressurePromise.promise, () => { - backpressurePromise = undefined; - }); } - } catch (error) { - // When the kDestroyOnSyncError flag is set (e.g. for - // CompressionStream), a sync throw must also destroy the - // stream so the readable side is errored too. Without this - // the readable side hangs forever. This replicates the - // TransformStream semantics: error both sides on any throw - // in the transform path. - if (options[kDestroyOnSyncError]) { - destroy(streamWritable, error); + const canAwaitWriteCallback = + streamWritable instanceof Writable && + writeMethod === writablePrototypeWrite && + streamWritable._readableState === undefined && + writableState?.corked === 0; + waitForWriteCallback = + isBufferSourceChunk && canAwaitWriteCallback; + if (!objectMode && + isArrayBufferViewChunk && + !waitForWriteCallback) { + // Duplex streams can retain chunks in their readable side or + // forward them after their own write callback. A corked stream does + // not invoke callbacks until uncork() or end(), and streams that + // override write() may not support callbacks at all. Preserve the + // existing completion timing for those cases by passing a private, + // same-type copy to the native stream. + chunk = cloneArrayBufferView(chunk); } - throw error; + } catch (error) { + throwSyncWriteError(error); } + + return writeChunk( + writeMethod, + chunk, + needDrainBeforeWrite, + waitForWriteCallback, + ); }, abort(reason) { + disposeAbortListener(); destroy(streamWritable, reason); }, close() { if (closed === undefined && !isWritableEnded(streamWritable)) { closed = PromiseWithResolvers(); - streamWritable.end(); + try { + streamWritable.end(); + } catch (error) { + closed = undefined; + disposeAbortListener(); + throw error; + } return closed.promise; } controller = undefined; + disposeAbortListener(); return PromiseResolve(); }, }, strategy); diff --git a/lib/internal/webstreams/writablestream.js b/lib/internal/webstreams/writablestream.js index 31f2a1f875dc..0ebce8dd60ae 100644 --- a/lib/internal/webstreams/writablestream.js +++ b/lib/internal/webstreams/writablestream.js @@ -609,6 +609,10 @@ class WritableStreamState { } ObjectSetPrototypeOf(WritableStreamState.prototype, null); +function getWritableStreamDefaultControllerSignal(controller) { + return controller[kState].abortController.signal; +} + function createWritableStreamState() { return new WritableStreamState(); } @@ -1420,6 +1424,7 @@ module.exports = { isWritableStreamDefaultController, isWritableStreamDefaultWriter, + getWritableStreamDefaultControllerSignal, isWritableStreamLocked, setupWritableStreamDefaultWriter, writableStreamAbort, From 0cc6d9908369a3bb1454558ac8133e8dd424a531 Mon Sep 17 00:00:00 2001 From: seungwoo Date: Sun, 19 Jul 2026 00:25:42 +0900 Subject: [PATCH 02/15] test: cover mutable Web Stream adapter writes Cover native callback completion, fallback copies, aborts, and error propagation for mutable BufferSource chunks passed to Node.js Web Streams adapters. Verify HTTP validation, SharedArrayBuffer backing, and Duplex and compression paths. Signed-off-by: seungwoo --- test/parallel/test-stream-duplex.js | 4 +- ...st-webstreams-adapters-sync-write-error.js | 75 +++- ...treams-adapters-writable-buffer-sources.js | 360 +++++++++++++++++- ...st-webstreams-compression-buffer-source.js | 29 ++ 4 files changed, 463 insertions(+), 5 deletions(-) diff --git a/test/parallel/test-stream-duplex.js b/test/parallel/test-stream-duplex.js index de67290a7e76..9908d000cb7e 100644 --- a/test/parallel/test-stream-duplex.js +++ b/test/parallel/test-stream-duplex.js @@ -120,7 +120,7 @@ process.on('exit', () => { this.push(null); }, write: common.mustCall((chunk) => { - assert.strictEqual(chunk, dataToWrite); + assert.deepStrictEqual(chunk, dataToWrite); }) }); @@ -143,7 +143,7 @@ process.on('exit', () => { this.push(null); }, write: common.mustCall((chunk) => { - assert.strictEqual(chunk, dataToWrite); + assert.deepStrictEqual(chunk, dataToWrite); }) }); diff --git a/test/parallel/test-webstreams-adapters-sync-write-error.js b/test/parallel/test-webstreams-adapters-sync-write-error.js index 748f682365ee..ff6dbddf377f 100644 --- a/test/parallel/test-webstreams-adapters-sync-write-error.js +++ b/test/parallel/test-webstreams-adapters-sync-write-error.js @@ -1,6 +1,6 @@ 'use strict'; // Flags: --no-warnings --expose-internals -require('../common'); +const common = require('../common'); const assert = require('assert'); const test = require('node:test'); const { Duplex, Writable } = require('stream'); @@ -9,6 +9,13 @@ const { newReadableWritablePairFromDuplex, } = require('internal/webstreams/adapters'); +function isSameError(expected) { + return common.mustCall((actual) => { + assert.strictEqual(actual, expected); + return true; + }); +} + // Verify that when the underlying Node.js stream throws synchronously from // write(), the writable web stream properly rejects but does not destroy // the stream (destroy-on-sync-throw is only used internally by @@ -34,6 +41,72 @@ test('WritableStream from Node.js stream handles sync write throw', async () => assert.strictEqual(writable.destroyed, false); }); +test('WritableStream from Node.js stream handles async write error', async () => { + const error = new Error('boom'); + const writable = new Writable({ + write(_chunk, _encoding, callback) { + setImmediate(callback, error); + }, + }); + const writer = Writable.toWeb(writable).getWriter(); + + await Promise.all([ + assert.rejects(writer.write(Buffer.from('hello')), isSameError(error)), + assert.rejects(writer.closed, isSameError(error)), + ]); +}); + +test('WritableStream aborts while a native write is pending', async () => { + const error = new Error('abort'); + let finishWrite; + let startWrite; + const writeStarted = new Promise((resolve) => { + startWrite = resolve; + }); + const writable = new Writable({ + write(_chunk, _encoding, callback) { + finishWrite = callback; + startWrite(); + }, + }); + const writer = Writable.toWeb(writable).getWriter(); + const writePromise = writer.write(Buffer.from('hello')); + await writeStarted; + + const writeRejected = assert.rejects(writePromise, isSameError(error)); + const closedRejected = assert.rejects(writer.closed, isSameError(error)); + await Promise.all([ + writer.abort(error), + writeRejected, + closedRejected, + ]); + + finishWrite(); + await new Promise(setImmediate); + assert.strictEqual(writable.destroyed, true); +}); + +test('WritableStream handles destruction while a write is pending', async () => { + const error = new Error('destroy'); + let startWrite; + const writeStarted = new Promise((resolve) => { + startWrite = resolve; + }); + const writable = new Writable({ + write(_chunk, _encoding, _callback) { + startWrite(); + }, + }); + const writer = Writable.toWeb(writable).getWriter(); + const writePromise = writer.write(Buffer.from('hello')); + await writeStarted; + + const writeRejected = assert.rejects(writePromise, isSameError(error)); + const closedRejected = assert.rejects(writer.closed, isSameError(error)); + writable.destroy(error); + await Promise.all([writeRejected, closedRejected]); +}); + test('Duplex-backed pair does NOT destroy on sync write throw', async () => { const error = new TypeError('invalid chunk'); const duplex = new Duplex({ diff --git a/test/parallel/test-webstreams-adapters-writable-buffer-sources.js b/test/parallel/test-webstreams-adapters-writable-buffer-sources.js index 995db97e7473..1a4a8eeb01f8 100644 --- a/test/parallel/test-webstreams-adapters-writable-buffer-sources.js +++ b/test/parallel/test-webstreams-adapters-writable-buffer-sources.js @@ -3,11 +3,46 @@ const common = require('../common'); const assert = require('assert'); const { Buffer } = require('buffer'); +const { ServerResponse } = require('http'); const { Duplex, Writable } = require('stream'); const { suite, test } = require('node:test'); const ctors = [ArrayBuffer, SharedArrayBuffer]; +function createServerResponse() { + const socket = new Writable({ + write(chunk, encoding, callback) { + callback(); + }, + }); + const response = new ServerResponse({ + method: 'GET', + httpVersionMajor: 1, + httpVersionMinor: 1, + }); + socket.on('error', common.mustNotCall()); + response.on('error', common.mustNotCall()); + response.assignSocket(socket); + return response; +} + +async function completesWithin(promise) { + let timer; + try { + return await Promise.race([ + promise, + new Promise((_, reject) => { + timer = setTimeout( + () => reject(new Error('write timed out')), + common.platformTimeout(1000), + ); + }), + ]); + } finally { + clearTimeout(timer); + } +} + suite('underlying Writable', () => { suite('in non-object mode', () => { for (const ctor of ctors) { @@ -26,6 +61,241 @@ suite('underlying Writable', () => { await writer.write(buffer); }); } + + for (const highWaterMark of [1, 16]) { + test(`waits for mutable chunks with highWaterMark ${highWaterMark}`, + async () => { + let finishWrite; + let consumed; + const writable = new Writable({ + highWaterMark, + write(chunk, encoding, callback) { + finishWrite = () => { + consumed = Buffer.from(chunk); + callback(); + }; + }, + }); + writable.on('error', common.mustNotCall()); + const writer = Writable.toWeb(writable).getWriter(); + const input = new Uint8Array([1, 2, 3, 4]); + let settled = false; + const writePromise = writer.write(input).then(() => { + settled = true; + }); + + await new Promise(setImmediate); + assert.strictEqual(settled, false); + finishWrite(); + await writePromise; + input.fill(9); + await writer.close(); + + assert.deepStrictEqual( + consumed, + Buffer.from([1, 2, 3, 4]), + ); + }); + } + + test('copies mutable chunks while the Writable is corked', async () => { + let consumed; + const writable = new Writable({ + write(chunk, encoding, callback) { + consumed = Buffer.from(chunk); + callback(); + }, + }); + writable.on('error', common.mustNotCall()); + writable.cork(); + const writer = Writable.toWeb(writable).getWriter(); + const input = new Uint8Array([1, 2, 3, 4]); + + await writer.write(input); + input.fill(9); + writable.uncork(); + await writer.close(); + + assert.deepStrictEqual(consumed, Buffer.from([1, 2, 3, 4])); + }); + + test('copies mutable chunks when write() is overridden', async () => { + let consumed; + let received; + let notifyConsumed; + const consumedPromise = new Promise((resolve) => { + notifyConsumed = resolve; + }); + const writable = new Writable({ + write(chunk, encoding, callback) { + callback(); + }, + }); + writable.on('error', common.mustNotCall()); + const writer = Writable.toWeb(writable).getWriter(); + writable.write = common.mustCall((chunk) => { + received = chunk; + setImmediate(() => { + consumed = Buffer.from(chunk); + notifyConsumed(); + }); + return true; + }); + const input = new Uint8Array([1, 2, 3, 4]); + + await writer.write(input); + input.fill(9); + await consumedPromise; + await writer.close(); + + assert.notStrictEqual(received.buffer, input.buffer); + assert.deepStrictEqual(consumed, Buffer.from([1, 2, 3, 4])); + }); + + test('does not trust a patched Writable.prototype.write', async () => { + const originalWrite = Writable.prototype.write; + let consumed; + let notifyConsumed; + const consumedPromise = new Promise((resolve) => { + notifyConsumed = resolve; + }); + const writable = new Writable({ + write(chunk, encoding, callback) { + setImmediate(() => { + consumed = Buffer.from(chunk); + callback(); + notifyConsumed(); + }); + }, + }); + writable.on('error', common.mustNotCall()); + let writer; + + Writable.prototype.write = function(chunk) { + return originalWrite.call(this, chunk); + }; + try { + writer = Writable.toWeb(writable).getWriter(); + const input = new Uint8Array([1, 2, 3, 4]); + await completesWithin(writer.write(input)); + input.fill(9); + await consumedPromise; + + assert.deepStrictEqual(consumed, Buffer.from([1, 2, 3, 4])); + } finally { + Writable.prototype.write = originalWrite; + } + await writer.close(); + }); + + test('reads write() once for classification and invocation', async () => { + const originalWrite = Writable.prototype.write; + let consumed; + const writable = new Writable({ + write(chunk, encoding, callback) { + setImmediate(() => { + consumed = Buffer.from(chunk); + callback(); + }); + }, + }); + writable.on('error', common.mustNotCall()); + const writer = Writable.toWeb(writable).getWriter(); + let writeAccesses = 0; + Object.defineProperty(writable, 'write', { + configurable: true, + get() { + writeAccesses++; + return writeAccesses === 1 ? + originalWrite : + function callbacklessWrite(chunk) { + return originalWrite.call(this, chunk); + }; + }, + }); + + const input = new Uint8Array([1, 2, 3, 4]); + await completesWithin(writer.write(input)); + input.fill(9); + + assert.strictEqual(writeAccesses, 1); + assert.deepStrictEqual(consumed, Buffer.from([1, 2, 3, 4])); + delete writable.write; + await writer.close(); + }); + + test('preserves cloned view brands and SharedArrayBuffer backing', + async () => { + const dataView = new DataView( + Uint8Array.from([0, 1, 2, 3, 4, 0]).buffer, + 1, + 4, + ); + const uint16Buffer = new ArrayBuffer(6); + const uint16 = new Uint16Array(uint16Buffer, 2, 2); + new Uint8Array(uint16Buffer, 2, 4).set([1, 2, 3, 4]); + const shared = new SharedArrayBuffer(6); + const sharedView = new Uint8Array(shared, 1, 4); + sharedView.set([1, 2, 3, 4]); + const inputs = [ + Buffer.from([1, 2, 3, 4]), + dataView, + uint16, + sharedView, + ]; + const expected = inputs.map((chunk) => ({ + brand: Buffer.isBuffer(chunk) ? + 'Buffer' : Object.prototype.toString.call(chunk), + bytes: Buffer.from(new Uint8Array( + chunk.buffer, + chunk.byteOffset, + chunk.byteLength, + )), + shared: chunk.buffer instanceof SharedArrayBuffer, + })); + const received = []; + const writable = new Writable({ + write(chunk, encoding, callback) { + callback(); + }, + }); + writable.on('error', common.mustNotCall()); + const writer = Writable.toWeb(writable).getWriter(); + writable.write = common.mustCall((chunk) => { + received.push(chunk); + return true; + }, inputs.length); + + for (const chunk of inputs) { + await writer.write(chunk); + new Uint8Array( + chunk.buffer, + chunk.byteOffset, + chunk.byteLength, + ).fill(9); + } + await writer.close(); + + for (let i = 0; i < received.length; i++) { + const actual = received[i]; + const actualBrand = Buffer.isBuffer(actual) ? + 'Buffer' : Object.prototype.toString.call(actual); + assert.strictEqual(actualBrand, expected[i].brand); + assert.notStrictEqual(actual.buffer, inputs[i].buffer); + assert.strictEqual( + actual.buffer instanceof SharedArrayBuffer, + expected[i].shared, + ); + assert.deepStrictEqual( + Buffer.from(new Uint8Array( + actual.buffer, + actual.byteOffset, + actual.byteLength, + )), + expected[i].bytes, + ); + } + }); }); suite('in object mode', () => { @@ -48,16 +318,71 @@ suite('underlying Writable', () => { }); }); +suite('underlying ServerResponse', () => { + test('rejects invalid view types before cloning', async () => { + const response = createServerResponse(); + const writer = Writable.toWeb(response).getWriter(); + + try { + await Promise.all([ + assert.rejects(writer.write(new DataView(new ArrayBuffer(4))), { + code: 'ERR_INVALID_ARG_TYPE', + }), + assert.rejects(writer.closed, { + code: 'ERR_INVALID_ARG_TYPE', + }), + ]); + } finally { + await new Promise((resolve) => response.end(resolve)); + } + }); + + test('preserves write() overrides', async () => { + const response = createServerResponse(); + const writer = Writable.toWeb(response).getWriter(); + const originalWrite = response.write; + let received; + response.write = common.mustCall((chunk) => { + received = chunk; + return true; + }); + const input = new DataView(Uint8Array.from([1, 2, 3, 4]).buffer); + + try { + await writer.write(input); + new Uint8Array(input.buffer).fill(9); + assert(received instanceof DataView); + assert.notStrictEqual(received.buffer, input.buffer); + assert.deepStrictEqual( + Buffer.from(received.buffer), + Buffer.from([1, 2, 3, 4]), + ); + } finally { + response.write = originalWrite; + await new Promise((resolve) => response.end(resolve)); + } + }); +}); + suite('underlying Duplex', () => { suite('in non-object mode', () => { for (const ctor of ctors) { - test(`converts ${ctor.name} chunks`, async () => { + test(`copies ${ctor.name} chunks`, async () => { const buffer = new ctor(4); + new Uint8Array(buffer).set([1, 2, 3, 4]); const duplex = new Duplex({ writableObjectMode: false, write: common.mustCall((chunk, encoding, callback) => { assert(Buffer.isBuffer(chunk)); - assert.strictEqual(chunk.buffer, buffer); + assert.notStrictEqual(chunk.buffer, buffer); + assert.strictEqual( + chunk.buffer instanceof SharedArrayBuffer, + buffer instanceof SharedArrayBuffer, + ); + assert.deepStrictEqual( + chunk, + Buffer.from([1, 2, 3, 4]), + ); callback(); }), read() { @@ -69,6 +394,37 @@ suite('underlying Duplex', () => { await writer.write(buffer); }); } + + test('copies mutable chunks without waiting for the write callback', + async () => { + let consumed; + let finishWrite; + const duplex = new Duplex({ + write(chunk, encoding, callback) { + finishWrite = () => { + consumed = Buffer.from(chunk); + callback(); + }; + }, + read() { + this.push(null); + }, + }); + duplex.on('error', common.mustNotCall()); + const writer = Duplex.toWeb(duplex).writable.getWriter(); + const input = new Uint8Array([1, 2, 3, 4]); + + await writer.write(input); + input.fill(9); + finishWrite(); + await new Promise(setImmediate); + + assert.deepStrictEqual( + consumed, + Buffer.from([1, 2, 3, 4]), + ); + await writer.close(); + }); }); suite('in object mode', () => { diff --git a/test/parallel/test-webstreams-compression-buffer-source.js b/test/parallel/test-webstreams-compression-buffer-source.js index 3304a8e64f31..a81d4b674b38 100644 --- a/test/parallel/test-webstreams-compression-buffer-source.js +++ b/test/parallel/test-webstreams-compression-buffer-source.js @@ -3,6 +3,7 @@ require('../common'); const assert = require('assert'); const test = require('node:test'); const { DecompressionStream, CompressionStream } = require('stream/web'); +const { gzipSync } = require('zlib'); // Minimal gzip-compressed bytes for "hello" const compressedGzip = new Uint8Array([ @@ -40,3 +41,31 @@ test('CompressionStream round-trip with ArrayBuffer input', async () => { const result = Buffer.concat(out.map((c) => Buffer.from(c))); assert.strictEqual(result.toString(), 'hello'); }); + +test('DecompressionStream writable completion is not coupled to readable ' + + 'backpressure', async () => { + const expected = Buffer.alloc(1024 * 1024, 0x61); + const compressed = gzipSync(expected); + const ds = new DecompressionStream('gzip'); + const writer = ds.writable.getWriter(); + let settled = false; + const writePromise = writer.write(compressed).then(() => { + settled = true; + }); + + await new Promise(setImmediate); + const settledBeforeRead = settled; + if (settledBeforeRead) { + compressed.fill(0); + } + + const outputPromise = Array.fromAsync(ds.readable); + await writePromise; + await writer.close(); + const output = Buffer.concat( + (await outputPromise).map((chunk) => Buffer.from(chunk)), + ); + + assert.strictEqual(settledBeforeRead, true); + assert.deepStrictEqual(output, expected); +}); From cb0aad6af7a9a2c137419dd428200861fe519b27 Mon Sep 17 00:00:00 2001 From: seungwoo Date: Mon, 28 Sep 2026 19:38:06 +0900 Subject: [PATCH 03/15] stream: simplify mutable chunk fallback Rely on native stream validation after copying array buffer views. Keep object-mode writes on their existing completion timing, and remove adapter-only HTTP detection and prototype capture. Signed-off-by: seungwoo --- lib/_http_outgoing.js | 3 --- lib/internal/streams/utils.js | 1 - lib/internal/webstreams/adapters.js | 16 +--------------- 3 files changed, 1 insertion(+), 19 deletions(-) diff --git a/lib/_http_outgoing.js b/lib/_http_outgoing.js index 817e4ddf49f2..2ca455eb5566 100644 --- a/lib/_http_outgoing.js +++ b/lib/_http_outgoing.js @@ -1038,8 +1038,6 @@ OutgoingMessage.prototype.write = function write(chunk, encoding, callback) { return ret; }; -const outgoingMessagePrototypeWrite = OutgoingMessage.prototype.write; - function onError(msg, err, callback) { if (msg.destroyed) { return; @@ -1490,5 +1488,4 @@ module.exports = { validateHeaderName, validateHeaderValue, OutgoingMessage, - outgoingMessagePrototypeWrite, }; diff --git a/lib/internal/streams/utils.js b/lib/internal/streams/utils.js index 06b05bab0184..45f55316104f 100644 --- a/lib/internal/streams/utils.js +++ b/lib/internal/streams/utils.js @@ -345,7 +345,6 @@ module.exports = { isWritableEnded, isWritableFinished, isWritableErrored, - isOutgoingMessage, isServerRequest, isServerResponse, willEmitClose, diff --git a/lib/internal/webstreams/adapters.js b/lib/internal/webstreams/adapters.js index 0fd99fafdb98..ffa30eafeacc 100644 --- a/lib/internal/webstreams/adapters.js +++ b/lib/internal/webstreams/adapters.js @@ -39,7 +39,6 @@ const { const { TextEncoder } = require('internal/encoding'); let addAbortListener; -let outgoingMessagePrototypeWrite; const { ReadableStream, @@ -71,7 +70,6 @@ const { const { isDestroyed, - isOutgoingMessage, isReadable, isWritable, isWritableEnded, @@ -93,7 +91,6 @@ const { isArrayBufferView, isDataView, isSharedArrayBuffer, - isUint8Array, } = require('internal/util/types'); const { @@ -615,24 +612,13 @@ function newWritableStreamFromStreamWritable( } needDrainBeforeWrite = streamWritable.writableNeedDrain; const isArrayBufferViewChunk = isArrayBufferView(chunk); - if (!objectMode && - isOutgoingMessage(streamWritable) && - isArrayBufferViewChunk && - !isUint8Array(chunk)) { - outgoingMessagePrototypeWrite ??= - require('_http_outgoing').outgoingMessagePrototypeWrite; - if (writeMethod === outgoingMessagePrototypeWrite) { - throw new ERR_INVALID_ARG_TYPE( - 'chunk', ['string', 'Buffer', 'Uint8Array'], chunk); - } - } const canAwaitWriteCallback = streamWritable instanceof Writable && writeMethod === writablePrototypeWrite && streamWritable._readableState === undefined && writableState?.corked === 0; waitForWriteCallback = - isBufferSourceChunk && canAwaitWriteCallback; + !objectMode && isBufferSourceChunk && canAwaitWriteCallback; if (!objectMode && isArrayBufferViewChunk && !waitForWriteCallback) { From 055036c5627c2a11d8aa55551ef43b80a902390a Mon Sep 17 00:00:00 2001 From: seungwoo Date: Mon, 28 Sep 2026 19:41:01 +0900 Subject: [PATCH 04/15] test: refine mutable Web Stream adapter coverage Remove overlapping override coverage. Verify object-mode writes settle independently of native callbacks. Exercise the Writable.toWeb(Duplex) path and keep native HTTP validation coverage. Assisted-by: Codex Signed-off-by: seungwoo --- ...treams-adapters-writable-buffer-sources.js | 60 +++++++------------ 1 file changed, 21 insertions(+), 39 deletions(-) diff --git a/test/parallel/test-webstreams-adapters-writable-buffer-sources.js b/test/parallel/test-webstreams-adapters-writable-buffer-sources.js index 1a4a8eeb01f8..7b6ed6a41026 100644 --- a/test/parallel/test-webstreams-adapters-writable-buffer-sources.js +++ b/test/parallel/test-webstreams-adapters-writable-buffer-sources.js @@ -119,39 +119,6 @@ suite('underlying Writable', () => { assert.deepStrictEqual(consumed, Buffer.from([1, 2, 3, 4])); }); - test('copies mutable chunks when write() is overridden', async () => { - let consumed; - let received; - let notifyConsumed; - const consumedPromise = new Promise((resolve) => { - notifyConsumed = resolve; - }); - const writable = new Writable({ - write(chunk, encoding, callback) { - callback(); - }, - }); - writable.on('error', common.mustNotCall()); - const writer = Writable.toWeb(writable).getWriter(); - writable.write = common.mustCall((chunk) => { - received = chunk; - setImmediate(() => { - consumed = Buffer.from(chunk); - notifyConsumed(); - }); - return true; - }); - const input = new Uint8Array([1, 2, 3, 4]); - - await writer.write(input); - input.fill(9); - await consumedPromise; - await writer.close(); - - assert.notStrictEqual(received.buffer, input.buffer); - assert.deepStrictEqual(consumed, Buffer.from([1, 2, 3, 4])); - }); - test('does not trust a patched Writable.prototype.write', async () => { const originalWrite = Writable.prototype.write; let consumed; @@ -302,24 +269,31 @@ suite('underlying Writable', () => { for (const ctor of ctors) { test(`passes through ${ctor.name} chunks`, async () => { const buffer = new ctor(4); + let finishWrite; const writable = new Writable({ objectMode: true, write: common.mustCall((chunk, encoding, callback) => { assert(chunk instanceof ctor); assert.strictEqual(chunk, buffer); - callback(); + finishWrite = callback; }), }); writable.on('error', common.mustNotCall()); const writer = Writable.toWeb(writable).getWriter(); - await writer.write(buffer); + const writePromise = writer.write(buffer); + try { + await completesWithin(writePromise); + } finally { + finishWrite(); + } + await writer.close(); }); } }); }); suite('underlying ServerResponse', () => { - test('rejects invalid view types before cloning', async () => { + test('rejects invalid view types', async () => { const response = createServerResponse(); const writer = Writable.toWeb(response).getWriter(); @@ -411,7 +385,8 @@ suite('underlying Duplex', () => { }, }); duplex.on('error', common.mustNotCall()); - const writer = Duplex.toWeb(duplex).writable.getWriter(); + const writer = Writable.toWeb(duplex).getWriter(); + duplex.resume(); const input = new Uint8Array([1, 2, 3, 4]); await writer.write(input); @@ -431,12 +406,13 @@ suite('underlying Duplex', () => { for (const ctor of ctors) { test(`passes through ${ctor.name} chunks`, async () => { const buffer = new ctor(4); + let finishWrite; const duplex = new Duplex({ writableObjectMode: true, write: common.mustCall((chunk, encoding, callback) => { assert(chunk instanceof ctor); assert.strictEqual(chunk, buffer); - callback(); + finishWrite = callback; }), read() { this.push(null); @@ -444,7 +420,13 @@ suite('underlying Duplex', () => { }); duplex.on('error', common.mustNotCall()); const writer = Duplex.toWeb(duplex).writable.getWriter(); - await writer.write(buffer); + const writePromise = writer.write(buffer); + try { + await completesWithin(writePromise); + } finally { + finishWrite(); + } + await writer.close(); }); } }); From 702bac3def8ece4c69ecffd301f86dbbbb827fa3 Mon Sep 17 00:00:00 2001 From: seungwoo Date: Tue, 29 Sep 2026 19:15:49 +0900 Subject: [PATCH 05/15] stream: use Float16Array from primordials Use the Float16Array constructor captured during Node.js initialization instead of reading it from globalThis when the adapter is loaded. This prevents changes to the global constructor from affecting chunk cloning. Signed-off-by: seungwoo --- lib/internal/webstreams/adapters.js | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/lib/internal/webstreams/adapters.js b/lib/internal/webstreams/adapters.js index ffa30eafeacc..cd4d8a3c8c00 100644 --- a/lib/internal/webstreams/adapters.js +++ b/lib/internal/webstreams/adapters.js @@ -6,6 +6,7 @@ const { BigInt64Array, BigUint64Array, DataView, + Float16Array, Float32Array, Float64Array, FunctionPrototypeCall, @@ -31,9 +32,6 @@ const { Uint32Array, Uint8Array, Uint8ClampedArray, - globalThis: { - Float16Array, - }, } = primordials; const { TextEncoder } = require('internal/encoding'); From 55118d54f6e059f1c5ea56cf2174bcffeaeda1cd Mon Sep 17 00:00:00 2001 From: seungwoo Date: Thu, 1 Oct 2026 20:26:23 +0900 Subject: [PATCH 06/15] stream: optimize private BufferSource copies Use private Buffer storage for ArrayBuffer-backed copies and reuse it for Buffer inputs. Read view metadata through intrinsic getters and avoid redundant type checks. Preserve view types, raw bytes, and independent backing. Cover initialized backing without requiring an exact allocation size. Signed-off-by: seungwoo --- lib/internal/webstreams/adapters.js | 50 +++-- test/parallel/test-stream-duplex.js | 2 + .../test-webstreams-adapters-writable-copy.js | 195 ++++++++++++++++++ 3 files changed, 230 insertions(+), 17 deletions(-) create mode 100644 test/parallel/test-webstreams-adapters-writable-copy.js diff --git a/lib/internal/webstreams/adapters.js b/lib/internal/webstreams/adapters.js index cd4d8a3c8c00..dd022b357718 100644 --- a/lib/internal/webstreams/adapters.js +++ b/lib/internal/webstreams/adapters.js @@ -6,6 +6,9 @@ const { BigInt64Array, BigUint64Array, DataView, + DataViewPrototypeGetBuffer, + DataViewPrototypeGetByteLength, + DataViewPrototypeGetByteOffset, Float16Array, Float32Array, Float64Array, @@ -77,10 +80,9 @@ const { Buffer, } = require('buffer'); +const { createUnsafeBuffer } = require('internal/buffer'); + const { - ArrayBufferViewGetBuffer, - ArrayBufferViewGetByteLength, - ArrayBufferViewGetByteOffset, kResolvedPromise, } = require('internal/webstreams/util'); @@ -269,33 +271,49 @@ function dispatchWritev(writer, chunks, tracker) { function cloneArrayBufferView(view) { const viewIsDataView = isDataView(view); - const buffer = ArrayBufferViewGetBuffer(view); + const buffer = viewIsDataView ? + DataViewPrototypeGetBuffer(view) : TypedArrayPrototypeGetBuffer(view); + const shared = isSharedArrayBuffer(buffer); // A detached backing buffer cannot be changed by the caller. Passing the // original view through also preserves the native stream's validation. - if (!isSharedArrayBuffer(buffer) && + if (!shared && ArrayBufferPrototypeGetDetached(buffer)) { return view; } - const byteLength = ArrayBufferViewGetByteLength(view); - const byteOffset = ArrayBufferViewGetByteOffset(view); - const copiedBytes = new Uint8Array( - isSharedArrayBuffer(buffer) ? - constructSharedArrayBuffer(byteLength) : - byteLength, - ); + const tag = viewIsDataView ? undefined : TypedArrayPrototypeGetSymbolToStringTag(view); + const byteLength = viewIsDataView ? + DataViewPrototypeGetByteLength(view) : TypedArrayPrototypeGetByteLength(view); + const byteOffset = viewIsDataView ? + DataViewPrototypeGetByteOffset(view) : TypedArrayPrototypeGetByteOffset(view); + // Every byte is overwritten before exposing the private, unpooled buffer. + const copiedBytes = shared ? + new Uint8Array(constructSharedArrayBuffer(byteLength)) : + createUnsafeBuffer(byteLength); + // Empty views can be out of bounds; shared views can grow. Keep a fixed + // byte range for both, rather than passing the original view to set(). TypedArrayPrototypeSet( copiedBytes, - new Uint8Array(buffer, byteOffset, byteLength), + !shared && tag === 'Uint8Array' && byteLength > 0 ? + view : new Uint8Array(buffer, byteOffset, byteLength), ); - const copiedBuffer = ArrayBufferViewGetBuffer(copiedBytes); + if (tag === 'Uint8Array') { + if (Buffer.isBuffer(view)) { + // Reuse the Buffer without materializing its backing for small chunks. + return shared ? + Buffer.from(TypedArrayPrototypeGetBuffer(copiedBytes)) : copiedBytes; + } + return shared ? copiedBytes : + new Uint8Array(TypedArrayPrototypeGetBuffer(copiedBytes)); + } + const copiedBuffer = TypedArrayPrototypeGetBuffer(copiedBytes); if (viewIsDataView) { return new DataView(copiedBuffer); } - switch (TypedArrayPrototypeGetSymbolToStringTag(view)) { + switch (tag) { case 'BigInt64Array': return new BigInt64Array(copiedBuffer); case 'BigUint64Array': @@ -312,8 +330,6 @@ function cloneArrayBufferView(view) { return new Int16Array(copiedBuffer); case 'Int32Array': return new Int32Array(copiedBuffer); - case 'Uint8Array': - return Buffer.isBuffer(view) ? Buffer.from(copiedBuffer) : copiedBytes; case 'Uint8ClampedArray': return new Uint8ClampedArray(copiedBuffer); case 'Uint16Array': diff --git a/test/parallel/test-stream-duplex.js b/test/parallel/test-stream-duplex.js index 9908d000cb7e..001cae48ca14 100644 --- a/test/parallel/test-stream-duplex.js +++ b/test/parallel/test-stream-duplex.js @@ -121,6 +121,7 @@ process.on('exit', () => { }, write: common.mustCall((chunk) => { assert.deepStrictEqual(chunk, dataToWrite); + assert.notStrictEqual(chunk, dataToWrite); }) }); @@ -144,6 +145,7 @@ process.on('exit', () => { }, write: common.mustCall((chunk) => { assert.deepStrictEqual(chunk, dataToWrite); + assert.notStrictEqual(chunk, dataToWrite); }) }); diff --git a/test/parallel/test-webstreams-adapters-writable-copy.js b/test/parallel/test-webstreams-adapters-writable-copy.js new file mode 100644 index 000000000000..139708c50f8f --- /dev/null +++ b/test/parallel/test-webstreams-adapters-writable-copy.js @@ -0,0 +1,195 @@ +'use strict'; +const common = require('../common'); +const assert = require('assert'); +const { PassThrough, Writable } = require('stream'); +const { test } = require('node:test'); +const { isDataView } = require('util').types; +const vm = require('vm'); + +const typedArrayPrototype = Object.getPrototypeOf(Uint8Array.prototype); +const typedArrayGetters = {}; +const dataViewGetters = {}; +for (const key of ['buffer', 'byteLength', 'byteOffset']) { + typedArrayGetters[key] = Function.call.bind( + Object.getOwnPropertyDescriptor(typedArrayPrototype, key).get); + dataViewGetters[key] = Function.call.bind( + Object.getOwnPropertyDescriptor(DataView.prototype, key).get); +} + +function metadata(view, key) { + return (isDataView(view) ? dataViewGetters : typedArrayGetters)[key](view); +} + +function bytes(view) { + return new Uint8Array( + metadata(view, 'buffer'), + metadata(view, 'byteOffset'), + metadata(view, 'byteLength'), + ); +} + +function assertInitializedCopy(view, expected) { + const buffer = metadata(view, 'buffer'); + assert.strictEqual(metadata(view, 'byteLength'), expected.length); + const initialized = new Uint8Array(buffer.byteLength); + initialized.set(expected, metadata(view, 'byteOffset')); + assert.deepStrictEqual(new Uint8Array(buffer), initialized); +} + +async function capture(input) { + let copied; + const writable = new Writable({ write: common.mustNotCall() }); + writable.write = common.mustCall((chunk) => { + copied = chunk; + return true; + }); + writable.on('error', common.mustNotCall()); + const writer = Writable.toWeb(writable).getWriter(); + await writer.write(input); + await writer.close(); + return copied; +} + +for (const size of [0, 64, 65]) { + test(`Buffer copy has independent, fully initialized backing: ${size} bytes`, async () => { + const input = Buffer.allocUnsafe(size + 32).subarray(16, size + 16); + for (let i = 0; i < size; i++) input[i] = i % 251; + const expected = Array.from(input); + const copied = await capture(input); + + assert(Buffer.isBuffer(copied)); + assert.notStrictEqual(copied, input); + assert.notStrictEqual(copied.buffer, input.buffer); + assertInitializedCopy(copied, expected); + new Uint8Array(input.buffer).fill(255); + assert.deepStrictEqual(Array.from(copied), expected); + }); +} + +const constructors = [ + Int8Array, Uint8Array, Uint8ClampedArray, Int16Array, Uint16Array, + Int32Array, Uint32Array, Float16Array, Float32Array, Float64Array, + BigInt64Array, BigUint64Array, DataView, Buffer, +]; +for (const Backing of [ArrayBuffer, SharedArrayBuffer]) { + for (const Constructor of constructors) { + test(`${Backing.name} copy preserves ${Constructor.name} and raw bytes`, async () => { + const buffer = new Backing(64); + const input = Constructor === Buffer ? Buffer.from(buffer, 8, 16) : + Constructor === DataView ? new DataView(buffer, 8, 16) : + new Constructor(buffer, 8, 16 / Constructor.BYTES_PER_ELEMENT); + const expected = Array.from({ length: 16 }, (_, i) => i + 1); + new Uint8Array(buffer, 8, 16).set(expected); + const copied = await capture(input); + + new Uint8Array(buffer).fill(255); + assert.notStrictEqual(metadata(copied, 'buffer'), buffer); + assertInitializedCopy(copied, expected); + assert.strictEqual(metadata(copied, 'buffer') instanceof Backing, true); + assert.strictEqual(Buffer.isBuffer(copied), Constructor === Buffer); + assert.strictEqual(Object.prototype.toString.call(copied), Object.prototype.toString.call(input)); + assert.deepStrictEqual(Array.from(bytes(copied)), expected); + }); + } +} + +for (const Constructor of [Uint8Array, Buffer]) { + for (const length of [0, 8]) { + test(`out-of-bounds resizable ${Constructor.name} remains an empty copy: ${length}`, async () => { + const buffer = new ArrayBuffer(32, { maxByteLength: 64 }); + const input = Constructor === Buffer ? Buffer.from(buffer, 16, length) : + new Uint8Array(buffer, 16, length); + buffer.resize(8); + const copied = await capture(input); + + assert.strictEqual(copied.byteLength, 0); + assert.notStrictEqual(copied.buffer, buffer); + assert.strictEqual(Buffer.isBuffer(copied), Constructor === Buffer); + }); + } +} + +for (const Backing of [ArrayBuffer, SharedArrayBuffer]) { + for (const tracking of [false, true]) { + test(`${Backing.name} copy survives resize or grow, tracking=${tracking}`, async () => { + const buffer = new Backing(32, { maxByteLength: 64 }); + const input = tracking ? new Uint8Array(buffer, 8) : new Uint8Array(buffer, 8, 16); + input.fill(7); + const expected = Array.from(input); + const copied = await capture(input); + + if (Backing === ArrayBuffer) buffer.resize(8); + else buffer.grow(64); + new Uint8Array(buffer).fill(255); + assert.notStrictEqual(copied.buffer, buffer); + assert.strictEqual(copied.buffer.resizable ?? copied.buffer.growable, false); + assert.deepStrictEqual(Array.from(copied), expected); + }); + } +} + +for (const Constructor of [Float16Array, Float32Array, Float64Array]) { + test(`${Constructor.name} copy preserves raw NaN payloads`, async () => { + const input = new Constructor(8); + bytes(input).fill(255); + const expected = Array.from(bytes(input)); + const copied = await capture(input); + + assert(copied instanceof Constructor); + assert.deepStrictEqual(Array.from(bytes(copied)), expected); + }); +} + +for (const Constructor of [Uint8Array, Buffer]) { + test(`${Constructor.name} copy ignores shadowed view metadata`, async () => { + const input = Constructor === Buffer ? Buffer.from([1, 2, 3, 4]) : new Uint8Array([1, 2, 3, 4]); + Object.defineProperty(input, 'byteLength', { value: 99 }); + for (const key of ['buffer', 'byteOffset', 'length', 'constructor']) { + Object.defineProperty(input, key, { get: common.mustNotCall(`unexpected ${key} access`) }); + } + const copied = await capture(input); + + assert.deepStrictEqual(Array.from(bytes(copied)), [1, 2, 3, 4]); + assert.notStrictEqual(metadata(copied, 'buffer'), metadata(input, 'buffer')); + }); +} + +for (const Constructor of [Uint16Array, Float32Array]) { + test(`${Constructor.name} with Buffer prototype retains its intrinsic type`, async () => { + const input = new Constructor([1, 2, 3, 4]); + Object.setPrototypeOf(input, Buffer.prototype); + assert(Buffer.isBuffer(input)); + const expected = Array.from(bytes(input)); + const copied = await capture(input); + + assert(copied instanceof Constructor); + assert.strictEqual(Buffer.isBuffer(copied), false); + assert.deepStrictEqual(Array.from(bytes(copied)), expected); + }); +} + +test('cross-realm Uint8Array copy has independent bytes', async () => { + const input = vm.runInNewContext('new Uint8Array([1, 2, 3, 4])'); + const copied = await capture(input); + input.fill(255); + + assert(copied instanceof Uint8Array); + assert.strictEqual(Buffer.isBuffer(copied), false); + assert.notStrictEqual(copied.buffer, input.buffer); + assert.deepStrictEqual(Array.from(copied), [1, 2, 3, 4]); +}); + +test('PassThrough retains private bytes after awaited write', async () => { + const stream = new PassThrough(); + const writer = Writable.toWeb(stream).getWriter(); + const input = Buffer.from([1, 2, 3, 4]); + await writer.write(input); + new Uint8Array(input.buffer).fill(255); + const result = stream.read(); + const closed = writer.close(); + stream.resume(); + await closed; + + assert.notStrictEqual(result.buffer, input.buffer); + assert.deepStrictEqual(Array.from(result), [1, 2, 3, 4]); +}); From 38e7a2f521d7b7463c04b529559b4fdd8bdd36f6 Mon Sep 17 00:00:00 2001 From: seungwoo Date: Thu, 1 Oct 2026 20:28:10 +0900 Subject: [PATCH 07/15] stream: optimize Web Stream write completion Observe synchronous native write completion without waiting for a public write callback tick. Resume pending writes without allocating a promise capability when async hooks and context propagation allow it. Coordinate callback completion, drain, and abort while preserving async context, Promise reaction order, and error delivery. Cover mixed completion modes, lifecycle changes, and modified Promise properties. Signed-off-by: seungwoo --- lib/internal/streams/writable.js | 5 +- lib/internal/webstreams/adapters.js | 116 +++++++---- lib/internal/webstreams/writablestream.js | 52 +++++ ...streams-adapters-writable-async-context.js | 174 ++++++++++++++++ ...-webstreams-adapters-writable-lifecycle.js | 196 ++++++++++++++++++ ...eams-adapters-writable-promise-mutation.js | 99 +++++++++ ...s-adapters-writable-prototype-pollution.js | 50 +++++ ...reams-adapters-writable-sync-completion.js | 196 ++++++++++++++++++ 8 files changed, 841 insertions(+), 47 deletions(-) create mode 100644 test/parallel/test-webstreams-adapters-writable-async-context.js create mode 100644 test/parallel/test-webstreams-adapters-writable-lifecycle.js create mode 100644 test/parallel/test-webstreams-adapters-writable-promise-mutation.js create mode 100644 test/parallel/test-webstreams-adapters-writable-prototype-pollution.js create mode 100644 test/parallel/test-webstreams-adapters-writable-sync-completion.js diff --git a/lib/internal/streams/writable.js b/lib/internal/streams/writable.js index d39981ee05f6..78a28c5e51ee 100644 --- a/lib/internal/streams/writable.js +++ b/lib/internal/streams/writable.js @@ -596,7 +596,7 @@ const kWriteFlowBlock = kObjectMode | kDestroyed | kErrored | kSync | kNeedDrain | kWriteCb | kExpectWriteCb | kBufferProcessing | kFinalCalled | kPrefinished | kOnFinished | kErrorEmitted; -function writeKnownBuffer(stream, chunk) { +function writeKnownBuffer(stream, chunk, onwrite) { const state = stream._writableState; if (state == null || state.length !== 0 || !(chunk instanceof Buffer)) return undefined; @@ -610,7 +610,7 @@ function writeKnownBuffer(stream, chunk) { state.length = len; state.writelen = len; state[kState] = bits | kWriting | kSync | kExpectWriteCb; - stream._write(chunk, 'buffer', state.onwrite); + stream._write(chunk, 'buffer', onwrite ?? state.onwrite); state[kState] &= ~kSync; const ret = state.length < state.highWaterMark || state.length === 0; @@ -1193,6 +1193,7 @@ Writable.toWeb = function(streamWritable) { streamWritable, undefined, writablePrototypeWrite, + writeKnownBuffer, ); }; diff --git a/lib/internal/webstreams/adapters.js b/lib/internal/webstreams/adapters.js index dd022b357718..b0cde3fd965e 100644 --- a/lib/internal/webstreams/adapters.js +++ b/lib/internal/webstreams/adapters.js @@ -38,6 +38,8 @@ const { } = primordials; const { TextEncoder } = require('internal/encoding'); +const AsyncContextFrame = require('internal/async_context_frame'); +const { enabledHooksExist } = require('internal/async_hooks'); let addAbortListener; @@ -53,6 +55,8 @@ const { const { WritableStream, getWritableStreamDefaultControllerSignal, + writableStreamDefaultControllerResolveWrite, + writableStreamDefaultControllerRejectWrite, isWritableStream, writableStreamDefaultWriterWriteWithRequest, } = require('internal/webstreams/writablestream'); @@ -80,10 +84,11 @@ const { Buffer, } = require('buffer'); -const { createUnsafeBuffer } = require('internal/buffer'); +const { FastBuffer, createUnsafeBuffer } = require('internal/buffer'); const { kResolvedPromise, + kParkedAlgorithmResult, } = require('internal/webstreams/util'); const { @@ -354,12 +359,14 @@ function cloneArrayBufferView(view) { * @param {Writable} streamWritable * @param {object} [options] * @param {Function} [writablePrototypeWrite] + * @param {Function} [writeKnownBufferForWeb] * @returns {WritableStream} */ function newWritableStreamFromStreamWritable( streamWritable, options = kEmptyObject, writablePrototypeWrite, + writeKnownBufferForWeb, ) { // Not using the internal/streams/utils isWritableNodeStream utility // here because it will return false if streamWritable is a Duplex @@ -384,6 +391,7 @@ function newWritableStreamFromStreamWritable( return writable; } + const hasValidationOptions = options !== kEmptyObject; const highWaterMark = streamWritable.writableHighWaterMark; const strategy = streamWritable.writableObjectMode ? @@ -402,6 +410,9 @@ function newWritableStreamFromStreamWritable( let closed; let abortSignal; let abortDisposable; + let writeInProgress = false; + // undefined: pending; null: success; otherwise: the write error. + let synchronousWriteResult; function disposeAbortListener() { abortDisposable?.[SymbolDispose](); @@ -416,12 +427,12 @@ function newWritableStreamFromStreamWritable( } function resolvePendingWrite(state) { - if (pendingWriteState !== state) { - return; - } pendingWriteState = undefined; - disposeAbortListener(); - state.deferred.resolve(); + if (state.deferred === undefined) { + writableStreamDefaultControllerResolveWrite(controller, state.contextFrame); + } else { + state.deferred.resolve(); + } } function rejectPendingWrite(error) { @@ -430,8 +441,11 @@ function newWritableStreamFromStreamWritable( return; } pendingWriteState = undefined; - disposeAbortListener(); - state.deferred.reject(error); + if (state.deferred === undefined) { + writableStreamDefaultControllerRejectWrite(controller, state.contextFrame, error); + } else { + state.deferred.reject(error); + } } function listenForPendingAbort() { @@ -464,14 +478,15 @@ function newWritableStreamFromStreamWritable( } function createPendingWriteState(waitForWriteCallback) { + // Keep promise resources for async hooks and legacy context propagation. + const canPark = AsyncContextFrame.enabled && !enabledHooksExist(); const state = { - __proto__: null, - deferred: PromiseWithResolvers(), + deferred: canPark ? undefined : PromiseWithResolvers(), + contextFrame: canPark ? AsyncContextFrame.current() : undefined, writeComplete: !waitForWriteCallback, backpressureCleared: false, }; pendingWriteState = state; - listenForPendingAbort(); return state; } @@ -485,6 +500,10 @@ function newWritableStreamFromStreamWritable( } function onWrite(error) { + if (writeInProgress) { + synchronousWriteResult = error ?? null; + return; + } const state = pendingWriteState; if (state === undefined) { return; @@ -495,6 +514,11 @@ function newWritableStreamFromStreamWritable( markWriteComplete(state, error); } + function onNativeWrite(error) { + streamWritable._writableState.onwrite(error); + onWrite(error || undefined); + } + function throwSyncWriteError(error) { disposeAbortListener(); // When the kDestroyOnSyncError flag is set (e.g. for @@ -510,9 +534,6 @@ function newWritableStreamFromStreamWritable( } function throwIfTerminated() { - if (abortSignal.aborted) { - recordTerminalReason(abortSignal.reason); - } if (isTerminated) { throw terminalReason; } @@ -524,47 +545,46 @@ function newWritableStreamFromStreamWritable( needDrainBeforeWrite, waitForWriteCallback, ) { - if (!waitForWriteCallback && !needDrainBeforeWrite) { - try { - const writeReturn = FunctionPrototypeCall( - writeMethod, - streamWritable, - chunk, - ); - throwIfTerminated(); - if (controller === undefined || - writeReturn || - !streamWritable.writableNeedDrain) { - return; - } - } catch (error) { - throwSyncWriteError(error); - } - - const state = createPendingWriteState(false); - return state.deferred.promise; - } - - const state = createPendingWriteState(waitForWriteCallback); + synchronousWriteResult = undefined; try { // The return value only reports backpressure. The callback reports when // the chunk has been handled and can safely be mutated by the caller. + writeInProgress = true; + if (waitForWriteCallback && writeKnownBufferForWeb !== undefined && + !(chunk instanceof Buffer)) { + // Match Writable.write()'s byte-mode view conversion before using + // the native Buffer fast path. This shares the original bytes. + chunk = new FastBuffer(chunk.buffer, chunk.byteOffset, chunk.byteLength); + } const writeReturn = waitForWriteCallback ? - FunctionPrototypeCall(writeMethod, streamWritable, chunk, onWrite) : + (writeKnownBufferForWeb === undefined ? + FunctionPrototypeCall(writeMethod, streamWritable, chunk, onWrite) : + (writeKnownBufferForWeb(streamWritable, chunk, onNativeWrite) ?? + FunctionPrototypeCall(writeMethod, streamWritable, chunk, onWrite))) : FunctionPrototypeCall(writeMethod, streamWritable, chunk); + writeInProgress = false; throwIfTerminated(); + if (synchronousWriteResult != null) { + throw handleKnownInternalErrors(synchronousWriteResult); + } const shouldWaitForDrain = (needDrainBeforeWrite || !writeReturn) && streamWritable.writableNeedDrain; + if ((synchronousWriteResult === null || + (!waitForWriteCallback && !needDrainBeforeWrite)) && + !shouldWaitForDrain) { + return; + } + const state = createPendingWriteState( + waitForWriteCallback && synchronousWriteResult === undefined, + ); state.backpressureCleared = !shouldWaitForDrain; maybeResolvePendingWrite(state); + return state.deferred === undefined ? kParkedAlgorithmResult : state.deferred.promise; } catch (error) { - resolvePendingWrite(state); - PromisePrototypeThen(state.deferred.promise, undefined, noop); + writeInProgress = false; throwSyncWriteError(error); } - - return state.deferred.promise; } const cleanup = eos(streamWritable, (error) => { @@ -608,6 +628,7 @@ function newWritableStreamFromStreamWritable( start(c) { controller = c; abortSignal = getWritableStreamDefaultControllerSignal(c); + listenForPendingAbort(); }, write(chunk) { @@ -615,17 +636,22 @@ function newWritableStreamFromStreamWritable( let needDrainBeforeWrite; let waitForWriteCallback; try { - options[kValidateChunk]?.(chunk); + if (hasValidationOptions) { + options[kValidateChunk]?.(chunk); + } writeMethod = streamWritable.write; const writableState = streamWritable._writableState; const objectMode = streamWritable.writableObjectMode; + let isArrayBufferViewChunk = isArrayBufferView(chunk); + const isArrayBufferChunk = + !isArrayBufferViewChunk && isAnyArrayBuffer(chunk); const isBufferSourceChunk = - isAnyArrayBuffer(chunk) || isArrayBufferView(chunk); - if (!objectMode && isAnyArrayBuffer(chunk)) { + isArrayBufferChunk || isArrayBufferViewChunk; + if (!objectMode && isArrayBufferChunk) { chunk = new Uint8Array(chunk); + isArrayBufferViewChunk = true; } needDrainBeforeWrite = streamWritable.writableNeedDrain; - const isArrayBufferViewChunk = isArrayBufferView(chunk); const canAwaitWriteCallback = streamWritable instanceof Writable && writeMethod === writablePrototypeWrite && diff --git a/lib/internal/webstreams/writablestream.js b/lib/internal/webstreams/writablestream.js index 0ebce8dd60ae..d9667f01d406 100644 --- a/lib/internal/webstreams/writablestream.js +++ b/lib/internal/webstreams/writablestream.js @@ -28,6 +28,10 @@ const { DOMException, } = internalBinding('messaging'); +const { + enqueueMicrotask, +} = internalBinding('task_queue'); + const { customInspectSymbol: kInspect, kEmptyObject, @@ -96,9 +100,11 @@ const { } = require('internal/process/task_queues'); const assert = require('internal/assert'); +const AsyncContextFrame = require('internal/async_context_frame'); const kAbort = Symbol('kAbort'); const kCloseSentinel = Symbol('kCloseSentinel'); + const kError = Symbol('kError'); const kSkipThrow = Symbol('kSkipThrow'); @@ -1166,6 +1172,49 @@ function writableStreamDefaultControllerWrite(controller, chunk, chunkSize) { writableStreamDefaultControllerAdvanceQueueIfNeeded(controller); } +// Trusted adapters can resume the parked write at the same microtask position +// as a pending sink promise, without allocating a promise capability per write. +function runWriteReaction(reaction, error) { + try { + reaction(error); + } catch (error) { + // Match an unhandled rejection from a throwing promise reaction rather + // than reporting the exception as an uncaught microtask exception. + PromiseReject(error); + } +} + +function writableStreamDefaultControllerResolveWrite(controller, contextFrame) { + const state = controller[kState]; + if (state.writeFulfilledTask === undefined) { + const reaction = state.writeFulfilled; + state.writeFulfilledTask = () => runWriteReaction(reaction); + } + const prior = AsyncContextFrame.current(); + if (prior === contextFrame) { + enqueueMicrotask(state.writeFulfilledTask); + return; + } + AsyncContextFrame.set(contextFrame); + try { + // The native microtask captures the frame without running promise init + // hooks or looking up a user-modified Promise constructor/species. + enqueueMicrotask(state.writeFulfilledTask); + } finally { + AsyncContextFrame.set(prior); + } +} + +function writableStreamDefaultControllerRejectWrite(controller, contextFrame, error) { + const prior = AsyncContextFrame.exchange(contextFrame); + try { + const reaction = controller[kState].writeRejected; + enqueueMicrotask(() => runWriteReaction(reaction, error)); + } finally { + AsyncContextFrame.set(prior); + } +} + function writableStreamDefaultControllerProcessWrite(controller, chunk) { const { stream, @@ -1376,6 +1425,7 @@ function setupWritableStreamDefaultController( stream, writeAlgorithm, writeFulfilled: undefined, + writeFulfilledTask: undefined, writeRejected: undefined, }; stream[kState].controller = controller; @@ -1425,6 +1475,8 @@ module.exports = { isWritableStreamDefaultWriter, getWritableStreamDefaultControllerSignal, + writableStreamDefaultControllerResolveWrite, + writableStreamDefaultControllerRejectWrite, isWritableStreamLocked, setupWritableStreamDefaultWriter, writableStreamAbort, diff --git a/test/parallel/test-webstreams-adapters-writable-async-context.js b/test/parallel/test-webstreams-adapters-writable-async-context.js new file mode 100644 index 000000000000..6066af60c85c --- /dev/null +++ b/test/parallel/test-webstreams-adapters-writable-async-context.js @@ -0,0 +1,174 @@ +'use strict'; + +const common = require('../common'); +const assert = require('assert'); +const { AsyncLocalStorage, AsyncResource, createHook } = require('async_hooks'); +const { Writable } = require('stream'); + +const mode = process.argv[2]; +const hook = mode === 'hooks' ? createHook({ init() {} }).enable() : undefined; + +async function checkContext(origin, highWaterMark) { + const storage = new AsyncLocalStorage(); + const started = [Promise.withResolvers(), Promise.withResolvers()]; + const callbacks = []; + const writable = new Writable({ + highWaterMark, + write: common.mustCall((chunk, encoding, callback) => { + // A completion callback from another request must not change the + // context in which the adapter invokes the next queued native write. + assert.strictEqual(storage.getStore(), origin); + const index = chunk[0] - 1; + callbacks[index] = callback; + started[index].resolve(); + }, 2), + }); + const writer = storage.run(origin, () => Writable.toWeb(writable).getWriter()); + const writes = []; + for (const id of [1, 2]) { + writes.push(storage.run(`writer-${id}`, common.mustCall(() => + writer.write(Buffer.from([id])).then(common.mustCall(() => { + assert.strictEqual(storage.getStore(), `writer-${id}`); + }))))); + } + + for (let index = 0; index < writes.length; index++) { + await started[index].promise; + storage.run(`callback-${index}`, common.mustCall(() => { + callbacks[index](); + assert.strictEqual(storage.getStore(), `callback-${index}`); + })); + await writes[index]; + } + await writer.close(); + storage.disable(); +} + +async function checkMicrotaskOrder(highWaterMark) { + const events = []; + const started = [Promise.withResolvers(), Promise.withResolvers()]; + const callbacks = []; + const writable = new Writable({ + highWaterMark, + write: common.mustCall((chunk, encoding, callback) => { + const index = chunk[0] - 1; + events.push(`write-${index + 1}`); + callbacks[index] = callback; + started[index].resolve(); + }, 2), + }); + const writer = Writable.toWeb(writable).getWriter(); + const first = writer.write(Buffer.from([1])).then(() => events.push('fulfilled-1')); + const second = writer.write(Buffer.from([2])); + await started[0].promise; + queueMicrotask(() => events.push('before-callback')); + callbacks[0](); + queueMicrotask(() => events.push('after-callback')); + await first; + assert.deepStrictEqual(events, [ + 'write-1', 'before-callback', 'write-2', 'after-callback', 'fulfilled-1', + ]); + await started[1].promise; + callbacks[1](); + await second; + await writer.close(); +} + +async function checkLateHook(sameContext, fail) { + const storage = new AsyncLocalStorage(); + const started = [Promise.withResolvers(), Promise.withResolvers()]; + const callbacks = []; + const error = fail ? new Error('late hook write failure') : undefined; + const writable = new Writable({ + write: common.mustCall((chunk, encoding, callback) => { + assert.strictEqual(storage.getStore(), 'stream-origin'); + const index = chunk[0] - 1; + callbacks[index] = callback; + started[index].resolve(); + }, fail ? 1 : 2), + }); + let writer; + let resource; + storage.run('stream-origin', () => { + resource = new AsyncResource('WritableCompletion'); + writer = Writable.toWeb(writable).getWriter(); + }); + const first = writer.write(Buffer.from([1])); + const completion = fail ? assert.rejects(first, (actual) => actual === error) : first; + const second = fail ? undefined : writer.write(Buffer.from([2])); + const closed = fail ? assert.rejects(writer.closed, (actual) => actual === error) : undefined; + await started[0].promise; + + let armed = false; + const lateHook = createHook({ + init(id, type) { + if (armed && type === 'PROMISE') { + storage.enterWith('hook-origin'); + } + }, + }).enable(); + + function finish(index) { + const invoke = common.mustCall(() => { + armed = true; + try { + callbacks[index](error); + } finally { + armed = false; + } + assert.strictEqual(storage.getStore(), sameContext ? 'stream-origin' : 'callback-origin'); + }); + if (sameContext) { + resource.runInAsyncScope(invoke); + } else { + storage.run('callback-origin', invoke); + } + } + + try { + finish(0); + await completion; + if (fail) { + await closed; + } else { + await started[1].promise; + finish(1); + await second; + await writer.close(); + } + } finally { + lateHook.disable(); + resource.emitDestroy(); + storage.disable(); + } +} + +async function main() { + // Keep this test outside node:test so its async hooks do not hide the + // adapter's default completion path. + for (const highWaterMark of [1, 65536]) { + for (const origin of [undefined, 'stream-origin']) { + await checkContext(origin, highWaterMark); + } + await checkMicrotaskOrder(highWaterMark); + } + for (const sameContext of [false, true]) { + for (const fail of [false, true]) { + await checkLateHook(sameContext, fail); + } + } + hook?.disable(); + + if (mode === undefined) { + for (const args of [ + [__filename, 'hooks'], + ['--no-async-context-frame', __filename, 'legacy'], + ]) { + const { code, signal, stderr } = await common.spawnPromisified(process.execPath, args); + assert.strictEqual(code, 0, stderr); + assert.strictEqual(signal, null); + } + } +} + +main().then(common.mustCall()); diff --git a/test/parallel/test-webstreams-adapters-writable-lifecycle.js b/test/parallel/test-webstreams-adapters-writable-lifecycle.js new file mode 100644 index 000000000000..b695440150bb --- /dev/null +++ b/test/parallel/test-webstreams-adapters-writable-lifecycle.js @@ -0,0 +1,196 @@ +'use strict'; + +const common = require('../common'); +const assert = require('assert'); +const { createHook } = require('async_hooks'); +const { Writable } = require('stream'); + +const mode = process.argv[2]; +const hook = mode === 'hooks' ? createHook({ init() {} }).enable() : undefined; + +async function checkBufferedWritev() { + const consumed = []; + let finishInitial; + let finishBuffered; + const writable = new Writable({ + highWaterMark: 1, + write: common.mustCall((chunk, encoding, callback) => { + consumed.push(Buffer.from(chunk)); + finishInitial = callback; + }), + writev: common.mustCall((chunks, callback) => { + consumed.push(...chunks.map(({ chunk }) => Buffer.from(chunk))); + finishBuffered = callback; + }), + }); + writable.write(Buffer.from('initial')); + writable.write(Buffer.from('external')); + const writer = Writable.toWeb(writable).getWriter(); + const input = Buffer.from('web'); + let fulfilled = false; + const pending = writer.write(input).then(common.mustCall(() => { fulfilled = true; })); + await new Promise(setImmediate); + assert.strictEqual(fulfilled, false); + finishInitial(); + assert.strictEqual(typeof finishBuffered, 'function'); + await new Promise(setImmediate); + assert.strictEqual(fulfilled, false); + finishBuffered(); + await pending; + input.fill(255); + await writer.close(); + assert.deepStrictEqual(consumed, [ + Buffer.from('initial'), Buffer.from('external'), Buffer.from('web'), + ]); +} + +async function checkPendingAbort(reason, highWaterMark) { + const started = Promise.withResolvers(); + let finishWrite; + const writable = new Writable({ + highWaterMark, + write: common.mustCall((chunk, encoding, callback) => { + finishWrite = callback; + started.resolve(); + }), + }); + const writer = Writable.toWeb(writable).getWriter(); + const pending = assert.rejects(writer.write(Buffer.from('x')), (actual) => actual === reason); + const closed = assert.rejects(writer.closed, (actual) => actual === reason); + await started.promise; + await Promise.all([pending, closed, writer.abort(reason)]); + assert.strictEqual(writable.destroyed, true); + // A callback arriving after abort must not complete another Web write. + finishWrite(); + await new Promise(setImmediate); +} + +async function checkAbortDuringDrain() { + const started = Promise.withResolvers(); + const error = new Error('abort during drain'); + let finishWrite; + let abort; + const writable = new Writable({ + highWaterMark: 1, + write: common.mustCall((chunk, encoding, callback) => { + finishWrite = callback; + started.resolve(); + }), + }); + const writer = Writable.toWeb(writable).getWriter(); + writable.on('drain', common.mustCall(() => { abort = writer.abort(error); })); + const pending = assert.rejects(writer.write(Buffer.from('x')), (actual) => actual === error); + const closed = assert.rejects(writer.closed, (actual) => actual === error); + await started.promise; + finishWrite(); + await Promise.all([abort, pending, closed]); +} + +async function checkReentrantWrite() { + const consumed = []; + const writable = new Writable({ + write: common.mustCall(function(chunk, encoding, callback) { + const consume = () => { + consumed.push(chunk[0]); + callback(); + }; + if (chunk[0] === 2) { + setImmediate(consume); + } else { + consume(); + if (chunk[0] === 1) this.write(Buffer.from([2]), common.mustCall()); + } + }, 3), + }); + const writer = Writable.toWeb(writable).getWriter(); + const input = Buffer.from([1]); + await writer.write(input); + input[0] = 3; + await writer.write(input); + input[0] = 9; + await writer.close(); + assert.deepStrictEqual(consumed, [1, 2, 3]); +} + +async function checkConstruction(fail, synchronous = false) { + const error = fail ? new Error('construction failed') : undefined; + let finishConstruct; + let consumed = false; + const writable = new Writable({ + construct: common.mustCall((callback) => { finishConstruct = callback; }), + write: fail ? common.mustNotCall() : common.mustCall((chunk, encoding, callback) => { + consumed = true; + assert.deepStrictEqual(chunk, Buffer.from('x')); + if (synchronous) callback(); + else setImmediate(callback); + }), + }); + const writer = Writable.toWeb(writable).getWriter(); + const write = writer.write(Buffer.from('x')); + const pending = fail ? assert.rejects(write, (actual) => actual === error) : write; + const closed = fail ? assert.rejects(writer.closed, (actual) => actual === error) : undefined; + await new Promise(setImmediate); + assert.strictEqual(consumed, false); + finishConstruct(error); + await pending; + assert.strictEqual(consumed, !fail); + if (fail) { + await closed; + } else { + await writer.close(); + } +} + +async function checkPrematureClose(error, highWaterMark) { + const started = Promise.withResolvers(); + let finishWrite; + const writable = new Writable({ + highWaterMark, + write: common.mustCall((chunk, encoding, callback) => { + finishWrite = callback; + started.resolve(); + }), + final: common.mustNotCall(), + }); + const expected = error === undefined ? { code: 'ABORT_ERR' } : (actual) => actual === error; + const writer = Writable.toWeb(writable).getWriter(); + const pending = assert.rejects(writer.write(Buffer.from('x')), expected); + const closed = assert.rejects(writer.closed, expected); + await started.promise; + writable.destroy(error); + const close = assert.rejects(writer.close(), expected); + await Promise.all([pending, closed, close]); + finishWrite(); + await new Promise(setImmediate); +} + +async function main() { + // node:test enables async hooks and would hide the default parked path. + await checkBufferedWritev(); + await checkAbortDuringDrain(); + await checkReentrantWrite(); + for (const highWaterMark of [1, 65536]) { + for (const reason of [null, false, 0, '', new Error('abort')]) { + await checkPendingAbort(reason, highWaterMark); + } + for (const error of [undefined, new Error('premature native close')]) { + await checkPrematureClose(error, highWaterMark); + } + } + for (const fail of [false, true]) await checkConstruction(fail); + await checkConstruction(false, true); + hook?.disable(); + + if (mode === undefined) { + for (const args of [ + [__filename, 'hooks'], + ['--no-async-context-frame', __filename, 'legacy'], + ]) { + const { code, signal, stderr } = await common.spawnPromisified(process.execPath, args); + assert.strictEqual(code, 0, stderr); + assert.strictEqual(signal, null); + } + } +} + +main().then(common.mustCall()); diff --git a/test/parallel/test-webstreams-adapters-writable-promise-mutation.js b/test/parallel/test-webstreams-adapters-writable-promise-mutation.js new file mode 100644 index 000000000000..a3310e104cd8 --- /dev/null +++ b/test/parallel/test-webstreams-adapters-writable-promise-mutation.js @@ -0,0 +1,99 @@ +'use strict'; + +const common = require('../common'); +const assert = require('assert'); +const { Writable } = require('stream'); + +async function checkCompletion(target, key, fail) { + const started = Promise.withResolvers(); + let finishWrite; + const error = fail ? new Error('native write failure') : undefined; + const writable = new Writable({ + write: common.mustCall((chunk, encoding, callback) => { + finishWrite = callback; + started.resolve(); + }), + }); + const writer = Writable.toWeb(writable).getWriter(); + const write = writer.write(Buffer.from([1])); + const completion = fail ? assert.rejects(write, (actual) => actual === error) : write; + const closed = fail ? assert.rejects(writer.closed, (actual) => actual === error) : undefined; + await started.promise; + + const descriptor = Object.getOwnPropertyDescriptor(target, key); + Object.defineProperty(target, key, { + __proto__: null, + configurable: true, + get: common.mustNotCall('completion must not look up the Promise constructor or species'), + }); + try { + // Registering the completion reaction later must not introduce a new + // user-code call or exception inside the native write callback. + finishWrite(error); + } finally { + Object.defineProperty(target, key, descriptor); + } + + await completion; + if (fail) { + await closed; + } else { + await writer.close(); + } +} + +async function checkReactionError() { + const started = Promise.withResolvers(); + const observed = Promise.withResolvers(); + const error = new Error('write reaction failure'); + const descriptor = Object.getOwnPropertyDescriptor(Promise.prototype, 'constructor'); + let finishWrite; + const writable = new Writable({ + write: common.mustCall((chunk, encoding, callback) => { + if (chunk[0] === 1) { + finishWrite = callback; + started.resolve(); + return; + } + // Force the next synchronous sink reaction to throw while the first + // asynchronous write's fulfillment reaction advances the queue. + Object.defineProperty(Promise.prototype, 'constructor', { + __proto__: null, + configurable: true, + get() { throw error; }, + }); + callback(); + }, 2), + }); + const writer = Writable.toWeb(writable).getWriter(); + const first = writer.write(Buffer.from([1])); + writer.write(Buffer.from([2])).catch(() => {}); + writer.closed.catch(() => {}); + const onUncaught = common.mustNotCall('a throwing write reaction must reject a promise'); + process.once('uncaughtException', onUncaught); + process.once('unhandledRejection', common.mustCall((actual) => { + Object.defineProperty(Promise.prototype, 'constructor', descriptor); + assert.strictEqual(actual, error); + writable.destroy(); + observed.resolve(); + })); + await started.promise; + finishWrite(); + await observed.promise; + await first; + process.removeListener('uncaughtException', onUncaught); +} + +async function main() { + for (const [target, key] of [ + [Promise.prototype, 'constructor'], + [Promise, Symbol.species], + ]) { + for (const fail of [false, true]) { + await checkCompletion(target, key, fail); + } + } + await checkReactionError(); +} + +main().then(common.mustCall()); diff --git a/test/parallel/test-webstreams-adapters-writable-prototype-pollution.js b/test/parallel/test-webstreams-adapters-writable-prototype-pollution.js new file mode 100644 index 000000000000..dd1092785877 --- /dev/null +++ b/test/parallel/test-webstreams-adapters-writable-prototype-pollution.js @@ -0,0 +1,50 @@ +'use strict'; +const common = require('../common'); +const assert = require('assert'); +const { Writable } = require('stream'); + +(async () => { + const consumed = []; + const writable = new Writable({ + highWaterMark: 1, + write: common.mustCall((chunk, encoding, callback) => { + queueMicrotask(() => { + consumed.push(chunk[0]); + callback(); + }); + }, 2), + }); + writable.on('error', common.mustNotCall()); + const writer = Writable.toWeb(writable).getWriter(); + const fields = ['deferred', 'writeComplete', 'backpressureCleared', 'then']; + const descriptors = fields.map((field) => + Object.getOwnPropertyDescriptor(Object.prototype, field)); + const input = Buffer.alloc(1); + + try { + for (const field of fields) { + Object.defineProperty(Object.prototype, field, { + __proto__: null, + configurable: true, + get: common.mustNotCall(`unexpected inherited get: ${field}`), + set: common.mustNotCall(`unexpected inherited set: ${field}`), + }); + } + for (let i = 0; i < 2; i++) { + input[0] = i; + await writer.write(input); + input[0] = 255; + } + await writer.close(); + } finally { + fields.forEach((field, i) => { + delete Object.prototype[field]; + if (descriptors[i]) { + Object.defineProperty(Object.prototype, field, descriptors[i]); + } + }); + } + + assert.deepStrictEqual(consumed, Array.from({ length: 2 }, (_, i) => i)); + assert.strictEqual(writable.writableFinished, true); +})().then(common.mustCall()); diff --git a/test/parallel/test-webstreams-adapters-writable-sync-completion.js b/test/parallel/test-webstreams-adapters-writable-sync-completion.js new file mode 100644 index 000000000000..2a4852e2f227 --- /dev/null +++ b/test/parallel/test-webstreams-adapters-writable-sync-completion.js @@ -0,0 +1,196 @@ +'use strict'; + +const common = require('../common'); +const assert = require('assert'); +const { Writable } = require('stream'); +const test = require('node:test'); + +test('synchronous Buffer writes on an active stream complete without another tick', async () => { + const writable = new Writable({ + write: common.mustCall((chunk, encoding, callback) => callback(), 2), + }); + const writer = Writable.toWeb(writable).getWriter(); + await writer.write(Buffer.from('hello')); + + // Run in a Promise reaction so a nextTick cannot run until these + // microtasks finish unless the write itself waits for that tick. + await Promise.resolve(); + let ticked = false; + process.nextTick(() => { ticked = true; }); + await writer.write(Buffer.from('hello')); + assert.strictEqual(ticked, false); + await writer.close(); +}); + +for (const ctor of [ArrayBuffer, SharedArrayBuffer]) { + test(`synchronous Uint8Array writes preserve the ${ctor.name} view without another tick`, async () => { + const buffer = new ctor(8); + const input = new Uint8Array(buffer, 2, 4); + input.set([1, 2, 3, 4]); + const writable = new Writable({ + highWaterMark: 1, + write: common.mustCall((chunk, encoding, callback) => { + assert(Buffer.isBuffer(chunk)); + assert.strictEqual(chunk.buffer, buffer); + assert.strictEqual(chunk.byteOffset, 2); + assert.deepStrictEqual(chunk, Buffer.from([1, 2, 3, 4])); + callback(); + }, 2), + }); + writable.on('drain', common.mustNotCall()); + const writer = Writable.toWeb(writable).getWriter(); + await writer.write(input); + + await Promise.resolve(); + let ticked = false; + process.nextTick(() => { ticked = true; }); + await writer.write(input); + assert.strictEqual(ticked, false); + await writer.close(); + }); +} + +test('mixed synchronous and asynchronous writes preserve reusable bytes', async () => { + const consumed = []; + const writable = new Writable({ + write(chunk, encoding, callback) { + const consume = () => { + consumed.push(chunk[0]); + callback(); + }; + if (chunk[0] % 2 === 0) { + setImmediate(consume); + } else { + consume(); + } + }, + }); + const writer = Writable.toWeb(writable).getWriter(); + const input = Buffer.alloc(1); + for (let i = 1; i <= 6; i++) { + input[0] = i; + await writer.write(input); + input[0] = 9; + } + await writer.close(); + assert.deepStrictEqual(consumed, [1, 2, 3, 4, 5, 6]); +}); + +test('synchronous write callback errors reject the write and closed promises', async () => { + const error = new Error('sync callback error'); + const writable = new Writable({ + write(chunk, encoding, callback) { + callback(chunk[0] === 1 ? undefined : error); + }, + }); + const writer = Writable.toWeb(writable).getWriter(); + await writer.write(Buffer.from([1])); + await Promise.all([ + assert.rejects(writer.write(Buffer.from([2])), (actual) => actual === error), + assert.rejects(writer.closed, (actual) => actual === error), + ]); +}); + +test('falsy native callback errors still complete the write', async () => { + const errors = [undefined, null, false, 0, '']; + const writable = new Writable({ + write: common.mustCall((chunk, encoding, callback) => { + callback(errors[chunk[0]]); + }, errors.length), + }); + const writer = Writable.toWeb(writable).getWriter(); + for (let i = 0; i < errors.length; i++) { + await writer.write(Buffer.from([i])); + } + await writer.close(); +}); + +test('abort releases an asynchronous write after synchronous writes', async () => { + let finishWrite; + let startWrite; + const started = new Promise((resolve) => { startWrite = resolve; }); + const error = new Error('abort after sync writes'); + const writable = new Writable({ + write(chunk, encoding, callback) { + if (chunk[0] === 1) { + callback(); + } else { + finishWrite = callback; + startWrite(); + } + }, + }); + const writer = Writable.toWeb(writable).getWriter(); + for (let i = 0; i < 2; i++) { + await writer.write(Buffer.from([1])); + } + const writePromise = writer.write(Buffer.from([2])); + await started; + await Promise.all([ + assert.rejects(writePromise, (actual) => actual === error), + assert.rejects(writer.closed, (actual) => actual === error), + writer.abort(error), + ]); + finishWrite(); + await new Promise(setImmediate); + assert.strictEqual(writable.destroyed, true); +}); + +test('abort during a synchronous write rejects its completion', async () => { + const error = new Error('abort during write'); + let writer; + let abortPromise; + const writable = new Writable({ + write: common.mustCall((chunk, encoding, callback) => { + if (chunk[0] === 2) { + abortPromise = writer.abort(error); + } + callback(); + }, 2), + }); + writer = Writable.toWeb(writable).getWriter(); + await writer.write(Buffer.from([1])); + await Promise.all([ + assert.rejects(writer.write(Buffer.from([2])), (actual) => actual === error), + assert.rejects(writer.closed, (actual) => actual === error), + ]); + await abortPromise; + assert.strictEqual(writable.destroyed, true); +}); + +for (const highWaterMark of [0, 1]) { + test(`synchronous and empty writes with highWaterMark ${highWaterMark}`, async () => { + const consumed = []; + const writable = new Writable({ + highWaterMark, + write(chunk, encoding, callback) { + consumed.push(Buffer.from(chunk)); + callback(); + }, + }); + const writer = Writable.toWeb(writable).getWriter(); + for (const input of [Buffer.from([1]), Buffer.alloc(0), Buffer.from([2])]) { + await writer.write(input); + input.fill(9); + assert.strictEqual(writable.writableNeedDrain, false); + } + await writer.close(); + assert.deepStrictEqual(consumed, [Buffer.from([1]), Buffer.alloc(0), Buffer.from([2])]); + }); +} + +test('asynchronous writes still wait for drain', async () => { + const writable = new Writable({ + highWaterMark: 1, + write(chunk, encoding, callback) { + setImmediate(callback); + }, + }); + writable.on('drain', common.mustCall(3)); + const writer = Writable.toWeb(writable).getWriter(); + for (let i = 0; i < 3; i++) { + await writer.write(Buffer.from([1])); + assert.strictEqual(writable.writableNeedDrain, false); + } + await writer.close(); +}); From 71da7b5a134a91e808d83401377d3bc9e5ae8c90 Mon Sep 17 00:00:00 2001 From: seungwoo Date: Thu, 1 Oct 2026 20:28:14 +0900 Subject: [PATCH 08/15] stream: preserve native HTTP write validation Preserve the original HTTP write type error when cloning an invalid DataView fails. Only restore validation for the unmodified native write method, leaving overridden methods on the copy fallback. Cover resizing before and after enqueue and changes to HTTP write methods. Signed-off-by: seungwoo --- lib/_http_outgoing.js | 3 + lib/internal/webstreams/adapters.js | 17 +++- ...bstreams-adapters-http-write-validation.js | 86 +++++++++++++++++++ 3 files changed, 105 insertions(+), 1 deletion(-) create mode 100644 test/parallel/test-webstreams-adapters-http-write-validation.js diff --git a/lib/_http_outgoing.js b/lib/_http_outgoing.js index 2ca455eb5566..817e4ddf49f2 100644 --- a/lib/_http_outgoing.js +++ b/lib/_http_outgoing.js @@ -1038,6 +1038,8 @@ OutgoingMessage.prototype.write = function write(chunk, encoding, callback) { return ret; }; +const outgoingMessagePrototypeWrite = OutgoingMessage.prototype.write; + function onError(msg, err, callback) { if (msg.destroyed) { return; @@ -1488,4 +1490,5 @@ module.exports = { validateHeaderName, validateHeaderValue, OutgoingMessage, + outgoingMessagePrototypeWrite, }; diff --git a/lib/internal/webstreams/adapters.js b/lib/internal/webstreams/adapters.js index b0cde3fd965e..dd99eb2aade3 100644 --- a/lib/internal/webstreams/adapters.js +++ b/lib/internal/webstreams/adapters.js @@ -42,6 +42,7 @@ const AsyncContextFrame = require('internal/async_context_frame'); const { enabledHooksExist } = require('internal/async_hooks'); let addAbortListener; +let outgoingMessagePrototypeWrite; const { ReadableStream, @@ -668,7 +669,21 @@ function newWritableStreamFromStreamWritable( // override write() may not support callbacks at all. Preserve the // existing completion timing for those cases by passing a private, // same-type copy to the native stream. - chunk = cloneArrayBufferView(chunk); + try { + chunk = cloneArrayBufferView(chunk); + } catch (error) { + if (isDataView(chunk)) { + outgoingMessagePrototypeWrite ??= + require('_http_outgoing').outgoingMessagePrototypeWrite; + if (writeMethod === outgoingMessagePrototypeWrite) { + // HTTP rejects DataView before inspecting its backing. Keep + // that native validation error when cloning an invalid view + // fails, without invoking an overridden write() method. + FunctionPrototypeCall(writeMethod, streamWritable, chunk); + } + } + throw error; + } } } catch (error) { throwSyncWriteError(error); diff --git a/test/parallel/test-webstreams-adapters-http-write-validation.js b/test/parallel/test-webstreams-adapters-http-write-validation.js new file mode 100644 index 000000000000..1499c9616d3b --- /dev/null +++ b/test/parallel/test-webstreams-adapters-http-write-validation.js @@ -0,0 +1,86 @@ +'use strict'; + +const common = require('../common'); +const assert = require('assert'); +const { IncomingMessage, OutgoingMessage, ServerResponse } = require('http'); +const { Writable } = require('stream'); + +function createResponse() { + const socket = new Writable({ + write(chunk, encoding, callback) { callback(); }, + }); + const request = new IncomingMessage(socket); + request.method = 'GET'; + request.httpVersionMajor = 1; + request.httpVersionMinor = 1; + const response = new ServerResponse(request); + response.assignSocket(socket); + response.on('error', common.mustNotCall()); + return response; +} + +async function checkNativeValidation(resizeAfterEnqueue) { + const backing = new ArrayBuffer(32, { maxByteLength: 64 }); + const input = new DataView(backing, 16, 8); + if (!resizeAfterEnqueue) { + backing.resize(8); + // Let the queuing strategy admit the view so the native type validation + // is exercised even though the view's intrinsic metadata is out of bounds. + Object.defineProperty(input, 'byteLength', { value: 8 }); + } + const writer = Writable.toWeb(createResponse()).getWriter(); + const pending = writer.write(input); + if (resizeAfterEnqueue) backing.resize(8); + await Promise.all([ + assert.rejects(pending, { code: 'ERR_INVALID_ARG_TYPE' }), + assert.rejects(writer.closed, { code: 'ERR_INVALID_ARG_TYPE' }), + ]); +} + +async function checkOverriddenWrite() { + const response = createResponse(); + const writer = Writable.toWeb(response).getWriter(); + let retained; + response.write = common.mustCall((chunk) => { + retained = chunk; + return true; + }); + const input = new DataView(new ArrayBuffer(4)); + input.setUint32(0, 0x01020304); + await writer.write(input); + input.setUint32(0, 0x09090909); + assert(retained instanceof DataView); + assert.notStrictEqual(retained.buffer, input.buffer); + assert.strictEqual(retained.getUint32(0), 0x01020304); + // A real HTTP server closes the response after finish. Let the test socket + // deliver that close event so the adapter can observe the completed stream. + response.once('finish', common.mustCall(() => response.socket.destroy())); + await writer.close(); +} + +async function checkPatchedWrite() { + const originalWrite = OutgoingMessage.prototype.write; + OutgoingMessage.prototype.write = common.mustNotCall(); + try { + const backing = new ArrayBuffer(32, { maxByteLength: 64 }); + const input = new DataView(backing, 16, 8); + backing.resize(8); + Object.defineProperty(input, 'byteLength', { value: 8 }); + const writer = Writable.toWeb(createResponse()).getWriter(); + const expected = common.mustCall((error) => error instanceof TypeError && error.code === undefined, 2); + await Promise.all([ + assert.rejects(writer.write(input), expected), + assert.rejects(writer.closed, expected), + ]); + } finally { + OutgoingMessage.prototype.write = originalWrite; + } +} + +async function main() { + for (const resizeAfterEnqueue of [false, true]) await checkNativeValidation(resizeAfterEnqueue); + await checkOverriddenWrite(); + await checkPatchedWrite(); +} + +main().then(common.mustCall()); From dba8982c6531b72aef303228b88d1ab18c2be614 Mon Sep 17 00:00:00 2001 From: seungwoo Date: Thu, 1 Oct 2026 20:28:17 +0900 Subject: [PATCH 09/15] stream: preserve errors on premature writable close Reject close() when the native stream was destroyed before finishing, so end() cannot mask a premature close before eos reports it. Preserve the native error when available and cover synchronous and asynchronous destruction, legacy streams, and normal autoDestroy. Signed-off-by: seungwoo --- lib/internal/webstreams/adapters.js | 11 ++ ...reams-adapters-writable-premature-close.js | 109 ++++++++++++++++++ 2 files changed, 120 insertions(+) create mode 100644 test/parallel/test-webstreams-adapters-writable-premature-close.js diff --git a/lib/internal/webstreams/adapters.js b/lib/internal/webstreams/adapters.js index dd99eb2aade3..2bfc2e0d6ac7 100644 --- a/lib/internal/webstreams/adapters.js +++ b/lib/internal/webstreams/adapters.js @@ -79,6 +79,8 @@ const { isReadable, isWritable, isWritableEnded, + isWritableErrored, + isWritableFinished, } = require('internal/streams/utils'); const { @@ -703,6 +705,15 @@ function newWritableStreamFromStreamWritable( }, close() { + // end() can mark a destroyed stream as ended before eos observes its + // premature close. Preserve that failure instead of masking it. + if (isDestroyed(streamWritable) && !isWritableFinished(streamWritable)) { + disposeAbortListener(); + const error = isWritableErrored(streamWritable); + throw handleKnownInternalErrors( + (typeof error !== 'boolean' && error) || new ERR_STREAM_PREMATURE_CLOSE(), + ); + } if (closed === undefined && !isWritableEnded(streamWritable)) { closed = PromiseWithResolvers(); try { diff --git a/test/parallel/test-webstreams-adapters-writable-premature-close.js b/test/parallel/test-webstreams-adapters-writable-premature-close.js new file mode 100644 index 000000000000..14fea1b80659 --- /dev/null +++ b/test/parallel/test-webstreams-adapters-writable-premature-close.js @@ -0,0 +1,109 @@ +'use strict'; + +const common = require('../common'); +const assert = require('assert'); +const { EventEmitter } = require('events'); +const { Writable } = require('stream'); +const test = require('node:test'); + +for (const highWaterMark of [0, 1]) { + for (const asynchronousDestroy of [false, true]) { + for (const error of [undefined, new Error('premature native close')]) { + test(`closing after premature destruction with highWaterMark ${highWaterMark}, ` + + `asynchronousDestroy ${asynchronousDestroy}, error ${error !== undefined}`, async () => { + const writable = new Writable({ + highWaterMark, + write: common.mustCall(function(chunk, encoding, callback) { + if (chunk[0] === 2) { + this.destroy(error); + } + callback(); + }, 2), + final: common.mustNotCall(), + destroy: common.mustCall((reason, callback) => { + if (asynchronousDestroy) { + setImmediate(callback, reason); + } else { + callback(reason); + } + }), + }); + writable.on('finish', common.mustNotCall()); + writable.on('close', common.mustCall()); + writable.on('error', error === undefined ? + common.mustNotCall() : + common.mustCall((actual) => assert.strictEqual(actual, error))); + + const writer = Writable.toWeb(writable).getWriter(); + await writer.write(Buffer.from([1])); + const expectedError = error === undefined ? + { code: 'ABORT_ERR' } : + common.mustCall((actual) => actual === error, 2); + const closed = assert.rejects(writer.closed, expectedError); + await writer.write(Buffer.from([2])); + assert.strictEqual(writable.destroyed, true); + assert.strictEqual(writable.writableFinished, false); + + // Close before the native close event can report premature destruction. + await Promise.all([ + assert.rejects(writer.close(), expectedError), + closed, + ]); + }); + } + } +} + +for (const error of [false, true]) { + test(`closing a prematurely destroyed legacy stream ignores boolean error ${error}`, async () => { + const writable = new EventEmitter(); + writable.writable = true; + writable.writableHighWaterMark = 16; + writable._writableState = { errored: null, finished: false }; + writable.write = common.mustCall((chunk) => { + if (chunk[0] === 2) { + writable.destroyed = true; + writable.writable = false; + writable._writableState.errored = error; + process.nextTick(() => writable.emit('close')); + } + return true; + }, 2); + writable.end = common.mustNotCall(); + writable.on('close', common.mustCall()); + + const writer = Writable.toWeb(writable).getWriter(); + await writer.write(Buffer.from([1])); + const closed = assert.rejects(writer.closed, { code: 'ABORT_ERR' }); + await writer.write(Buffer.from([2])); + await Promise.all([ + assert.rejects(writer.close(), { code: 'ABORT_ERR' }), + closed, + ]); + }); +} + +test('closing a normally finished native stream accepts autoDestroy', async () => { + let writer; + let closePromise; + const writable = new Writable({ + autoDestroy: true, + write: common.mustCall((chunk, encoding, callback) => callback()), + final: common.mustCall((callback) => callback()), + destroy: common.mustCall(function(error, callback) { + assert.strictEqual(error, null); + assert.strictEqual(this.destroyed, true); + assert.strictEqual(this.writableFinished, true); + closePromise = writer.close(); + callback(error); + }), + }); + writable.on('finish', common.mustCall()); + writable.on('close', common.mustCall()); + writable.on('error', common.mustNotCall()); + writer = Writable.toWeb(writable).getWriter(); + await writer.write(Buffer.from([1])); + writable.end(); + await writer.closed; + await closePromise; +}); From 64a173d3ff8b90c6a7f922824baf834067c2546e Mon Sep 17 00:00:00 2001 From: seungwoo Date: Thu, 1 Oct 2026 20:28:20 +0900 Subject: [PATCH 10/15] test: deduplicate Web Stream adapter regressions Remove cases now covered by the dedicated copy and lifecycle tests. Keep compression coverage focused on completing a write before reading output and preserving bytes after the input is reused. Signed-off-by: seungwoo --- ...st-webstreams-adapters-sync-write-error.js | 51 ---------- ...treams-adapters-writable-buffer-sources.js | 97 ------------------- ...st-webstreams-compression-buffer-source.js | 26 +++-- 3 files changed, 16 insertions(+), 158 deletions(-) diff --git a/test/parallel/test-webstreams-adapters-sync-write-error.js b/test/parallel/test-webstreams-adapters-sync-write-error.js index ff6dbddf377f..39c970be7616 100644 --- a/test/parallel/test-webstreams-adapters-sync-write-error.js +++ b/test/parallel/test-webstreams-adapters-sync-write-error.js @@ -56,57 +56,6 @@ test('WritableStream from Node.js stream handles async write error', async () => ]); }); -test('WritableStream aborts while a native write is pending', async () => { - const error = new Error('abort'); - let finishWrite; - let startWrite; - const writeStarted = new Promise((resolve) => { - startWrite = resolve; - }); - const writable = new Writable({ - write(_chunk, _encoding, callback) { - finishWrite = callback; - startWrite(); - }, - }); - const writer = Writable.toWeb(writable).getWriter(); - const writePromise = writer.write(Buffer.from('hello')); - await writeStarted; - - const writeRejected = assert.rejects(writePromise, isSameError(error)); - const closedRejected = assert.rejects(writer.closed, isSameError(error)); - await Promise.all([ - writer.abort(error), - writeRejected, - closedRejected, - ]); - - finishWrite(); - await new Promise(setImmediate); - assert.strictEqual(writable.destroyed, true); -}); - -test('WritableStream handles destruction while a write is pending', async () => { - const error = new Error('destroy'); - let startWrite; - const writeStarted = new Promise((resolve) => { - startWrite = resolve; - }); - const writable = new Writable({ - write(_chunk, _encoding, _callback) { - startWrite(); - }, - }); - const writer = Writable.toWeb(writable).getWriter(); - const writePromise = writer.write(Buffer.from('hello')); - await writeStarted; - - const writeRejected = assert.rejects(writePromise, isSameError(error)); - const closedRejected = assert.rejects(writer.closed, isSameError(error)); - writable.destroy(error); - await Promise.all([writeRejected, closedRejected]); -}); - test('Duplex-backed pair does NOT destroy on sync write throw', async () => { const error = new TypeError('invalid chunk'); const duplex = new Duplex({ diff --git a/test/parallel/test-webstreams-adapters-writable-buffer-sources.js b/test/parallel/test-webstreams-adapters-writable-buffer-sources.js index 7b6ed6a41026..45c5bfcd0c2a 100644 --- a/test/parallel/test-webstreams-adapters-writable-buffer-sources.js +++ b/test/parallel/test-webstreams-adapters-writable-buffer-sources.js @@ -191,78 +191,6 @@ suite('underlying Writable', () => { await writer.close(); }); - test('preserves cloned view brands and SharedArrayBuffer backing', - async () => { - const dataView = new DataView( - Uint8Array.from([0, 1, 2, 3, 4, 0]).buffer, - 1, - 4, - ); - const uint16Buffer = new ArrayBuffer(6); - const uint16 = new Uint16Array(uint16Buffer, 2, 2); - new Uint8Array(uint16Buffer, 2, 4).set([1, 2, 3, 4]); - const shared = new SharedArrayBuffer(6); - const sharedView = new Uint8Array(shared, 1, 4); - sharedView.set([1, 2, 3, 4]); - const inputs = [ - Buffer.from([1, 2, 3, 4]), - dataView, - uint16, - sharedView, - ]; - const expected = inputs.map((chunk) => ({ - brand: Buffer.isBuffer(chunk) ? - 'Buffer' : Object.prototype.toString.call(chunk), - bytes: Buffer.from(new Uint8Array( - chunk.buffer, - chunk.byteOffset, - chunk.byteLength, - )), - shared: chunk.buffer instanceof SharedArrayBuffer, - })); - const received = []; - const writable = new Writable({ - write(chunk, encoding, callback) { - callback(); - }, - }); - writable.on('error', common.mustNotCall()); - const writer = Writable.toWeb(writable).getWriter(); - writable.write = common.mustCall((chunk) => { - received.push(chunk); - return true; - }, inputs.length); - - for (const chunk of inputs) { - await writer.write(chunk); - new Uint8Array( - chunk.buffer, - chunk.byteOffset, - chunk.byteLength, - ).fill(9); - } - await writer.close(); - - for (let i = 0; i < received.length; i++) { - const actual = received[i]; - const actualBrand = Buffer.isBuffer(actual) ? - 'Buffer' : Object.prototype.toString.call(actual); - assert.strictEqual(actualBrand, expected[i].brand); - assert.notStrictEqual(actual.buffer, inputs[i].buffer); - assert.strictEqual( - actual.buffer instanceof SharedArrayBuffer, - expected[i].shared, - ); - assert.deepStrictEqual( - Buffer.from(new Uint8Array( - actual.buffer, - actual.byteOffset, - actual.byteLength, - )), - expected[i].bytes, - ); - } - }); }); suite('in object mode', () => { @@ -311,31 +239,6 @@ suite('underlying ServerResponse', () => { } }); - test('preserves write() overrides', async () => { - const response = createServerResponse(); - const writer = Writable.toWeb(response).getWriter(); - const originalWrite = response.write; - let received; - response.write = common.mustCall((chunk) => { - received = chunk; - return true; - }); - const input = new DataView(Uint8Array.from([1, 2, 3, 4]).buffer); - - try { - await writer.write(input); - new Uint8Array(input.buffer).fill(9); - assert(received instanceof DataView); - assert.notStrictEqual(received.buffer, input.buffer); - assert.deepStrictEqual( - Buffer.from(received.buffer), - Buffer.from([1, 2, 3, 4]), - ); - } finally { - response.write = originalWrite; - await new Promise((resolve) => response.end(resolve)); - } - }); }); suite('underlying Duplex', () => { diff --git a/test/parallel/test-webstreams-compression-buffer-source.js b/test/parallel/test-webstreams-compression-buffer-source.js index a81d4b674b38..22bff450b297 100644 --- a/test/parallel/test-webstreams-compression-buffer-source.js +++ b/test/parallel/test-webstreams-compression-buffer-source.js @@ -1,5 +1,5 @@ 'use strict'; -require('../common'); +const common = require('../common'); const assert = require('assert'); const test = require('node:test'); const { DecompressionStream, CompressionStream } = require('stream/web'); @@ -48,14 +48,20 @@ test('DecompressionStream writable completion is not coupled to readable ' + const compressed = gzipSync(expected); const ds = new DecompressionStream('gzip'); const writer = ds.writable.getWriter(); - let settled = false; - const writePromise = writer.write(compressed).then(() => { - settled = true; - }); - - await new Promise(setImmediate); - const settledBeforeRead = settled; - if (settledBeforeRead) { + const writePromise = writer.write(compressed); + let timer; + let completedBeforeRead; + try { + completedBeforeRead = await Promise.race([ + writePromise.then(() => true), + new Promise((resolve) => { + timer = setTimeout(() => resolve(false), common.platformTimeout(2000)); + }), + ]); + } finally { + clearTimeout(timer); + } + if (completedBeforeRead) { compressed.fill(0); } @@ -66,6 +72,6 @@ test('DecompressionStream writable completion is not coupled to readable ' + (await outputPromise).map((chunk) => Buffer.from(chunk)), ); - assert.strictEqual(settledBeforeRead, true); + assert.strictEqual(completedBeforeRead, true); assert.deepStrictEqual(output, expected); }); From 971077448df740d80e4ca383f87570af8dd234b3 Mon Sep 17 00:00:00 2001 From: seungwoo Date: Thu, 1 Oct 2026 20:28:22 +0900 Subject: [PATCH 11/15] doc: clarify mutable Web Stream chunk reuse Document when a mutable chunk can be reused after writer.write() fulfills, and prohibit mutation while the write is pending. Explain copying in byte mode, reference passing in object mode, and preservation of bytes retained by a Duplex readable side. Signed-off-by: seungwoo --- doc/api/stream.md | 15 +++++++++++++++ 1 file changed, 15 insertions(+) diff --git a/doc/api/stream.md b/doc/api/stream.md index 1e8a2bd1ef93..6d2cd0fe05d6 100644 --- a/doc/api/stream.md +++ b/doc/api/stream.md @@ -3339,6 +3339,14 @@ changes: * `streamWritable` {stream.Writable} * Returns: {WritableStream} +For streams not operating in object mode, a mutable {Buffer}, {TypedArray}, +{DataView}, {ArrayBuffer}, or {SharedArrayBuffer} passed to the returned stream's +writer can be reused after the promise returned by `writer.write()` is +fulfilled. Do not modify the chunk or its underlying bytes while that promise +is pending. The adapter may copy the chunk, so the underlying Node.js stream is +not guaranteed to receive the same object. Chunks written to streams operating +in object mode are passed by reference without copying. + ### `stream.Duplex.from(src)`