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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 24 additions & 0 deletions benchmark/webstreams/adapters.js
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ const bench = common.createBenchmark(main, {
'readable-to-web',
'readable-from-web',
'writable-to-web',
'writable-to-web-async-writev',
'writable-from-web',
],
});
Expand Down Expand Up @@ -68,6 +69,26 @@ async function writableToWeb(n) {
bench.end(n);
}

async function writableToWebAsyncWritev(n) {
const chunk = Buffer.alloc(1024);
const streamWritable = new Writable({
highWaterMark: 64 * 1024,
// Defer completion so subsequent writes can be batched in writev().
write(chunk, encoding, callback) {
setImmediate(callback);
},
writev(chunks, callback) {
setImmediate(callback);
},
});
const writer = Writable.toWeb(streamWritable).getWriter();
bench.start();
for (let i = 0; i < n; i++)
await writer.write(chunk);
await writer.close();
bench.end(n);
}

function writableFromWeb(n) {
const chunk = Buffer.alloc(1024);
const writableStream = new WritableStream({
Expand Down Expand Up @@ -99,6 +120,9 @@ function main({ n, kind }) {
case 'writable-to-web':
writableToWeb(n);
break;
case 'writable-to-web-async-writev':
writableToWebAsyncWritev(n);
break;
case 'writable-from-web':
writableFromWeb(n);
break;
Expand Down
15 changes: 15 additions & 0 deletions doc/api/stream.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)`

<!-- YAML
Expand Down Expand Up @@ -3507,6 +3515,12 @@ changes:
* `readable` {ReadableStream}
* `writable` {WritableStream}

The `writable` stream has the same mutable chunk reuse behavior as
[`stream.Writable.toWeb()`][]. For streams not operating in writable object mode,
reusing the chunk after `writer.write()` is fulfilled does not modify bytes
retained by the Node.js stream's readable side. The adapter may copy chunks
before passing them to the Node.js stream.

```mjs
import { Duplex } from 'node:stream';

Expand Down Expand Up @@ -5088,6 +5102,7 @@ contain multi-byte characters.
[`readable.push('')`]: #readablepush
[`readable.setEncoding()`]: #readablesetencodingencoding
[`stream.Readable.from()`]: #streamreadablefromiterable-options
[`stream.Writable.toWeb()`]: #streamwritabletowebstreamwritable
[`stream.addAbortSignal()`]: #streamaddabortsignalsignal-stream
[`stream.compose(...streams)`]: #streamcomposestreams
[`stream.cork()`]: #writablecork
Expand Down
3 changes: 3 additions & 0 deletions lib/_http_outgoing.js
Original file line number Diff line number Diff line change
Expand Up @@ -1017,6 +1017,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;
Expand Down Expand Up @@ -1467,4 +1469,5 @@ module.exports = {
validateHeaderName,
validateHeaderValue,
OutgoingMessage,
outgoingMessagePrototypeWrite,
};
45 changes: 30 additions & 15 deletions lib/internal/streams/writable.js
Original file line number Diff line number Diff line change
Expand Up @@ -449,7 +449,7 @@ Writable.prototype.pipe = function() {
errorOrDestroy(this, new ERR_STREAM_CANNOT_PIPE());
};

function _write(stream, chunk, encoding, cb) {
function _write(stream, chunk, encoding, cb, onwrite) {
const state = stream._writableState;

if (cb == null || typeof cb !== 'function') {
Expand Down Expand Up @@ -500,7 +500,7 @@ function _write(stream, chunk, encoding, cb) {
}

state.pendingcb++;
return writeOrBuffer(stream, state, chunk, encoding, cb);
return writeOrBuffer(stream, state, chunk, encoding, cb, onwrite);
}

Writable.prototype.write = function(chunk, encoding, cb) {
Expand All @@ -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;

Expand Down Expand Up @@ -547,7 +549,7 @@ Writable.prototype.setDefaultEncoding = function setDefaultEncoding(encoding) {
// If we're already writing something, then just put this
// in the queue, and wait our turn. Otherwise, call _write
// If we return false, then we need a drain event, so set that flag.
function writeOrBuffer(stream, state, chunk, encoding, callback) {
function writeOrBuffer(stream, state, chunk, encoding, callback, onwrite) {
const len = (state[kState] & kObjectMode) !== 0 ? 1 : chunk.length;

state.length += len;
Expand All @@ -558,7 +560,7 @@ function writeOrBuffer(stream, state, chunk, encoding, callback) {
state[kBufferedValue] = [];
}

state[kBufferedValue].push({ chunk, encoding, callback });
state[kBufferedValue].push({ chunk, encoding, callback, onwrite });
if ((state[kState] & kAllBuffers) !== 0 && encoding !== 'buffer') {
state[kState] &= ~kAllBuffers;
}
Expand All @@ -571,7 +573,7 @@ function writeOrBuffer(stream, state, chunk, encoding, callback) {
state.writecb = callback;
}
state[kState] |= kWriting | kSync | kExpectWriteCb;
stream._write(chunk, encoding, state.onwrite);
stream._write(chunk, encoding, onwrite ?? state.onwrite);
state[kState] &= ~kSync;
}

Expand All @@ -594,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;
Expand All @@ -608,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;
Expand All @@ -623,18 +625,25 @@ ObjectDefineProperty(Writable, 'writeKnownBuffer', {
enumerable: false,
});

function doWrite(stream, state, writev, len, chunk, encoding, cb) {
// Adapter-owned copies can be queued without waiting for their individual
// callbacks. Carry their completion context through native buffering so the
// next batch does not inherit another request's completion context.
function writeWithContext(stream, chunk, onwrite) {
return _write(stream, chunk, undefined, undefined, onwrite) === true;
}

function doWrite(stream, state, writev, len, chunk, encoding, cb, onwrite) {
state.writelen = len;
if (cb !== nop) {
state.writecb = cb;
}
state[kState] |= kWriting | kSync | kExpectWriteCb;
if ((state[kState] & kDestroyed) !== 0)
state.onwrite(new ERR_STREAM_DESTROYED('write'));
(onwrite ?? state.onwrite)(new ERR_STREAM_DESTROYED('write'));
else if (writev)
stream._writev(chunk, state.onwrite);
stream._writev(chunk, onwrite ?? state.onwrite);
else
stream._write(chunk, encoding, state.onwrite);
stream._write(chunk, encoding, onwrite ?? state.onwrite);
state[kState] &= ~kSync;
}

Expand Down Expand Up @@ -811,15 +820,15 @@ function clearBuffer(stream, state) {
buffered : ArrayPrototypeSlice(buffered, i);
chunks.allBuffers = (state[kState] & kAllBuffers) !== 0;

doWrite(stream, state, true, state.length, chunks, '', callback);
doWrite(stream, state, true, state.length, chunks, '', callback, buffered[i].onwrite);

resetBuffer(state);
} else {
do {
const { chunk, encoding, callback } = buffered[i];
const { chunk, encoding, callback, onwrite } = buffered[i];
buffered[i++] = null;
const len = objectMode ? 1 : chunk.length;
doWrite(stream, state, false, len, chunk, encoding, callback);
doWrite(stream, state, false, len, chunk, encoding, callback, onwrite);
} while (i < buffered.length && (state[kState] & kWriting) === 0);

if (i === buffered.length) {
Expand Down Expand Up @@ -1187,7 +1196,13 @@ Writable.fromWeb = function(writableStream, options) {
};

Writable.toWeb = function(streamWritable) {
return lazyWebStreams().newWritableStreamFromStreamWritable(streamWritable);
return lazyWebStreams().newWritableStreamFromStreamWritable(
streamWritable,
undefined,
writablePrototypeWrite,
writeKnownBuffer,
writeWithContext,
);
};

Writable.prototype[SymbolAsyncDispose] = async function() {
Expand Down
Loading