Commit 02a51a7

Browse files
trivikraduh95
authored andcommitted
stream: add sync iterable fast path to pipeTo
Avoid normalizing sync iterable sources through from() when pipeTo() has no transforms or signal and the writer can accept sync writes. This keeps writes incremental while preserving async fallback for values that still need it. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: openai:gpt-5.5 PR-URL: #63318 Backport-PR-URL: #64675 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Ethan Arrowood <ethan@arrowood.dev>
1 parent 4c9f3ed commit 02a51a7

3 files changed

Lines changed: 155 additions & 6 deletions

File tree

β€Žbenchmark/streams/iter-throughput-pipeto.jsβ€Ž

Lines changed: 26 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ const common = require('../common.js');
66
const{ Readable, Writable, pipeline }=require('stream');
77

88
constbench=common.createBenchmark(main,{
9-
api: ['classic','webstream','iter','iter-sync'],
9+
api: ['classic','webstream','iter','iter-sync-source','iter-sync'],
1010
datasize: [1024*1024,16*1024*1024,64*1024*1024],
1111
n: [5],
1212
},{
@@ -26,6 +26,8 @@ function main({ api, datasize, n }) {
2626
returnbenchWebStream(chunk,datasize,n,totalOps);
2727
case'iter':
2828
returnbenchIter(chunk,datasize,n,totalOps);
29+
case'iter-sync-source':
30+
returnbenchIterSyncSource(chunk,datasize,n,totalOps);
2931
case'iter-sync':
3032
returnbenchIterSync(chunk,datasize,n,totalOps);
3133
}
@@ -101,6 +103,29 @@ function benchIter(chunk, datasize, n, totalOps) {
101103
})();
102104
}
103105

106+
functionbenchIterSyncSource(chunk,datasize,n,totalOps){
107+
const{ pipeTo }=require('stream/iter');
108+
109+
asyncfunctionrun(){
110+
letremaining=datasize;
111+
function*source(){
112+
while(remaining>0){
113+
constsize=Math.min(remaining,chunk.length);
114+
remaining-=size;
115+
yieldsize===chunk.length ? chunk : chunk.subarray(0,size);
116+
}
117+
}
118+
constwriter={write(){},writeSync(){returntrue;}};
119+
awaitpipeTo(source(),writer);
120+
}
121+
122+
(async()=>{
123+
bench.start();
124+
for(leti=0;i<n;i++)awaitrun();
125+
bench.end(totalOps);
126+
})();
127+
}
128+
104129
functionbenchIterSync(chunk,datasize,n,totalOps){
105130
const{ pipeToSync }=require('stream/iter');
106131

β€Žlib/internal/streams/iter/pull.jsβ€Ž

Lines changed: 53 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,8 @@
88

99
const{
1010
ArrayBufferIsView,
11+
ArrayFromAsync,
12+
ArrayIsArray,
1113
ArrayPrototypePush,
1214
ArrayPrototypeSlice,
1315
PromisePrototypeThen,
@@ -38,7 +40,9 @@ const {
3840
fromSync,
3941
isSyncIterable,
4042
isAsyncIterable,
43+
isPrimitiveChunk,
4144
isUint8ArrayBatch,
45+
normalizeAsyncValue,
4246
}=require('internal/streams/iter/from');
4347

4448
const{
@@ -53,7 +57,10 @@ const {
5357
const{
5458
drainableProtocol,
5559
kSyncWriteAcceptedOnFalse,
60+
kValidatedSource,
5661
kValidatedTransform,
62+
toAsyncStreamable,
63+
toStreamable,
5764
}=require('internal/streams/iter/types');
5865

5966
// =============================================================================
@@ -116,6 +123,22 @@ function parsePipeToArgs(args, requiredMethod) {
116123
};
117124
}
118125

126+
functioncanUseSyncIterablePipeToFastPath(source,transforms,signal){
127+
if(signal!==undefined||
128+
transforms.length!==0||
129+
isPrimitiveChunk(source)||
130+
ArrayIsArray(source)||
131+
source?.[kValidatedSource]||
132+
!isSyncIterable(source)||
133+
isAsyncIterable(source)){
134+
returnfalse;
135+
}
136+
137+
// Preserve from()'s top-level protocol precedence for custom iterables.
138+
returntypeofsource[toAsyncStreamable]!=='function'&&
139+
typeofsource[toStreamable]!=='function';
140+
}
141+
119142
// =============================================================================
120143
// Transform Output Flattening
121144
// =============================================================================
@@ -822,12 +845,13 @@ async function pipeTo(source, ...args) {
822845
// Check for abort
823846
signal?.throwIfAborted();
824847

825-
// Normalize source via from()
826-
constnormalized=from(source);
848+
consthasWriteSync=typeofwriter.writeSync==='function';
849+
constuseSyncIterableFastPath=
850+
hasWriteSync&&canUseSyncIterablePipeToFastPath(source,transforms,signal);
851+
constnormalized=useSyncIterableFastPath ? undefined : from(source);
827852

828853
lettotalBytes=0;
829854
consthasWritev=typeofwriter.writev==='function';
830-
consthasWriteSync=typeofwriter.writeSync==='function';
831855
consthasWritevSync=typeofwriter.writevSync==='function';
832856
consthasEndSync=typeofwriter.endSync==='function';
833857
constsyncFalseCanBeAccepted=writer[kSyncWriteAcceptedOnFalse]===true;
@@ -908,8 +932,32 @@ async function pipeTo(source, ...args) {
908932
}
909933

910934
try{
911-
// Fast path: no transforms - iterate normalized source directly
912-
if(transforms.length===0){
935+
if(useSyncIterableFastPath){
936+
// Avoid from()'s async sync-iterable batching path. This keeps writes
937+
// incremental for synchronous sources while preserving async
938+
// normalization for non-primitive yielded values.
939+
for(constvalueofsource){
940+
if(isUint8ArrayBatch(value)){
941+
if(value.length>0){
942+
constp=writeBatch(value);
943+
if(p)awaitp;
944+
}
945+
continue;
946+
}
947+
if(isUint8Array(value)){
948+
constp=writeBatch([value]);
949+
if(p)awaitp;
950+
continue;
951+
}
952+
953+
constbatch=awaitArrayFromAsync(normalizeAsyncValue(value));
954+
if(batch.length>0){
955+
constp=writeBatch(batch);
956+
if(p)awaitp;
957+
}
958+
}
959+
}elseif(transforms.length===0){
960+
// Fast path: no transforms - iterate normalized source directly
913961
if(signal){
914962
forawait(constbatchofnormalized){
915963
signal.throwIfAborted();

β€Žtest/parallel/test-stream-iter-pipeto.jsβ€Ž

Lines changed: 76 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -219,6 +219,79 @@ async function testPipeToSyncMinimalWriter() {
219219
assert.strictEqual(chunks.length>0,true);
220220
}
221221

222+
asyncfunctiontestPipeToSyncIterableFastPathWritesIncrementally(){
223+
letpulled=0;
224+
letfirstWritePulled=0;
225+
constchunks=[];
226+
function*source(){
227+
for(leti=0;i<3;i++){
228+
pulled++;
229+
yieldnewUint8Array([0x61+i]);
230+
}
231+
}
232+
constwriter={
233+
write: common.mustNotCall(),
234+
writeSync(chunk){
235+
if(firstWritePulled===0){
236+
firstWritePulled=pulled;
237+
}
238+
chunks.push(chunk);
239+
returntrue;
240+
},
241+
};
242+
243+
consttotalBytes=awaitpipeTo(source(),writer);
244+
assert.strictEqual(totalBytes,3);
245+
assert.strictEqual(firstWritePulled,1);
246+
assert.deepStrictEqual(chunks,[
247+
newUint8Array([0x61]),
248+
newUint8Array([0x62]),
249+
newUint8Array([0x63]),
250+
]);
251+
}
252+
253+
asyncfunctiontestPipeToSyncIterableFastPathWriteFallback(){
254+
constasyncWrites=[];
255+
constwriter={
256+
writeSync(chunk){
257+
returnchunk[0]!==0x62;
258+
},
259+
asyncwrite(chunk){
260+
asyncWrites.push(chunk);
261+
},
262+
};
263+
function*source(){
264+
yieldnewUint8Array([0x61]);
265+
yieldnewUint8Array([0x62]);
266+
yieldnewUint8Array([0x63]);
267+
}
268+
269+
consttotalBytes=awaitpipeTo(source(),writer);
270+
assert.strictEqual(totalBytes,3);
271+
assert.deepStrictEqual(asyncWrites,[newUint8Array([0x62])]);
272+
}
273+
274+
asyncfunctiontestPipeToSyncIterableFastPathAsyncValue(){
275+
constchunks=[];
276+
constwriter={
277+
write: common.mustNotCall(),
278+
writeSync(chunk){
279+
chunks.push(chunk);
280+
returntrue;
281+
},
282+
};
283+
function*source(){
284+
yieldPromise.resolve('a');
285+
yieldnewUint8Array([0x62]);
286+
}
287+
288+
consttotalBytes=awaitpipeTo(source(),writer);
289+
assert.strictEqual(totalBytes,2);
290+
constresult=newTextDecoder().decode(
291+
newUint8Array(chunks.reduce((acc,c)=>[...acc, ...c],[])));
292+
assert.strictEqual(result,'ab');
293+
}
294+
222295
Promise.all([
223296
testPipeToSync(),
224297
testPipeTo(),
@@ -234,4 +307,7 @@ Promise.all([
234307
testPipeToSyncPreventClose(),
235308
testPipeToMinimalWriter(),
236309
testPipeToSyncMinimalWriter(),
310+
testPipeToSyncIterableFastPathWritesIncrementally(),
311+
testPipeToSyncIterableFastPathWriteFallback(),
312+
testPipeToSyncIterableFastPathAsyncValue(),
237313
]).then(common.mustCall());

0 commit comments

Comments
Β (0)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content

Commit 02a51a7

Browse files
trivikraduh95
authored andcommitted
stream: add sync iterable fast path to pipeTo
Avoid normalizing sync iterable sources through from() when pipeTo() has no transforms or signal and the writer can accept sync writes. This keeps writes incremental while preserving async fallback for values that still need it. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: openai:gpt-5.5 PR-URL: #63318 Backport-PR-URL: #64675 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Ethan Arrowood <ethan@arrowood.dev>
1 parent 4c9f3ed commit 02a51a7

3 files changed

Lines changed: 155 additions & 6 deletions

File tree

β€Žbenchmark/streams/iter-throughput-pipeto.jsβ€Ž

Lines changed: 26 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ const common = require('../common.js');
66
const{ Readable, Writable, pipeline }=require('stream');
77

88
constbench=common.createBenchmark(main,{
9-
api: ['classic','webstream','iter','iter-sync'],
9+
api: ['classic','webstream','iter','iter-sync-source','iter-sync'],
1010
datasize: [1024*1024,16*1024*1024,64*1024*1024],
1111
n: [5],
1212
},{
@@ -26,6 +26,8 @@ function main({ api, datasize, n }) {
2626
returnbenchWebStream(chunk,datasize,n,totalOps);
2727
case'iter':
2828
returnbenchIter(chunk,datasize,n,totalOps);
29+
case'iter-sync-source':
30+
returnbenchIterSyncSource(chunk,datasize,n,totalOps);
2931
case'iter-sync':
3032
returnbenchIterSync(chunk,datasize,n,totalOps);
3133
}
@@ -101,6 +103,29 @@ function benchIter(chunk, datasize, n, totalOps) {
101103
})();
102104
}
103105

106+
functionbenchIterSyncSource(chunk,datasize,n,totalOps){
107+
const{ pipeTo }=require('stream/iter');
108+
109+
asyncfunctionrun(){
110+
letremaining=datasize;
111+
function*source(){
112+
while(remaining>0){
113+
constsize=Math.min(remaining,chunk.length);
114+
remaining-=size;
115+
yieldsize===chunk.length ? chunk : chunk.subarray(0,size);
116+
}
117+
}
118+
constwriter={write(){},writeSync(){returntrue;}};
119+
awaitpipeTo(source(),writer);
120+
}
121+
122+
(async()=>{
123+
bench.start();
124+
for(leti=0;i<n;i++)awaitrun();
125+
bench.end(totalOps);
126+
})();
127+
}
128+
104129
functionbenchIterSync(chunk,datasize,n,totalOps){
105130
const{ pipeToSync }=require('stream/iter');
106131

β€Žlib/internal/streams/iter/pull.jsβ€Ž

Lines changed: 53 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,8 @@
88

99
const{
1010
ArrayBufferIsView,
11+
ArrayFromAsync,
12+
ArrayIsArray,
1113
ArrayPrototypePush,
1214
ArrayPrototypeSlice,
1315
PromisePrototypeThen,
@@ -38,7 +40,9 @@ const {
3840
fromSync,
3941
isSyncIterable,
4042
isAsyncIterable,
43+
isPrimitiveChunk,
4144
isUint8ArrayBatch,
45+
normalizeAsyncValue,
4246
}=require('internal/streams/iter/from');
4347

4448
const{
@@ -53,7 +57,10 @@ const {
5357
const{
5458
drainableProtocol,
5559
kSyncWriteAcceptedOnFalse,
60+
kValidatedSource,
5661
kValidatedTransform,
62+
toAsyncStreamable,
63+
toStreamable,
5764
}=require('internal/streams/iter/types');
5865

5966
// =============================================================================
@@ -116,6 +123,22 @@ function parsePipeToArgs(args, requiredMethod) {
116123
};
117124
}
118125

126+
functioncanUseSyncIterablePipeToFastPath(source,transforms,signal){
127+
if(signal!==undefined||
128+
transforms.length!==0||
129+
isPrimitiveChunk(source)||
130+
ArrayIsArray(source)||
131+
source?.[kValidatedSource]||
132+
!isSyncIterable(source)||
133+
isAsyncIterable(source)){
134+
returnfalse;
135+
}
136+
137+
// Preserve from()'s top-level protocol precedence for custom iterables.
138+
returntypeofsource[toAsyncStreamable]!=='function'&&
139+
typeofsource[toStreamable]!=='function';
140+
}
141+
119142
// =============================================================================
120143
// Transform Output Flattening
121144
// =============================================================================
@@ -822,12 +845,13 @@ async function pipeTo(source, ...args) {
822845
// Check for abort
823846
signal?.throwIfAborted();
824847

825-
// Normalize source via from()
826-
constnormalized=from(source);
848+
consthasWriteSync=typeofwriter.writeSync==='function';
849+
constuseSyncIterableFastPath=
850+
hasWriteSync&&canUseSyncIterablePipeToFastPath(source,transforms,signal);
851+
constnormalized=useSyncIterableFastPath ? undefined : from(source);
827852

828853
lettotalBytes=0;
829854
consthasWritev=typeofwriter.writev==='function';
830-
consthasWriteSync=typeofwriter.writeSync==='function';
831855
consthasWritevSync=typeofwriter.writevSync==='function';
832856
consthasEndSync=typeofwriter.endSync==='function';
833857
constsyncFalseCanBeAccepted=writer[kSyncWriteAcceptedOnFalse]===true;
@@ -908,8 +932,32 @@ async function pipeTo(source, ...args) {
908932
}
909933

910934
try{
911-
// Fast path: no transforms - iterate normalized source directly
912-
if(transforms.length===0){
935+
if(useSyncIterableFastPath){
936+
// Avoid from()'s async sync-iterable batching path. This keeps writes
937+
// incremental for synchronous sources while preserving async
938+
// normalization for non-primitive yielded values.
939+
for(constvalueofsource){
940+
if(isUint8ArrayBatch(value)){
941+
if(value.length>0){
942+
constp=writeBatch(value);
943+
if(p)awaitp;
944+
}
945+
continue;
946+
}
947+
if(isUint8Array(value)){
948+
constp=writeBatch([value]);
949+
if(p)awaitp;
950+
continue;
951+
}
952+
953+
constbatch=awaitArrayFromAsync(normalizeAsyncValue(value));
954+
if(batch.length>0){
955+
constp=writeBatch(batch);
956+
if(p)awaitp;
957+
}
958+
}
959+
}elseif(transforms.length===0){
960+
// Fast path: no transforms - iterate normalized source directly
913961
if(signal){
914962
forawait(constbatchofnormalized){
915963
signal.throwIfAborted();

β€Žtest/parallel/test-stream-iter-pipeto.jsβ€Ž

Lines changed: 76 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -219,6 +219,79 @@ async function testPipeToSyncMinimalWriter() {
219219
assert.strictEqual(chunks.length>0,true);
220220
}
221221

222+
asyncfunctiontestPipeToSyncIterableFastPathWritesIncrementally(){
223+
letpulled=0;
224+
letfirstWritePulled=0;
225+
constchunks=[];
226+
function*source(){
227+
for(leti=0;i<3;i++){
228+
pulled++;
229+
yieldnewUint8Array([0x61+i]);
230+
}
231+
}
232+
constwriter={
233+
write: common.mustNotCall(),
234+
writeSync(chunk){
235+
if(firstWritePulled===0){
236+
firstWritePulled=pulled;
237+
}
238+
chunks.push(chunk);
239+
returntrue;
240+
},
241+
};
242+
243+
consttotalBytes=awaitpipeTo(source(),writer);
244+
assert.strictEqual(totalBytes,3);
245+
assert.strictEqual(firstWritePulled,1);
246+
assert.deepStrictEqual(chunks,[
247+
newUint8Array([0x61]),
248+
newUint8Array([0x62]),
249+
newUint8Array([0x63]),
250+
]);
251+
}
252+
253+
asyncfunctiontestPipeToSyncIterableFastPathWriteFallback(){
254+
constasyncWrites=[];
255+
constwriter={
256+
writeSync(chunk){
257+
returnchunk[0]!==0x62;
258+
},
259+
asyncwrite(chunk){
260+
asyncWrites.push(chunk);
261+
},
262+
};
263+
function*source(){
264+
yieldnewUint8Array([0x61]);
265+
yieldnewUint8Array([0x62]);
266+
yieldnewUint8Array([0x63]);
267+
}
268+
269+
consttotalBytes=awaitpipeTo(source(),writer);
270+
assert.strictEqual(totalBytes,3);
271+
assert.deepStrictEqual(asyncWrites,[newUint8Array([0x62])]);
272+
}
273+
274+
asyncfunctiontestPipeToSyncIterableFastPathAsyncValue(){
275+
constchunks=[];
276+
constwriter={
277+
write: common.mustNotCall(),
278+
writeSync(chunk){
279+
chunks.push(chunk);
280+
returntrue;
281+
},
282+
};
283+
function*source(){
284+
yieldPromise.resolve('a');
285+
yieldnewUint8Array([0x62]);
286+
}
287+
288+
consttotalBytes=awaitpipeTo(source(),writer);
289+
assert.strictEqual(totalBytes,2);
290+
constresult=newTextDecoder().decode(
291+
newUint8Array(chunks.reduce((acc,c)=>[...acc, ...c],[])));
292+
assert.strictEqual(result,'ab');
293+
}
294+
222295
Promise.all([
223296
testPipeToSync(),
224297
testPipeTo(),
@@ -234,4 +307,7 @@ Promise.all([
234307
testPipeToSyncPreventClose(),
235308
testPipeToMinimalWriter(),
236309
testPipeToSyncMinimalWriter(),
310+
testPipeToSyncIterableFastPathWritesIncrementally(),
311+
testPipeToSyncIterableFastPathWriteFallback(),
312+
testPipeToSyncIterableFastPathAsyncValue(),
237313
]).then(common.mustCall());

0 commit comments

Comments
Β (0)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Commit 02a51a7

Browse files
trivikraduh95
authored andcommitted
stream: add sync iterable fast path to pipeTo
Avoid normalizing sync iterable sources through from() when pipeTo() has no transforms or signal and the writer can accept sync writes. This keeps writes incremental while preserving async fallback for values that still need it. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: openai:gpt-5.5 PR-URL: #63318 Backport-PR-URL: #64675 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Ethan Arrowood <ethan@arrowood.dev>
1 parent 4c9f3ed commit 02a51a7

3 files changed

Lines changed: 155 additions & 6 deletions

File tree

β€Žbenchmark/streams/iter-throughput-pipeto.jsβ€Ž

Lines changed: 26 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ const common = require('../common.js');
66
const{ Readable, Writable, pipeline }=require('stream');
77

88
constbench=common.createBenchmark(main,{
9-
api: ['classic','webstream','iter','iter-sync'],
9+
api: ['classic','webstream','iter','iter-sync-source','iter-sync'],
1010
datasize: [1024*1024,16*1024*1024,64*1024*1024],
1111
n: [5],
1212
},{
@@ -26,6 +26,8 @@ function main({ api, datasize, n }) {
2626
returnbenchWebStream(chunk,datasize,n,totalOps);
2727
case'iter':
2828
returnbenchIter(chunk,datasize,n,totalOps);
29+
case'iter-sync-source':
30+
returnbenchIterSyncSource(chunk,datasize,n,totalOps);
2931
case'iter-sync':
3032
returnbenchIterSync(chunk,datasize,n,totalOps);
3133
}
@@ -101,6 +103,29 @@ function benchIter(chunk, datasize, n, totalOps) {
101103
})();
102104
}
103105

106+
functionbenchIterSyncSource(chunk,datasize,n,totalOps){
107+
const{ pipeTo }=require('stream/iter');
108+
109+
asyncfunctionrun(){
110+
letremaining=datasize;
111+
function*source(){
112+
while(remaining>0){
113+
constsize=Math.min(remaining,chunk.length);
114+
remaining-=size;
115+
yieldsize===chunk.length ? chunk : chunk.subarray(0,size);
116+
}
117+
}
118+
constwriter={write(){},writeSync(){returntrue;}};
119+
awaitpipeTo(source(),writer);
120+
}
121+
122+
(async()=>{
123+
bench.start();
124+
for(leti=0;i<n;i++)awaitrun();
125+
bench.end(totalOps);
126+
})();
127+
}
128+
104129
functionbenchIterSync(chunk,datasize,n,totalOps){
105130
const{ pipeToSync }=require('stream/iter');
106131

β€Žlib/internal/streams/iter/pull.jsβ€Ž

Lines changed: 53 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,8 @@
88

99
const{
1010
ArrayBufferIsView,
11+
ArrayFromAsync,
12+
ArrayIsArray,
1113
ArrayPrototypePush,
1214
ArrayPrototypeSlice,
1315
PromisePrototypeThen,
@@ -38,7 +40,9 @@ const {
3840
fromSync,
3941
isSyncIterable,
4042
isAsyncIterable,
43+
isPrimitiveChunk,
4144
isUint8ArrayBatch,
45+
normalizeAsyncValue,
4246
}=require('internal/streams/iter/from');
4347

4448
const{
@@ -53,7 +57,10 @@ const {
5357
const{
5458
drainableProtocol,
5559
kSyncWriteAcceptedOnFalse,
60+
kValidatedSource,
5661
kValidatedTransform,
62+
toAsyncStreamable,
63+
toStreamable,
5764
}=require('internal/streams/iter/types');
5865

5966
// =============================================================================
@@ -116,6 +123,22 @@ function parsePipeToArgs(args, requiredMethod) {
116123
};
117124
}
118125

126+
functioncanUseSyncIterablePipeToFastPath(source,transforms,signal){
127+
if(signal!==undefined||
128+
transforms.length!==0||
129+
isPrimitiveChunk(source)||
130+
ArrayIsArray(source)||
131+
source?.[kValidatedSource]||
132+
!isSyncIterable(source)||
133+
isAsyncIterable(source)){
134+
returnfalse;
135+
}
136+
137+
// Preserve from()'s top-level protocol precedence for custom iterables.
138+
returntypeofsource[toAsyncStreamable]!=='function'&&
139+
typeofsource[toStreamable]!=='function';
140+
}
141+
119142
// =============================================================================
120143
// Transform Output Flattening
121144
// =============================================================================
@@ -822,12 +845,13 @@ async function pipeTo(source, ...args) {
822845
// Check for abort
823846
signal?.throwIfAborted();
824847

825-
// Normalize source via from()
826-
constnormalized=from(source);
848+
consthasWriteSync=typeofwriter.writeSync==='function';
849+
constuseSyncIterableFastPath=
850+
hasWriteSync&&canUseSyncIterablePipeToFastPath(source,transforms,signal);
851+
constnormalized=useSyncIterableFastPath ? undefined : from(source);
827852

828853
lettotalBytes=0;
829854
consthasWritev=typeofwriter.writev==='function';
830-
consthasWriteSync=typeofwriter.writeSync==='function';
831855
consthasWritevSync=typeofwriter.writevSync==='function';
832856
consthasEndSync=typeofwriter.endSync==='function';
833857
constsyncFalseCanBeAccepted=writer[kSyncWriteAcceptedOnFalse]===true;
@@ -908,8 +932,32 @@ async function pipeTo(source, ...args) {
908932
}
909933

910934
try{
911-
// Fast path: no transforms - iterate normalized source directly
912-
if(transforms.length===0){
935+
if(useSyncIterableFastPath){
936+
// Avoid from()'s async sync-iterable batching path. This keeps writes
937+
// incremental for synchronous sources while preserving async
938+
// normalization for non-primitive yielded values.
939+
for(constvalueofsource){
940+
if(isUint8ArrayBatch(value)){
941+
if(value.length>0){
942+
constp=writeBatch(value);
943+
if(p)awaitp;
944+
}
945+
continue;
946+
}
947+
if(isUint8Array(value)){
948+
constp=writeBatch([value]);
949+
if(p)awaitp;
950+
continue;
951+
}
952+
953+
constbatch=awaitArrayFromAsync(normalizeAsyncValue(value));
954+
if(batch.length>0){
955+
constp=writeBatch(batch);
956+
if(p)awaitp;
957+
}
958+
}
959+
}elseif(transforms.length===0){
960+
// Fast path: no transforms - iterate normalized source directly
913961
if(signal){
914962
forawait(constbatchofnormalized){
915963
signal.throwIfAborted();

β€Žtest/parallel/test-stream-iter-pipeto.jsβ€Ž

Lines changed: 76 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -219,6 +219,79 @@ async function testPipeToSyncMinimalWriter() {
219219
assert.strictEqual(chunks.length>0,true);
220220
}
221221

222+
asyncfunctiontestPipeToSyncIterableFastPathWritesIncrementally(){
223+
letpulled=0;
224+
letfirstWritePulled=0;
225+
constchunks=[];
226+
function*source(){
227+
for(leti=0;i<3;i++){
228+
pulled++;
229+
yieldnewUint8Array([0x61+i]);
230+
}
231+
}
232+
constwriter={
233+
write: common.mustNotCall(),
234+
writeSync(chunk){
235+
if(firstWritePulled===0){
236+
firstWritePulled=pulled;
237+
}
238+
chunks.push(chunk);
239+
returntrue;
240+
},
241+
};
242+
243+
consttotalBytes=awaitpipeTo(source(),writer);
244+
assert.strictEqual(totalBytes,3);
245+
assert.strictEqual(firstWritePulled,1);
246+
assert.deepStrictEqual(chunks,[
247+
newUint8Array([0x61]),
248+
newUint8Array([0x62]),
249+
newUint8Array([0x63]),
250+
]);
251+
}
252+
253+
asyncfunctiontestPipeToSyncIterableFastPathWriteFallback(){
254+
constasyncWrites=[];
255+
constwriter={
256+
writeSync(chunk){
257+
returnchunk[0]!==0x62;
258+
},
259+
asyncwrite(chunk){
260+
asyncWrites.push(chunk);
261+
},
262+
};
263+
function*source(){
264+
yieldnewUint8Array([0x61]);
265+
yieldnewUint8Array([0x62]);
266+
yieldnewUint8Array([0x63]);
267+
}
268+
269+
consttotalBytes=awaitpipeTo(source(),writer);
270+
assert.strictEqual(totalBytes,3);
271+
assert.deepStrictEqual(asyncWrites,[newUint8Array([0x62])]);
272+
}
273+
274+
asyncfunctiontestPipeToSyncIterableFastPathAsyncValue(){
275+
constchunks=[];
276+
constwriter={
277+
write: common.mustNotCall(),
278+
writeSync(chunk){
279+
chunks.push(chunk);
280+
returntrue;
281+
},
282+
};
283+
function*source(){
284+
yieldPromise.resolve('a');
285+
yieldnewUint8Array([0x62]);
286+
}
287+
288+
consttotalBytes=awaitpipeTo(source(),writer);
289+
assert.strictEqual(totalBytes,2);
290+
constresult=newTextDecoder().decode(
291+
newUint8Array(chunks.reduce((acc,c)=>[...acc, ...c],[])));
292+
assert.strictEqual(result,'ab');
293+
}
294+
222295
Promise.all([
223296
testPipeToSync(),
224297
testPipeTo(),
@@ -234,4 +307,7 @@ Promise.all([
234307
testPipeToSyncPreventClose(),
235308
testPipeToMinimalWriter(),
236309
testPipeToSyncMinimalWriter(),
310+
testPipeToSyncIterableFastPathWritesIncrementally(),
311+
testPipeToSyncIterableFastPathWriteFallback(),
312+
testPipeToSyncIterableFastPathAsyncValue(),
237313
]).then(common.mustCall());

0 commit comments

Comments
Β (0)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Commit 02a51a7

Browse files
trivikraduh95
authored andcommitted
stream: add sync iterable fast path to pipeTo
Avoid normalizing sync iterable sources through from() when pipeTo() has no transforms or signal and the writer can accept sync writes. This keeps writes incremental while preserving async fallback for values that still need it. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: openai:gpt-5.5 PR-URL: #63318 Backport-PR-URL: #64675 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Ethan Arrowood <ethan@arrowood.dev>
1 parent 4c9f3ed commit 02a51a7

3 files changed

Lines changed: 155 additions & 6 deletions

File tree

β€Žbenchmark/streams/iter-throughput-pipeto.jsβ€Ž

Lines changed: 26 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ const common = require('../common.js');
66
const{ Readable, Writable, pipeline }=require('stream');
77

88
constbench=common.createBenchmark(main,{
9-
api: ['classic','webstream','iter','iter-sync'],
9+
api: ['classic','webstream','iter','iter-sync-source','iter-sync'],
1010
datasize: [1024*1024,16*1024*1024,64*1024*1024],
1111
n: [5],
1212
},{
@@ -26,6 +26,8 @@ function main({ api, datasize, n }) {
2626
returnbenchWebStream(chunk,datasize,n,totalOps);
2727
case'iter':
2828
returnbenchIter(chunk,datasize,n,totalOps);
29+
case'iter-sync-source':
30+
returnbenchIterSyncSource(chunk,datasize,n,totalOps);
2931
case'iter-sync':
3032
returnbenchIterSync(chunk,datasize,n,totalOps);
3133
}
@@ -101,6 +103,29 @@ function benchIter(chunk, datasize, n, totalOps) {
101103
})();
102104
}
103105

106+
functionbenchIterSyncSource(chunk,datasize,n,totalOps){
107+
const{ pipeTo }=require('stream/iter');
108+
109+
asyncfunctionrun(){
110+
letremaining=datasize;
111+
function*source(){
112+
while(remaining>0){
113+
constsize=Math.min(remaining,chunk.length);
114+
remaining-=size;
115+
yieldsize===chunk.length ? chunk : chunk.subarray(0,size);
116+
}
117+
}
118+
constwriter={write(){},writeSync(){returntrue;}};
119+
awaitpipeTo(source(),writer);
120+
}
121+
122+
(async()=>{
123+
bench.start();
124+
for(leti=0;i<n;i++)awaitrun();
125+
bench.end(totalOps);
126+
})();
127+
}
128+
104129
functionbenchIterSync(chunk,datasize,n,totalOps){
105130
const{ pipeToSync }=require('stream/iter');
106131

β€Žlib/internal/streams/iter/pull.jsβ€Ž

Lines changed: 53 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,8 @@
88

99
const{
1010
ArrayBufferIsView,
11+
ArrayFromAsync,
12+
ArrayIsArray,
1113
ArrayPrototypePush,
1214
ArrayPrototypeSlice,
1315
PromisePrototypeThen,
@@ -38,7 +40,9 @@ const {
3840
fromSync,
3941
isSyncIterable,
4042
isAsyncIterable,
43+
isPrimitiveChunk,
4144
isUint8ArrayBatch,
45+
normalizeAsyncValue,
4246
}=require('internal/streams/iter/from');
4347

4448
const{
@@ -53,7 +57,10 @@ const {
5357
const{
5458
drainableProtocol,
5559
kSyncWriteAcceptedOnFalse,
60+
kValidatedSource,
5661
kValidatedTransform,
62+
toAsyncStreamable,
63+
toStreamable,
5764
}=require('internal/streams/iter/types');
5865

5966
// =============================================================================
@@ -116,6 +123,22 @@ function parsePipeToArgs(args, requiredMethod) {
116123
};
117124
}
118125

126+
functioncanUseSyncIterablePipeToFastPath(source,transforms,signal){
127+
if(signal!==undefined||
128+
transforms.length!==0||
129+
isPrimitiveChunk(source)||
130+
ArrayIsArray(source)||
131+
source?.[kValidatedSource]||
132+
!isSyncIterable(source)||
133+
isAsyncIterable(source)){
134+
returnfalse;
135+
}
136+
137+
// Preserve from()'s top-level protocol precedence for custom iterables.
138+
returntypeofsource[toAsyncStreamable]!=='function'&&
139+
typeofsource[toStreamable]!=='function';
140+
}
141+
119142
// =============================================================================
120143
// Transform Output Flattening
121144
// =============================================================================
@@ -822,12 +845,13 @@ async function pipeTo(source, ...args) {
822845
// Check for abort
823846
signal?.throwIfAborted();
824847

825-
// Normalize source via from()
826-
constnormalized=from(source);
848+
consthasWriteSync=typeofwriter.writeSync==='function';
849+
constuseSyncIterableFastPath=
850+
hasWriteSync&&canUseSyncIterablePipeToFastPath(source,transforms,signal);
851+
constnormalized=useSyncIterableFastPath ? undefined : from(source);
827852

828853
lettotalBytes=0;
829854
consthasWritev=typeofwriter.writev==='function';
830-
consthasWriteSync=typeofwriter.writeSync==='function';
831855
consthasWritevSync=typeofwriter.writevSync==='function';
832856
consthasEndSync=typeofwriter.endSync==='function';
833857
constsyncFalseCanBeAccepted=writer[kSyncWriteAcceptedOnFalse]===true;
@@ -908,8 +932,32 @@ async function pipeTo(source, ...args) {
908932
}
909933

910934
try{
911-
// Fast path: no transforms - iterate normalized source directly
912-
if(transforms.length===0){
935+
if(useSyncIterableFastPath){
936+
// Avoid from()'s async sync-iterable batching path. This keeps writes
937+
// incremental for synchronous sources while preserving async
938+
// normalization for non-primitive yielded values.
939+
for(constvalueofsource){
940+
if(isUint8ArrayBatch(value)){
941+
if(value.length>0){
942+
constp=writeBatch(value);
943+
if(p)awaitp;
944+
}
945+
continue;
946+
}
947+
if(isUint8Array(value)){
948+
constp=writeBatch([value]);
949+
if(p)awaitp;
950+
continue;
951+
}
952+
953+
constbatch=awaitArrayFromAsync(normalizeAsyncValue(value));
954+
if(batch.length>0){
955+
constp=writeBatch(batch);
956+
if(p)awaitp;
957+
}
958+
}
959+
}elseif(transforms.length===0){
960+
// Fast path: no transforms - iterate normalized source directly
913961
if(signal){
914962
forawait(constbatchofnormalized){
915963
signal.throwIfAborted();

β€Žtest/parallel/test-stream-iter-pipeto.jsβ€Ž

Lines changed: 76 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -219,6 +219,79 @@ async function testPipeToSyncMinimalWriter() {
219219
assert.strictEqual(chunks.length>0,true);
220220
}
221221

222+
asyncfunctiontestPipeToSyncIterableFastPathWritesIncrementally(){
223+
letpulled=0;
224+
letfirstWritePulled=0;
225+
constchunks=[];
226+
function*source(){
227+
for(leti=0;i<3;i++){
228+
pulled++;
229+
yieldnewUint8Array([0x61+i]);
230+
}
231+
}
232+
constwriter={
233+
write: common.mustNotCall(),
234+
writeSync(chunk){
235+
if(firstWritePulled===0){
236+
firstWritePulled=pulled;
237+
}
238+
chunks.push(chunk);
239+
returntrue;
240+
},
241+
};
242+
243+
consttotalBytes=awaitpipeTo(source(),writer);
244+
assert.strictEqual(totalBytes,3);
245+
assert.strictEqual(firstWritePulled,1);
246+
assert.deepStrictEqual(chunks,[
247+
newUint8Array([0x61]),
248+
newUint8Array([0x62]),
249+
newUint8Array([0x63]),
250+
]);
251+
}
252+
253+
asyncfunctiontestPipeToSyncIterableFastPathWriteFallback(){
254+
constasyncWrites=[];
255+
constwriter={
256+
writeSync(chunk){
257+
returnchunk[0]!==0x62;
258+
},
259+
asyncwrite(chunk){
260+
asyncWrites.push(chunk);
261+
},
262+
};
263+
function*source(){
264+
yieldnewUint8Array([0x61]);
265+
yieldnewUint8Array([0x62]);
266+
yieldnewUint8Array([0x63]);
267+
}
268+
269+
consttotalBytes=awaitpipeTo(source(),writer);
270+
assert.strictEqual(totalBytes,3);
271+
assert.deepStrictEqual(asyncWrites,[newUint8Array([0x62])]);
272+
}
273+
274+
asyncfunctiontestPipeToSyncIterableFastPathAsyncValue(){
275+
constchunks=[];
276+
constwriter={
277+
write: common.mustNotCall(),
278+
writeSync(chunk){
279+
chunks.push(chunk);
280+
returntrue;
281+
},
282+
};
283+
function*source(){
284+
yieldPromise.resolve('a');
285+
yieldnewUint8Array([0x62]);
286+
}
287+
288+
consttotalBytes=awaitpipeTo(source(),writer);
289+
assert.strictEqual(totalBytes,2);
290+
constresult=newTextDecoder().decode(
291+
newUint8Array(chunks.reduce((acc,c)=>[...acc, ...c],[])));
292+
assert.strictEqual(result,'ab');
293+
}
294+
222295
Promise.all([
223296
testPipeToSync(),
224297
testPipeTo(),
@@ -234,4 +307,7 @@ Promise.all([
234307
testPipeToSyncPreventClose(),
235308
testPipeToMinimalWriter(),
236309
testPipeToSyncMinimalWriter(),
310+
testPipeToSyncIterableFastPathWritesIncrementally(),
311+
testPipeToSyncIterableFastPathWriteFallback(),
312+
testPipeToSyncIterableFastPathAsyncValue(),
237313
]).then(common.mustCall());

0 commit comments

Comments
Β (0)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content

Commit 02a51a7

Browse files
trivikraduh95
authored andcommitted
stream: add sync iterable fast path to pipeTo
Avoid normalizing sync iterable sources through from() when pipeTo() has no transforms or signal and the writer can accept sync writes. This keeps writes incremental while preserving async fallback for values that still need it. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: openai:gpt-5.5 PR-URL: #63318 Backport-PR-URL: #64675 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Ethan Arrowood <ethan@arrowood.dev>
1 parent 4c9f3ed commit 02a51a7

3 files changed

Lines changed: 155 additions & 6 deletions

File tree

β€Žbenchmark/streams/iter-throughput-pipeto.jsβ€Ž

Lines changed: 26 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ const common = require('../common.js');
66
const{ Readable, Writable, pipeline }=require('stream');
77

88
constbench=common.createBenchmark(main,{
9-
api: ['classic','webstream','iter','iter-sync'],
9+
api: ['classic','webstream','iter','iter-sync-source','iter-sync'],
1010
datasize: [1024*1024,16*1024*1024,64*1024*1024],
1111
n: [5],
1212
},{
@@ -26,6 +26,8 @@ function main({ api, datasize, n }) {
2626
returnbenchWebStream(chunk,datasize,n,totalOps);
2727
case'iter':
2828
returnbenchIter(chunk,datasize,n,totalOps);
29+
case'iter-sync-source':
30+
returnbenchIterSyncSource(chunk,datasize,n,totalOps);
2931
case'iter-sync':
3032
returnbenchIterSync(chunk,datasize,n,totalOps);
3133
}
@@ -101,6 +103,29 @@ function benchIter(chunk, datasize, n, totalOps) {
101103
})();
102104
}
103105

106+
functionbenchIterSyncSource(chunk,datasize,n,totalOps){
107+
const{ pipeTo }=require('stream/iter');
108+
109+
asyncfunctionrun(){
110+
letremaining=datasize;
111+
function*source(){
112+
while(remaining>0){
113+
constsize=Math.min(remaining,chunk.length);
114+
remaining-=size;
115+
yieldsize===chunk.length ? chunk : chunk.subarray(0,size);
116+
}
117+
}
118+
constwriter={write(){},writeSync(){returntrue;}};
119+
awaitpipeTo(source(),writer);
120+
}
121+
122+
(async()=>{
123+
bench.start();
124+
for(leti=0;i<n;i++)awaitrun();
125+
bench.end(totalOps);
126+
})();
127+
}
128+
104129
functionbenchIterSync(chunk,datasize,n,totalOps){
105130
const{ pipeToSync }=require('stream/iter');
106131

β€Žlib/internal/streams/iter/pull.jsβ€Ž

Lines changed: 53 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,8 @@
88

99
const{
1010
ArrayBufferIsView,
11+
ArrayFromAsync,
12+
ArrayIsArray,
1113
ArrayPrototypePush,
1214
ArrayPrototypeSlice,
1315
PromisePrototypeThen,
@@ -38,7 +40,9 @@ const {
3840
fromSync,
3941
isSyncIterable,
4042
isAsyncIterable,
43+
isPrimitiveChunk,
4144
isUint8ArrayBatch,
45+
normalizeAsyncValue,
4246
}=require('internal/streams/iter/from');
4347

4448
const{
@@ -53,7 +57,10 @@ const {
5357
const{
5458
drainableProtocol,
5559
kSyncWriteAcceptedOnFalse,
60+
kValidatedSource,
5661
kValidatedTransform,
62+
toAsyncStreamable,
63+
toStreamable,
5764
}=require('internal/streams/iter/types');
5865

5966
// =============================================================================
@@ -116,6 +123,22 @@ function parsePipeToArgs(args, requiredMethod) {
116123
};
117124
}
118125

126+
functioncanUseSyncIterablePipeToFastPath(source,transforms,signal){
127+
if(signal!==undefined||
128+
transforms.length!==0||
129+
isPrimitiveChunk(source)||
130+
ArrayIsArray(source)||
131+
source?.[kValidatedSource]||
132+
!isSyncIterable(source)||
133+
isAsyncIterable(source)){
134+
returnfalse;
135+
}
136+
137+
// Preserve from()'s top-level protocol precedence for custom iterables.
138+
returntypeofsource[toAsyncStreamable]!=='function'&&
139+
typeofsource[toStreamable]!=='function';
140+
}
141+
119142
// =============================================================================
120143
// Transform Output Flattening
121144
// =============================================================================
@@ -822,12 +845,13 @@ async function pipeTo(source, ...args) {
822845
// Check for abort
823846
signal?.throwIfAborted();
824847

825-
// Normalize source via from()
826-
constnormalized=from(source);
848+
consthasWriteSync=typeofwriter.writeSync==='function';
849+
constuseSyncIterableFastPath=
850+
hasWriteSync&&canUseSyncIterablePipeToFastPath(source,transforms,signal);
851+
constnormalized=useSyncIterableFastPath ? undefined : from(source);
827852

828853
lettotalBytes=0;
829854
consthasWritev=typeofwriter.writev==='function';
830-
consthasWriteSync=typeofwriter.writeSync==='function';
831855
consthasWritevSync=typeofwriter.writevSync==='function';
832856
consthasEndSync=typeofwriter.endSync==='function';
833857
constsyncFalseCanBeAccepted=writer[kSyncWriteAcceptedOnFalse]===true;
@@ -908,8 +932,32 @@ async function pipeTo(source, ...args) {
908932
}
909933

910934
try{
911-
// Fast path: no transforms - iterate normalized source directly
912-
if(transforms.length===0){
935+
if(useSyncIterableFastPath){
936+
// Avoid from()'s async sync-iterable batching path. This keeps writes
937+
// incremental for synchronous sources while preserving async
938+
// normalization for non-primitive yielded values.
939+
for(constvalueofsource){
940+
if(isUint8ArrayBatch(value)){
941+
if(value.length>0){
942+
constp=writeBatch(value);
943+
if(p)awaitp;
944+
}
945+
continue;
946+
}
947+
if(isUint8Array(value)){
948+
constp=writeBatch([value]);
949+
if(p)awaitp;
950+
continue;
951+
}
952+
953+
constbatch=awaitArrayFromAsync(normalizeAsyncValue(value));
954+
if(batch.length>0){
955+
constp=writeBatch(batch);
956+
if(p)awaitp;
957+
}
958+
}
959+
}elseif(transforms.length===0){
960+
// Fast path: no transforms - iterate normalized source directly
913961
if(signal){
914962
forawait(constbatchofnormalized){
915963
signal.throwIfAborted();

β€Žtest/parallel/test-stream-iter-pipeto.jsβ€Ž

Lines changed: 76 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -219,6 +219,79 @@ async function testPipeToSyncMinimalWriter() {
219219
assert.strictEqual(chunks.length>0,true);
220220
}
221221

222+
asyncfunctiontestPipeToSyncIterableFastPathWritesIncrementally(){
223+
letpulled=0;
224+
letfirstWritePulled=0;
225+
constchunks=[];
226+
function*source(){
227+
for(leti=0;i<3;i++){
228+
pulled++;
229+
yieldnewUint8Array([0x61+i]);
230+
}
231+
}
232+
constwriter={
233+
write: common.mustNotCall(),
234+
writeSync(chunk){
235+
if(firstWritePulled===0){
236+
firstWritePulled=pulled;
237+
}
238+
chunks.push(chunk);
239+
returntrue;
240+
},
241+
};
242+
243+
consttotalBytes=awaitpipeTo(source(),writer);
244+
assert.strictEqual(totalBytes,3);
245+
assert.strictEqual(firstWritePulled,1);
246+
assert.deepStrictEqual(chunks,[
247+
newUint8Array([0x61]),
248+
newUint8Array([0x62]),
249+
newUint8Array([0x63]),
250+
]);
251+
}
252+
253+
asyncfunctiontestPipeToSyncIterableFastPathWriteFallback(){
254+
constasyncWrites=[];
255+
constwriter={
256+
writeSync(chunk){
257+
returnchunk[0]!==0x62;
258+
},
259+
asyncwrite(chunk){
260+
asyncWrites.push(chunk);
261+
},
262+
};
263+
function*source(){
264+
yieldnewUint8Array([0x61]);
265+
yieldnewUint8Array([0x62]);
266+
yieldnewUint8Array([0x63]);
267+
}
268+
269+
consttotalBytes=awaitpipeTo(source(),writer);
270+
assert.strictEqual(totalBytes,3);
271+
assert.deepStrictEqual(asyncWrites,[newUint8Array([0x62])]);
272+
}
273+
274+
asyncfunctiontestPipeToSyncIterableFastPathAsyncValue(){
275+
constchunks=[];
276+
constwriter={
277+
write: common.mustNotCall(),
278+
writeSync(chunk){
279+
chunks.push(chunk);
280+
returntrue;
281+
},
282+
};
283+
function*source(){
284+
yieldPromise.resolve('a');
285+
yieldnewUint8Array([0x62]);
286+
}
287+
288+
consttotalBytes=awaitpipeTo(source(),writer);
289+
assert.strictEqual(totalBytes,2);
290+
constresult=newTextDecoder().decode(
291+
newUint8Array(chunks.reduce((acc,c)=>[...acc, ...c],[])));
292+
assert.strictEqual(result,'ab');
293+
}
294+
222295
Promise.all([
223296
testPipeToSync(),
224297
testPipeTo(),
@@ -234,4 +307,7 @@ Promise.all([
234307
testPipeToSyncPreventClose(),
235308
testPipeToMinimalWriter(),
236309
testPipeToSyncMinimalWriter(),
310+
testPipeToSyncIterableFastPathWritesIncrementally(),
311+
testPipeToSyncIterableFastPathWriteFallback(),
312+
testPipeToSyncIterableFastPathAsyncValue(),
237313
]).then(common.mustCall());

0 commit comments

Comments
Β (0)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Commit 02a51a7

Browse files
trivikraduh95
authored andcommitted
stream: add sync iterable fast path to pipeTo
Avoid normalizing sync iterable sources through from() when pipeTo() has no transforms or signal and the writer can accept sync writes. This keeps writes incremental while preserving async fallback for values that still need it. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: openai:gpt-5.5 PR-URL: #63318 Backport-PR-URL: #64675 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Ethan Arrowood <ethan@arrowood.dev>
1 parent 4c9f3ed commit 02a51a7

3 files changed

Lines changed: 155 additions & 6 deletions

File tree

β€Žbenchmark/streams/iter-throughput-pipeto.jsβ€Ž

Lines changed: 26 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ const common = require('../common.js');
66
const{ Readable, Writable, pipeline }=require('stream');
77

88
constbench=common.createBenchmark(main,{
9-
api: ['classic','webstream','iter','iter-sync'],
9+
api: ['classic','webstream','iter','iter-sync-source','iter-sync'],
1010
datasize: [1024*1024,16*1024*1024,64*1024*1024],
1111
n: [5],
1212
},{
@@ -26,6 +26,8 @@ function main({ api, datasize, n }) {
2626
returnbenchWebStream(chunk,datasize,n,totalOps);
2727
case'iter':
2828
returnbenchIter(chunk,datasize,n,totalOps);
29+
case'iter-sync-source':
30+
returnbenchIterSyncSource(chunk,datasize,n,totalOps);
2931
case'iter-sync':
3032
returnbenchIterSync(chunk,datasize,n,totalOps);
3133
}
@@ -101,6 +103,29 @@ function benchIter(chunk, datasize, n, totalOps) {
101103
})();
102104
}
103105

106+
functionbenchIterSyncSource(chunk,datasize,n,totalOps){
107+
const{ pipeTo }=require('stream/iter');
108+
109+
asyncfunctionrun(){
110+
letremaining=datasize;
111+
function*source(){
112+
while(remaining>0){
113+
constsize=Math.min(remaining,chunk.length);
114+
remaining-=size;
115+
yieldsize===chunk.length ? chunk : chunk.subarray(0,size);
116+
}
117+
}
118+
constwriter={write(){},writeSync(){returntrue;}};
119+
awaitpipeTo(source(),writer);
120+
}
121+
122+
(async()=>{
123+
bench.start();
124+
for(leti=0;i<n;i++)awaitrun();
125+
bench.end(totalOps);
126+
})();
127+
}
128+
104129
functionbenchIterSync(chunk,datasize,n,totalOps){
105130
const{ pipeToSync }=require('stream/iter');
106131

β€Žlib/internal/streams/iter/pull.jsβ€Ž

Lines changed: 53 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,8 @@
88

99
const{
1010
ArrayBufferIsView,
11+
ArrayFromAsync,
12+
ArrayIsArray,
1113
ArrayPrototypePush,
1214
ArrayPrototypeSlice,
1315
PromisePrototypeThen,
@@ -38,7 +40,9 @@ const {
3840
fromSync,
3941
isSyncIterable,
4042
isAsyncIterable,
43+
isPrimitiveChunk,
4144
isUint8ArrayBatch,
45+
normalizeAsyncValue,
4246
}=require('internal/streams/iter/from');
4347

4448
const{
@@ -53,7 +57,10 @@ const {
5357
const{
5458
drainableProtocol,
5559
kSyncWriteAcceptedOnFalse,
60+
kValidatedSource,
5661
kValidatedTransform,
62+
toAsyncStreamable,
63+
toStreamable,
5764
}=require('internal/streams/iter/types');
5865

5966
// =============================================================================
@@ -116,6 +123,22 @@ function parsePipeToArgs(args, requiredMethod) {
116123
};
117124
}
118125

126+
functioncanUseSyncIterablePipeToFastPath(source,transforms,signal){
127+
if(signal!==undefined||
128+
transforms.length!==0||
129+
isPrimitiveChunk(source)||
130+
ArrayIsArray(source)||
131+
source?.[kValidatedSource]||
132+
!isSyncIterable(source)||
133+
isAsyncIterable(source)){
134+
returnfalse;
135+
}
136+
137+
// Preserve from()'s top-level protocol precedence for custom iterables.
138+
returntypeofsource[toAsyncStreamable]!=='function'&&
139+
typeofsource[toStreamable]!=='function';
140+
}
141+
119142
// =============================================================================
120143
// Transform Output Flattening
121144
// =============================================================================
@@ -822,12 +845,13 @@ async function pipeTo(source, ...args) {
822845
// Check for abort
823846
signal?.throwIfAborted();
824847

825-
// Normalize source via from()
826-
constnormalized=from(source);
848+
consthasWriteSync=typeofwriter.writeSync==='function';
849+
constuseSyncIterableFastPath=
850+
hasWriteSync&&canUseSyncIterablePipeToFastPath(source,transforms,signal);
851+
constnormalized=useSyncIterableFastPath ? undefined : from(source);
827852

828853
lettotalBytes=0;
829854
consthasWritev=typeofwriter.writev==='function';
830-
consthasWriteSync=typeofwriter.writeSync==='function';
831855
consthasWritevSync=typeofwriter.writevSync==='function';
832856
consthasEndSync=typeofwriter.endSync==='function';
833857
constsyncFalseCanBeAccepted=writer[kSyncWriteAcceptedOnFalse]===true;
@@ -908,8 +932,32 @@ async function pipeTo(source, ...args) {
908932
}
909933

910934
try{
911-
// Fast path: no transforms - iterate normalized source directly
912-
if(transforms.length===0){
935+
if(useSyncIterableFastPath){
936+
// Avoid from()'s async sync-iterable batching path. This keeps writes
937+
// incremental for synchronous sources while preserving async
938+
// normalization for non-primitive yielded values.
939+
for(constvalueofsource){
940+
if(isUint8ArrayBatch(value)){
941+
if(value.length>0){
942+
constp=writeBatch(value);
943+
if(p)awaitp;
944+
}
945+
continue;
946+
}
947+
if(isUint8Array(value)){
948+
constp=writeBatch([value]);
949+
if(p)awaitp;
950+
continue;
951+
}
952+
953+
constbatch=awaitArrayFromAsync(normalizeAsyncValue(value));
954+
if(batch.length>0){
955+
constp=writeBatch(batch);
956+
if(p)awaitp;
957+
}
958+
}
959+
}elseif(transforms.length===0){
960+
// Fast path: no transforms - iterate normalized source directly
913961
if(signal){
914962
forawait(constbatchofnormalized){
915963
signal.throwIfAborted();

β€Žtest/parallel/test-stream-iter-pipeto.jsβ€Ž

Lines changed: 76 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -219,6 +219,79 @@ async function testPipeToSyncMinimalWriter() {
219219
assert.strictEqual(chunks.length>0,true);
220220
}
221221

222+
asyncfunctiontestPipeToSyncIterableFastPathWritesIncrementally(){
223+
letpulled=0;
224+
letfirstWritePulled=0;
225+
constchunks=[];
226+
function*source(){
227+
for(leti=0;i<3;i++){
228+
pulled++;
229+
yieldnewUint8Array([0x61+i]);
230+
}
231+
}
232+
constwriter={
233+
write: common.mustNotCall(),
234+
writeSync(chunk){
235+
if(firstWritePulled===0){
236+
firstWritePulled=pulled;
237+
}
238+
chunks.push(chunk);
239+
returntrue;
240+
},
241+
};
242+
243+
consttotalBytes=awaitpipeTo(source(),writer);
244+
assert.strictEqual(totalBytes,3);
245+
assert.strictEqual(firstWritePulled,1);
246+
assert.deepStrictEqual(chunks,[
247+
newUint8Array([0x61]),
248+
newUint8Array([0x62]),
249+
newUint8Array([0x63]),
250+
]);
251+
}
252+
253+
asyncfunctiontestPipeToSyncIterableFastPathWriteFallback(){
254+
constasyncWrites=[];
255+
constwriter={
256+
writeSync(chunk){
257+
returnchunk[0]!==0x62;
258+
},
259+
asyncwrite(chunk){
260+
asyncWrites.push(chunk);
261+
},
262+
};
263+
function*source(){
264+
yieldnewUint8Array([0x61]);
265+
yieldnewUint8Array([0x62]);
266+
yieldnewUint8Array([0x63]);
267+
}
268+
269+
consttotalBytes=awaitpipeTo(source(),writer);
270+
assert.strictEqual(totalBytes,3);
271+
assert.deepStrictEqual(asyncWrites,[newUint8Array([0x62])]);
272+
}
273+
274+
asyncfunctiontestPipeToSyncIterableFastPathAsyncValue(){
275+
constchunks=[];
276+
constwriter={
277+
write: common.mustNotCall(),
278+
writeSync(chunk){
279+
chunks.push(chunk);
280+
returntrue;
281+
},
282+
};
283+
function*source(){
284+
yieldPromise.resolve('a');
285+
yieldnewUint8Array([0x62]);
286+
}
287+
288+
consttotalBytes=awaitpipeTo(source(),writer);
289+
assert.strictEqual(totalBytes,2);
290+
constresult=newTextDecoder().decode(
291+
newUint8Array(chunks.reduce((acc,c)=>[...acc, ...c],[])));
292+
assert.strictEqual(result,'ab');
293+
}
294+
222295
Promise.all([
223296
testPipeToSync(),
224297
testPipeTo(),
@@ -234,4 +307,7 @@ Promise.all([
234307
testPipeToSyncPreventClose(),
235308
testPipeToMinimalWriter(),
236309
testPipeToSyncMinimalWriter(),
310+
testPipeToSyncIterableFastPathWritesIncrementally(),
311+
testPipeToSyncIterableFastPathWriteFallback(),
312+
testPipeToSyncIterableFastPathAsyncValue(),
237313
]).then(common.mustCall());

0 commit comments

Comments
Β (0)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Commit 02a51a7

Browse files
trivikraduh95
authored andcommitted
stream: add sync iterable fast path to pipeTo
Avoid normalizing sync iterable sources through from() when pipeTo() has no transforms or signal and the writer can accept sync writes. This keeps writes incremental while preserving async fallback for values that still need it. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: openai:gpt-5.5 PR-URL: #63318 Backport-PR-URL: #64675 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Ethan Arrowood <ethan@arrowood.dev>
1 parent 4c9f3ed commit 02a51a7

3 files changed

Lines changed: 155 additions & 6 deletions

File tree

β€Žbenchmark/streams/iter-throughput-pipeto.jsβ€Ž

Lines changed: 26 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ const common = require('../common.js');
66
const{ Readable, Writable, pipeline }=require('stream');
77

88
constbench=common.createBenchmark(main,{
9-
api: ['classic','webstream','iter','iter-sync'],
9+
api: ['classic','webstream','iter','iter-sync-source','iter-sync'],
1010
datasize: [1024*1024,16*1024*1024,64*1024*1024],
1111
n: [5],
1212
},{
@@ -26,6 +26,8 @@ function main({ api, datasize, n }) {
2626
returnbenchWebStream(chunk,datasize,n,totalOps);
2727
case'iter':
2828
returnbenchIter(chunk,datasize,n,totalOps);
29+
case'iter-sync-source':
30+
returnbenchIterSyncSource(chunk,datasize,n,totalOps);
2931
case'iter-sync':
3032
returnbenchIterSync(chunk,datasize,n,totalOps);
3133
}
@@ -101,6 +103,29 @@ function benchIter(chunk, datasize, n, totalOps) {
101103
})();
102104
}
103105

106+
functionbenchIterSyncSource(chunk,datasize,n,totalOps){
107+
const{ pipeTo }=require('stream/iter');
108+
109+
asyncfunctionrun(){
110+
letremaining=datasize;
111+
function*source(){
112+
while(remaining>0){
113+
constsize=Math.min(remaining,chunk.length);
114+
remaining-=size;
115+
yieldsize===chunk.length ? chunk : chunk.subarray(0,size);
116+
}
117+
}
118+
constwriter={write(){},writeSync(){returntrue;}};
119+
awaitpipeTo(source(),writer);
120+
}
121+
122+
(async()=>{
123+
bench.start();
124+
for(leti=0;i<n;i++)awaitrun();
125+
bench.end(totalOps);
126+
})();
127+
}
128+
104129
functionbenchIterSync(chunk,datasize,n,totalOps){
105130
const{ pipeToSync }=require('stream/iter');
106131

β€Žlib/internal/streams/iter/pull.jsβ€Ž

Lines changed: 53 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,8 @@
88

99
const{
1010
ArrayBufferIsView,
11+
ArrayFromAsync,
12+
ArrayIsArray,
1113
ArrayPrototypePush,
1214
ArrayPrototypeSlice,
1315
PromisePrototypeThen,
@@ -38,7 +40,9 @@ const {
3840
fromSync,
3941
isSyncIterable,
4042
isAsyncIterable,
43+
isPrimitiveChunk,
4144
isUint8ArrayBatch,
45+
normalizeAsyncValue,
4246
}=require('internal/streams/iter/from');
4347

4448
const{
@@ -53,7 +57,10 @@ const {
5357
const{
5458
drainableProtocol,
5559
kSyncWriteAcceptedOnFalse,
60+
kValidatedSource,
5661
kValidatedTransform,
62+
toAsyncStreamable,
63+
toStreamable,
5764
}=require('internal/streams/iter/types');
5865

5966
// =============================================================================
@@ -116,6 +123,22 @@ function parsePipeToArgs(args, requiredMethod) {
116123
};
117124
}
118125

126+
functioncanUseSyncIterablePipeToFastPath(source,transforms,signal){
127+
if(signal!==undefined||
128+
transforms.length!==0||
129+
isPrimitiveChunk(source)||
130+
ArrayIsArray(source)||
131+
source?.[kValidatedSource]||
132+
!isSyncIterable(source)||
133+
isAsyncIterable(source)){
134+
returnfalse;
135+
}
136+
137+
// Preserve from()'s top-level protocol precedence for custom iterables.
138+
returntypeofsource[toAsyncStreamable]!=='function'&&
139+
typeofsource[toStreamable]!=='function';
140+
}
141+
119142
// =============================================================================
120143
// Transform Output Flattening
121144
// =============================================================================
@@ -822,12 +845,13 @@ async function pipeTo(source, ...args) {
822845
// Check for abort
823846
signal?.throwIfAborted();
824847

825-
// Normalize source via from()
826-
constnormalized=from(source);
848+
consthasWriteSync=typeofwriter.writeSync==='function';
849+
constuseSyncIterableFastPath=
850+
hasWriteSync&&canUseSyncIterablePipeToFastPath(source,transforms,signal);
851+
constnormalized=useSyncIterableFastPath ? undefined : from(source);
827852

828853
lettotalBytes=0;
829854
consthasWritev=typeofwriter.writev==='function';
830-
consthasWriteSync=typeofwriter.writeSync==='function';
831855
consthasWritevSync=typeofwriter.writevSync==='function';
832856
consthasEndSync=typeofwriter.endSync==='function';
833857
constsyncFalseCanBeAccepted=writer[kSyncWriteAcceptedOnFalse]===true;
@@ -908,8 +932,32 @@ async function pipeTo(source, ...args) {
908932
}
909933

910934
try{
911-
// Fast path: no transforms - iterate normalized source directly
912-
if(transforms.length===0){
935+
if(useSyncIterableFastPath){
936+
// Avoid from()'s async sync-iterable batching path. This keeps writes
937+
// incremental for synchronous sources while preserving async
938+
// normalization for non-primitive yielded values.
939+
for(constvalueofsource){
940+
if(isUint8ArrayBatch(value)){
941+
if(value.length>0){
942+
constp=writeBatch(value);
943+
if(p)awaitp;
944+
}
945+
continue;
946+
}
947+
if(isUint8Array(value)){
948+
constp=writeBatch([value]);
949+
if(p)awaitp;
950+
continue;
951+
}
952+
953+
constbatch=awaitArrayFromAsync(normalizeAsyncValue(value));
954+
if(batch.length>0){
955+
constp=writeBatch(batch);
956+
if(p)awaitp;
957+
}
958+
}
959+
}elseif(transforms.length===0){
960+
// Fast path: no transforms - iterate normalized source directly
913961
if(signal){
914962
forawait(constbatchofnormalized){
915963
signal.throwIfAborted();

β€Žtest/parallel/test-stream-iter-pipeto.jsβ€Ž

Lines changed: 76 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -219,6 +219,79 @@ async function testPipeToSyncMinimalWriter() {
219219
assert.strictEqual(chunks.length>0,true);
220220
}
221221

222+
asyncfunctiontestPipeToSyncIterableFastPathWritesIncrementally(){
223+
letpulled=0;
224+
letfirstWritePulled=0;
225+
constchunks=[];
226+
function*source(){
227+
for(leti=0;i<3;i++){
228+
pulled++;
229+
yieldnewUint8Array([0x61+i]);
230+
}
231+
}
232+
constwriter={
233+
write: common.mustNotCall(),
234+
writeSync(chunk){
235+
if(firstWritePulled===0){
236+
firstWritePulled=pulled;
237+
}
238+
chunks.push(chunk);
239+
returntrue;
240+
},
241+
};
242+
243+
consttotalBytes=awaitpipeTo(source(),writer);
244+
assert.strictEqual(totalBytes,3);
245+
assert.strictEqual(firstWritePulled,1);
246+
assert.deepStrictEqual(chunks,[
247+
newUint8Array([0x61]),
248+
newUint8Array([0x62]),
249+
newUint8Array([0x63]),
250+
]);
251+
}
252+
253+
asyncfunctiontestPipeToSyncIterableFastPathWriteFallback(){
254+
constasyncWrites=[];
255+
constwriter={
256+
writeSync(chunk){
257+
returnchunk[0]!==0x62;
258+
},
259+
asyncwrite(chunk){
260+
asyncWrites.push(chunk);
261+
},
262+
};
263+
function*source(){
264+
yieldnewUint8Array([0x61]);
265+
yieldnewUint8Array([0x62]);
266+
yieldnewUint8Array([0x63]);
267+
}
268+
269+
consttotalBytes=awaitpipeTo(source(),writer);
270+
assert.strictEqual(totalBytes,3);
271+
assert.deepStrictEqual(asyncWrites,[newUint8Array([0x62])]);
272+
}
273+
274+
asyncfunctiontestPipeToSyncIterableFastPathAsyncValue(){
275+
constchunks=[];
276+
constwriter={
277+
write: common.mustNotCall(),
278+
writeSync(chunk){
279+
chunks.push(chunk);
280+
returntrue;
281+
},
282+
};
283+
function*source(){
284+
yieldPromise.resolve('a');
285+
yieldnewUint8Array([0x62]);
286+
}
287+
288+
consttotalBytes=awaitpipeTo(source(),writer);
289+
assert.strictEqual(totalBytes,2);
290+
constresult=newTextDecoder().decode(
291+
newUint8Array(chunks.reduce((acc,c)=>[...acc, ...c],[])));
292+
assert.strictEqual(result,'ab');
293+
}
294+
222295
Promise.all([
223296
testPipeToSync(),
224297
testPipeTo(),
@@ -234,4 +307,7 @@ Promise.all([
234307
testPipeToSyncPreventClose(),
235308
testPipeToMinimalWriter(),
236309
testPipeToSyncMinimalWriter(),
310+
testPipeToSyncIterableFastPathWritesIncrementally(),
311+
testPipeToSyncIterableFastPathWriteFallback(),
312+
testPipeToSyncIterableFastPathAsyncValue(),
237313
]).then(common.mustCall());

0 commit comments

Comments
Β (0)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content

Commit 02a51a7

Browse files
trivikraduh95
authored andcommitted
stream: add sync iterable fast path to pipeTo
Avoid normalizing sync iterable sources through from() when pipeTo() has no transforms or signal and the writer can accept sync writes. This keeps writes incremental while preserving async fallback for values that still need it. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: openai:gpt-5.5 PR-URL: #63318 Backport-PR-URL: #64675 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Ethan Arrowood <ethan@arrowood.dev>
1 parent 4c9f3ed commit 02a51a7

3 files changed

Lines changed: 155 additions & 6 deletions

File tree

β€Žbenchmark/streams/iter-throughput-pipeto.jsβ€Ž

Lines changed: 26 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ const common = require('../common.js');
66
const{ Readable, Writable, pipeline }=require('stream');
77

88
constbench=common.createBenchmark(main,{
9-
api: ['classic','webstream','iter','iter-sync'],
9+
api: ['classic','webstream','iter','iter-sync-source','iter-sync'],
1010
datasize: [1024*1024,16*1024*1024,64*1024*1024],
1111
n: [5],
1212
},{
@@ -26,6 +26,8 @@ function main({ api, datasize, n }) {
2626
returnbenchWebStream(chunk,datasize,n,totalOps);
2727
case'iter':
2828
returnbenchIter(chunk,datasize,n,totalOps);
29+
case'iter-sync-source':
30+
returnbenchIterSyncSource(chunk,datasize,n,totalOps);
2931
case'iter-sync':
3032
returnbenchIterSync(chunk,datasize,n,totalOps);
3133
}
@@ -101,6 +103,29 @@ function benchIter(chunk, datasize, n, totalOps) {
101103
})();
102104
}
103105

106+
functionbenchIterSyncSource(chunk,datasize,n,totalOps){
107+
const{ pipeTo }=require('stream/iter');
108+
109+
asyncfunctionrun(){
110+
letremaining=datasize;
111+
function*source(){
112+
while(remaining>0){
113+
constsize=Math.min(remaining,chunk.length);
114+
remaining-=size;
115+
yieldsize===chunk.length ? chunk : chunk.subarray(0,size);
116+
}
117+
}
118+
constwriter={write(){},writeSync(){returntrue;}};
119+
awaitpipeTo(source(),writer);
120+
}
121+
122+
(async()=>{
123+
bench.start();
124+
for(leti=0;i<n;i++)awaitrun();
125+
bench.end(totalOps);
126+
})();
127+
}
128+
104129
functionbenchIterSync(chunk,datasize,n,totalOps){
105130
const{ pipeToSync }=require('stream/iter');
106131

β€Žlib/internal/streams/iter/pull.jsβ€Ž

Lines changed: 53 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,8 @@
88

99
const{
1010
ArrayBufferIsView,
11+
ArrayFromAsync,
12+
ArrayIsArray,
1113
ArrayPrototypePush,
1214
ArrayPrototypeSlice,
1315
PromisePrototypeThen,
@@ -38,7 +40,9 @@ const {
3840
fromSync,
3941
isSyncIterable,
4042
isAsyncIterable,
43+
isPrimitiveChunk,
4144
isUint8ArrayBatch,
45+
normalizeAsyncValue,
4246
}=require('internal/streams/iter/from');
4347

4448
const{
@@ -53,7 +57,10 @@ const {
5357
const{
5458
drainableProtocol,
5559
kSyncWriteAcceptedOnFalse,
60+
kValidatedSource,
5661
kValidatedTransform,
62+
toAsyncStreamable,
63+
toStreamable,
5764
}=require('internal/streams/iter/types');
5865

5966
// =============================================================================
@@ -116,6 +123,22 @@ function parsePipeToArgs(args, requiredMethod) {
116123
};
117124
}
118125

126+
functioncanUseSyncIterablePipeToFastPath(source,transforms,signal){
127+
if(signal!==undefined||
128+
transforms.length!==0||
129+
isPrimitiveChunk(source)||
130+
ArrayIsArray(source)||
131+
source?.[kValidatedSource]||
132+
!isSyncIterable(source)||
133+
isAsyncIterable(source)){
134+
returnfalse;
135+
}
136+
137+
// Preserve from()'s top-level protocol precedence for custom iterables.
138+
returntypeofsource[toAsyncStreamable]!=='function'&&
139+
typeofsource[toStreamable]!=='function';
140+
}
141+
119142
// =============================================================================
120143
// Transform Output Flattening
121144
// =============================================================================
@@ -822,12 +845,13 @@ async function pipeTo(source, ...args) {
822845
// Check for abort
823846
signal?.throwIfAborted();
824847

825-
// Normalize source via from()
826-
constnormalized=from(source);
848+
consthasWriteSync=typeofwriter.writeSync==='function';
849+
constuseSyncIterableFastPath=
850+
hasWriteSync&&canUseSyncIterablePipeToFastPath(source,transforms,signal);
851+
constnormalized=useSyncIterableFastPath ? undefined : from(source);
827852

828853
lettotalBytes=0;
829854
consthasWritev=typeofwriter.writev==='function';
830-
consthasWriteSync=typeofwriter.writeSync==='function';
831855
consthasWritevSync=typeofwriter.writevSync==='function';
832856
consthasEndSync=typeofwriter.endSync==='function';
833857
constsyncFalseCanBeAccepted=writer[kSyncWriteAcceptedOnFalse]===true;
@@ -908,8 +932,32 @@ async function pipeTo(source, ...args) {
908932
}
909933

910934
try{
911-
// Fast path: no transforms - iterate normalized source directly
912-
if(transforms.length===0){
935+
if(useSyncIterableFastPath){
936+
// Avoid from()'s async sync-iterable batching path. This keeps writes
937+
// incremental for synchronous sources while preserving async
938+
// normalization for non-primitive yielded values.
939+
for(constvalueofsource){
940+
if(isUint8ArrayBatch(value)){
941+
if(value.length>0){
942+
constp=writeBatch(value);
943+
if(p)awaitp;
944+
}
945+
continue;
946+
}
947+
if(isUint8Array(value)){
948+
constp=writeBatch([value]);
949+
if(p)awaitp;
950+
continue;
951+
}
952+
953+
constbatch=awaitArrayFromAsync(normalizeAsyncValue(value));
954+
if(batch.length>0){
955+
constp=writeBatch(batch);
956+
if(p)awaitp;
957+
}
958+
}
959+
}elseif(transforms.length===0){
960+
// Fast path: no transforms - iterate normalized source directly
913961
if(signal){
914962
forawait(constbatchofnormalized){
915963
signal.throwIfAborted();

β€Žtest/parallel/test-stream-iter-pipeto.jsβ€Ž

Lines changed: 76 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -219,6 +219,79 @@ async function testPipeToSyncMinimalWriter() {
219219
assert.strictEqual(chunks.length>0,true);
220220
}
221221

222+
asyncfunctiontestPipeToSyncIterableFastPathWritesIncrementally(){
223+
letpulled=0;
224+
letfirstWritePulled=0;
225+
constchunks=[];
226+
function*source(){
227+
for(leti=0;i<3;i++){
228+
pulled++;
229+
yieldnewUint8Array([0x61+i]);
230+
}
231+
}
232+
constwriter={
233+
write: common.mustNotCall(),
234+
writeSync(chunk){
235+
if(firstWritePulled===0){
236+
firstWritePulled=pulled;
237+
}
238+
chunks.push(chunk);
239+
returntrue;
240+
},
241+
};
242+
243+
consttotalBytes=awaitpipeTo(source(),writer);
244+
assert.strictEqual(totalBytes,3);
245+
assert.strictEqual(firstWritePulled,1);
246+
assert.deepStrictEqual(chunks,[
247+
newUint8Array([0x61]),
248+
newUint8Array([0x62]),
249+
newUint8Array([0x63]),
250+
]);
251+
}
252+
253+
asyncfunctiontestPipeToSyncIterableFastPathWriteFallback(){
254+
constasyncWrites=[];
255+
constwriter={
256+
writeSync(chunk){
257+
returnchunk[0]!==0x62;
258+
},
259+
asyncwrite(chunk){
260+
asyncWrites.push(chunk);
261+
},
262+
};
263+
function*source(){
264+
yieldnewUint8Array([0x61]);
265+
yieldnewUint8Array([0x62]);
266+
yieldnewUint8Array([0x63]);
267+
}
268+
269+
consttotalBytes=awaitpipeTo(source(),writer);
270+
assert.strictEqual(totalBytes,3);
271+
assert.deepStrictEqual(asyncWrites,[newUint8Array([0x62])]);
272+
}
273+
274+
asyncfunctiontestPipeToSyncIterableFastPathAsyncValue(){
275+
constchunks=[];
276+
constwriter={
277+
write: common.mustNotCall(),
278+
writeSync(chunk){
279+
chunks.push(chunk);
280+
returntrue;
281+
},
282+
};
283+
function*source(){
284+
yieldPromise.resolve('a');
285+
yieldnewUint8Array([0x62]);
286+
}
287+
288+
consttotalBytes=awaitpipeTo(source(),writer);
289+
assert.strictEqual(totalBytes,2);
290+
constresult=newTextDecoder().decode(
291+
newUint8Array(chunks.reduce((acc,c)=>[...acc, ...c],[])));
292+
assert.strictEqual(result,'ab');
293+
}
294+
222295
Promise.all([
223296
testPipeToSync(),
224297
testPipeTo(),
@@ -234,4 +307,7 @@ Promise.all([
234307
testPipeToSyncPreventClose(),
235308
testPipeToMinimalWriter(),
236309
testPipeToSyncMinimalWriter(),
310+
testPipeToSyncIterableFastPathWritesIncrementally(),
311+
testPipeToSyncIterableFastPathWriteFallback(),
312+
testPipeToSyncIterableFastPathAsyncValue(),
237313
]).then(common.mustCall());

0 commit comments

Comments
Β (0)