Commit ec2666e

Browse files
trivikraduh95
authored andcommitted
stream: avoid retrying accepted pipeTo writes
PushWriter in block backpressure mode can return false from writeSync() and writevSync() after accepting data. Treat that false return as backpressure and wait for drain instead of retrying the same chunks asynchronously. Fixes: #63296 Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: openai:gpt-5.5 PR-URL: #63297 Backport-PR-URL: #64675Fixes: #63296 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Ethan Arrowood <ethan@arrowood.dev>
1 parent f63143f commit ec2666e

4 files changed

Lines changed: 76 additions & 3 deletions

File tree

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

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,8 @@ const {
5151
}=require('internal/streams/iter/utils');
5252

5353
const{
54+
drainableProtocol,
55+
kSyncWriteAcceptedOnFalse,
5456
kValidatedTransform,
5557
}=require('internal/streams/iter/types');
5658

@@ -828,13 +830,33 @@ async function pipeTo(source, ...args) {
828830
consthasWriteSync=typeofwriter.writeSync==='function';
829831
consthasWritevSync=typeofwriter.writevSync==='function';
830832
consthasEndSync=typeofwriter.endSync==='function';
833+
constsyncFalseCanBeAccepted=writer[kSyncWriteAcceptedOnFalse]===true;
834+
835+
functionsyncFalseWasAccepted(){
836+
returnsyncFalseCanBeAccepted&&writer.desiredSize===0;
837+
}
838+
839+
functionwaitForSyncBackpressure(){
840+
constondrain=writer[drainableProtocol];
841+
returnondrain?.call(writer);
842+
}
843+
844+
asyncfunctionwriteBatchAfterAcceptedBackpressure(batch,startIndex){
845+
awaitwaitForSyncBackpressure();
846+
awaitwriteBatchAsyncFallback(batch,startIndex);
847+
}
848+
831849
// Async fallback for writeBatch when sync write fails partway through.
832850
// Continues writing from batch[startIndex] using async write().
833851
asyncfunctionwriteBatchAsyncFallback(batch,startIndex){
834852
for(leti=startIndex;i<batch.length;i++){
835853
constchunk=batch[i];
836854
if(hasWriteSync&&writer.writeSync(chunk)){
837855
// Sync retry succeeded
856+
}elseif(syncFalseWasAccepted()){
857+
totalBytes+=TypedArrayPrototypeGetByteLength(chunk);
858+
awaitwaitForSyncBackpressure();
859+
continue;
838860
}else{
839861
constresult=writer.write(
840862
chunk,signal ? {__proto__: null, signal } : undefined);
@@ -852,6 +874,12 @@ async function pipeTo(source, ...args) {
852874
functionwriteBatch(batch){
853875
if(hasWritev&&batch.length>1){
854876
if(!hasWritevSync||!writer.writevSync(batch)){
877+
if(hasWritevSync&&syncFalseWasAccepted()){
878+
for(leti=0;i<batch.length;i++){
879+
totalBytes+=TypedArrayPrototypeGetByteLength(batch[i]);
880+
}
881+
returnwaitForSyncBackpressure();
882+
}
855883
constopts=signal ? {__proto__: null, signal } : undefined;
856884
returnPromisePrototypeThen(writer.writev(batch,opts),()=>{
857885
for(leti=0;i<batch.length;i++){
@@ -867,6 +895,10 @@ async function pipeTo(source, ...args) {
867895
for(leti=0;i<batch.length;i++){
868896
constchunk=batch[i];
869897
if(!hasWriteSync||!writer.writeSync(chunk)){
898+
if(hasWriteSync&&syncFalseWasAccepted()){
899+
totalBytes+=TypedArrayPrototypeGetByteLength(chunk);
900+
returnwriteBatchAfterAcceptedBackpressure(batch,i+1);
901+
}
870902
// Sync path failed at index i - fall back to async for the rest.
871903
// Count bytes for chunks already written synchronously (0..i-1).
872904
returnwriteBatchAsyncFallback(batch,i);

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

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@ const {
3232

3333
const{
3434
drainableProtocol,
35+
kSyncWriteAcceptedOnFalse,
3536
}=require('internal/streams/iter/types');
3637

3738
const{
@@ -560,6 +561,10 @@ class PushWriter {
560561
returnthis.#queue.desiredSize;
561562
}
562563

564+
get[kSyncWriteAcceptedOnFalse](){
565+
returnthis.#queue.backpressurePolicy==='block';
566+
}
567+
563568
write(chunk,options){
564569
if(!options?.signal&&this.#queue.canWriteSync()){
565570
constbytes=toUint8Array(chunk);
@@ -586,7 +591,8 @@ class PushWriter {
586591
writeSync(chunk){
587592
constbytes=toUint8Array(chunk);
588593
constresult=this.#queue.writeSync([bytes]);
589-
if(!result&&this.#queue.backpressurePolicy==='block'){
594+
if(!result&&this.#queue.backpressurePolicy==='block'&&
595+
this.#queue.desiredSize===0){
590596
// Block policy: force-enqueue and return false as backpressure signal.
591597
// Data IS accepted; false tells caller to slow down.
592598
this.#queue.forceEnqueue([bytes]);
@@ -601,7 +607,8 @@ class PushWriter {
601607
}
602608
constbytes=convertChunks(chunks);
603609
constresult=this.#queue.writeSync(bytes);
604-
if(!result&&this.#queue.backpressurePolicy==='block'){
610+
if(!result&&this.#queue.backpressurePolicy==='block'&&
611+
this.#queue.desiredSize===0){
605612
this.#queue.forceEnqueue(bytes);
606613
returnfalse;
607614
}

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,9 +64,12 @@ const kValidatedTransform = Symbol('kValidatedTransform');
6464
*/
6565
constkValidatedSource=Symbol('kValidatedSource');
6666

67+
constkSyncWriteAcceptedOnFalse=Symbol('kSyncWriteAcceptedOnFalse');
68+
6769
module.exports={
6870
broadcastProtocol,
6971
drainableProtocol,
72+
kSyncWriteAcceptedOnFalse,
7073
kValidatedSource,
7174
kValidatedTransform,
7275
shareProtocol,

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

Lines changed: 32 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,8 @@
55

66
constcommon=require('../common');
77
constassert=require('assert');
8-
const{ pipeTo, pipeToSync }=require('stream/iter');
8+
const{setImmediate: setImmediatePromise}=require('timers/promises');
9+
const{ pipeTo, pipeToSync, push, text }=require('stream/iter');
910

1011
// Multi-chunk batch with writevSync (sync success path)
1112
asyncfunctiontestWritevSyncSuccess(){
@@ -104,6 +105,35 @@ async function testWriteSyncAlwaysFails() {
104105
assert.strictEqual(total,2);
105106
}
106107

108+
// PushWriter block mode accepts sync writes even when returning false for
109+
// backpressure. pipeTo must wait for drain, not retry the same write.
110+
asyncfunctionassertPushWriterBlockPipeTo(source,expected,expectedTotal){
111+
const{ writer, readable }=push({
112+
highWaterMark: 1,
113+
backpressure: 'block',
114+
});
115+
116+
constpipe=pipeTo(source,writer);
117+
awaitsetImmediatePromise();
118+
constdata=awaittext(readable);
119+
consttotal=awaitpipe;
120+
121+
assert.strictEqual(data,expected);
122+
assert.strictEqual(total,expectedTotal);
123+
}
124+
125+
asyncfunctiontestPushWriterBlockSyncFalseAccepted(){
126+
awaitassertPushWriterBlockPipeTo((asyncfunction*(){
127+
yield[newUint8Array([97])];
128+
yield[newUint8Array([98])];
129+
})(),'ab',2);
130+
131+
awaitassertPushWriterBlockPipeTo((asyncfunction*(){
132+
yield[newUint8Array([97,98])];
133+
yield[newUint8Array([99]),newUint8Array([100])];
134+
})(),'abcd',4);
135+
}
136+
107137
// pipeToSync with writevSync
108138
asyncfunctiontestPipeToSyncWritev(){
109139
constbatches=[];
@@ -142,6 +172,7 @@ Promise.all([
142172
testWritevSyncFails(),
143173
testWriteSyncFailsMidBatch(),
144174
testWriteSyncAlwaysFails(),
175+
testPushWriterBlockSyncFalseAccepted(),
145176
testPipeToSyncWritev(),
146177
testPipeToSyncWriteFallback(),
147178
]).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 ec2666e

Browse files
trivikraduh95
authored andcommitted
stream: avoid retrying accepted pipeTo writes
PushWriter in block backpressure mode can return false from writeSync() and writevSync() after accepting data. Treat that false return as backpressure and wait for drain instead of retrying the same chunks asynchronously. Fixes: #63296 Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: openai:gpt-5.5 PR-URL: #63297 Backport-PR-URL: #64675Fixes: #63296 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Ethan Arrowood <ethan@arrowood.dev>
1 parent f63143f commit ec2666e

4 files changed

Lines changed: 76 additions & 3 deletions

File tree

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

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,8 @@ const {
5151
}=require('internal/streams/iter/utils');
5252

5353
const{
54+
drainableProtocol,
55+
kSyncWriteAcceptedOnFalse,
5456
kValidatedTransform,
5557
}=require('internal/streams/iter/types');
5658

@@ -828,13 +830,33 @@ async function pipeTo(source, ...args) {
828830
consthasWriteSync=typeofwriter.writeSync==='function';
829831
consthasWritevSync=typeofwriter.writevSync==='function';
830832
consthasEndSync=typeofwriter.endSync==='function';
833+
constsyncFalseCanBeAccepted=writer[kSyncWriteAcceptedOnFalse]===true;
834+
835+
functionsyncFalseWasAccepted(){
836+
returnsyncFalseCanBeAccepted&&writer.desiredSize===0;
837+
}
838+
839+
functionwaitForSyncBackpressure(){
840+
constondrain=writer[drainableProtocol];
841+
returnondrain?.call(writer);
842+
}
843+
844+
asyncfunctionwriteBatchAfterAcceptedBackpressure(batch,startIndex){
845+
awaitwaitForSyncBackpressure();
846+
awaitwriteBatchAsyncFallback(batch,startIndex);
847+
}
848+
831849
// Async fallback for writeBatch when sync write fails partway through.
832850
// Continues writing from batch[startIndex] using async write().
833851
asyncfunctionwriteBatchAsyncFallback(batch,startIndex){
834852
for(leti=startIndex;i<batch.length;i++){
835853
constchunk=batch[i];
836854
if(hasWriteSync&&writer.writeSync(chunk)){
837855
// Sync retry succeeded
856+
}elseif(syncFalseWasAccepted()){
857+
totalBytes+=TypedArrayPrototypeGetByteLength(chunk);
858+
awaitwaitForSyncBackpressure();
859+
continue;
838860
}else{
839861
constresult=writer.write(
840862
chunk,signal ? {__proto__: null, signal } : undefined);
@@ -852,6 +874,12 @@ async function pipeTo(source, ...args) {
852874
functionwriteBatch(batch){
853875
if(hasWritev&&batch.length>1){
854876
if(!hasWritevSync||!writer.writevSync(batch)){
877+
if(hasWritevSync&&syncFalseWasAccepted()){
878+
for(leti=0;i<batch.length;i++){
879+
totalBytes+=TypedArrayPrototypeGetByteLength(batch[i]);
880+
}
881+
returnwaitForSyncBackpressure();
882+
}
855883
constopts=signal ? {__proto__: null, signal } : undefined;
856884
returnPromisePrototypeThen(writer.writev(batch,opts),()=>{
857885
for(leti=0;i<batch.length;i++){
@@ -867,6 +895,10 @@ async function pipeTo(source, ...args) {
867895
for(leti=0;i<batch.length;i++){
868896
constchunk=batch[i];
869897
if(!hasWriteSync||!writer.writeSync(chunk)){
898+
if(hasWriteSync&&syncFalseWasAccepted()){
899+
totalBytes+=TypedArrayPrototypeGetByteLength(chunk);
900+
returnwriteBatchAfterAcceptedBackpressure(batch,i+1);
901+
}
870902
// Sync path failed at index i - fall back to async for the rest.
871903
// Count bytes for chunks already written synchronously (0..i-1).
872904
returnwriteBatchAsyncFallback(batch,i);

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

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@ const {
3232

3333
const{
3434
drainableProtocol,
35+
kSyncWriteAcceptedOnFalse,
3536
}=require('internal/streams/iter/types');
3637

3738
const{
@@ -560,6 +561,10 @@ class PushWriter {
560561
returnthis.#queue.desiredSize;
561562
}
562563

564+
get[kSyncWriteAcceptedOnFalse](){
565+
returnthis.#queue.backpressurePolicy==='block';
566+
}
567+
563568
write(chunk,options){
564569
if(!options?.signal&&this.#queue.canWriteSync()){
565570
constbytes=toUint8Array(chunk);
@@ -586,7 +591,8 @@ class PushWriter {
586591
writeSync(chunk){
587592
constbytes=toUint8Array(chunk);
588593
constresult=this.#queue.writeSync([bytes]);
589-
if(!result&&this.#queue.backpressurePolicy==='block'){
594+
if(!result&&this.#queue.backpressurePolicy==='block'&&
595+
this.#queue.desiredSize===0){
590596
// Block policy: force-enqueue and return false as backpressure signal.
591597
// Data IS accepted; false tells caller to slow down.
592598
this.#queue.forceEnqueue([bytes]);
@@ -601,7 +607,8 @@ class PushWriter {
601607
}
602608
constbytes=convertChunks(chunks);
603609
constresult=this.#queue.writeSync(bytes);
604-
if(!result&&this.#queue.backpressurePolicy==='block'){
610+
if(!result&&this.#queue.backpressurePolicy==='block'&&
611+
this.#queue.desiredSize===0){
605612
this.#queue.forceEnqueue(bytes);
606613
returnfalse;
607614
}

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,9 +64,12 @@ const kValidatedTransform = Symbol('kValidatedTransform');
6464
*/
6565
constkValidatedSource=Symbol('kValidatedSource');
6666

67+
constkSyncWriteAcceptedOnFalse=Symbol('kSyncWriteAcceptedOnFalse');
68+
6769
module.exports={
6870
broadcastProtocol,
6971
drainableProtocol,
72+
kSyncWriteAcceptedOnFalse,
7073
kValidatedSource,
7174
kValidatedTransform,
7275
shareProtocol,

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

Lines changed: 32 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,8 @@
55

66
constcommon=require('../common');
77
constassert=require('assert');
8-
const{ pipeTo, pipeToSync }=require('stream/iter');
8+
const{setImmediate: setImmediatePromise}=require('timers/promises');
9+
const{ pipeTo, pipeToSync, push, text }=require('stream/iter');
910

1011
// Multi-chunk batch with writevSync (sync success path)
1112
asyncfunctiontestWritevSyncSuccess(){
@@ -104,6 +105,35 @@ async function testWriteSyncAlwaysFails() {
104105
assert.strictEqual(total,2);
105106
}
106107

108+
// PushWriter block mode accepts sync writes even when returning false for
109+
// backpressure. pipeTo must wait for drain, not retry the same write.
110+
asyncfunctionassertPushWriterBlockPipeTo(source,expected,expectedTotal){
111+
const{ writer, readable }=push({
112+
highWaterMark: 1,
113+
backpressure: 'block',
114+
});
115+
116+
constpipe=pipeTo(source,writer);
117+
awaitsetImmediatePromise();
118+
constdata=awaittext(readable);
119+
consttotal=awaitpipe;
120+
121+
assert.strictEqual(data,expected);
122+
assert.strictEqual(total,expectedTotal);
123+
}
124+
125+
asyncfunctiontestPushWriterBlockSyncFalseAccepted(){
126+
awaitassertPushWriterBlockPipeTo((asyncfunction*(){
127+
yield[newUint8Array([97])];
128+
yield[newUint8Array([98])];
129+
})(),'ab',2);
130+
131+
awaitassertPushWriterBlockPipeTo((asyncfunction*(){
132+
yield[newUint8Array([97,98])];
133+
yield[newUint8Array([99]),newUint8Array([100])];
134+
})(),'abcd',4);
135+
}
136+
107137
// pipeToSync with writevSync
108138
asyncfunctiontestPipeToSyncWritev(){
109139
constbatches=[];
@@ -142,6 +172,7 @@ Promise.all([
142172
testWritevSyncFails(),
143173
testWriteSyncFailsMidBatch(),
144174
testWriteSyncAlwaysFails(),
175+
testPushWriterBlockSyncFalseAccepted(),
145176
testPipeToSyncWritev(),
146177
testPipeToSyncWriteFallback(),
147178
]).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 ec2666e

Browse files
trivikraduh95
authored andcommitted
stream: avoid retrying accepted pipeTo writes
PushWriter in block backpressure mode can return false from writeSync() and writevSync() after accepting data. Treat that false return as backpressure and wait for drain instead of retrying the same chunks asynchronously. Fixes: #63296 Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: openai:gpt-5.5 PR-URL: #63297 Backport-PR-URL: #64675Fixes: #63296 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Ethan Arrowood <ethan@arrowood.dev>
1 parent f63143f commit ec2666e

4 files changed

Lines changed: 76 additions & 3 deletions

File tree

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

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,8 @@ const {
5151
}=require('internal/streams/iter/utils');
5252

5353
const{
54+
drainableProtocol,
55+
kSyncWriteAcceptedOnFalse,
5456
kValidatedTransform,
5557
}=require('internal/streams/iter/types');
5658

@@ -828,13 +830,33 @@ async function pipeTo(source, ...args) {
828830
consthasWriteSync=typeofwriter.writeSync==='function';
829831
consthasWritevSync=typeofwriter.writevSync==='function';
830832
consthasEndSync=typeofwriter.endSync==='function';
833+
constsyncFalseCanBeAccepted=writer[kSyncWriteAcceptedOnFalse]===true;
834+
835+
functionsyncFalseWasAccepted(){
836+
returnsyncFalseCanBeAccepted&&writer.desiredSize===0;
837+
}
838+
839+
functionwaitForSyncBackpressure(){
840+
constondrain=writer[drainableProtocol];
841+
returnondrain?.call(writer);
842+
}
843+
844+
asyncfunctionwriteBatchAfterAcceptedBackpressure(batch,startIndex){
845+
awaitwaitForSyncBackpressure();
846+
awaitwriteBatchAsyncFallback(batch,startIndex);
847+
}
848+
831849
// Async fallback for writeBatch when sync write fails partway through.
832850
// Continues writing from batch[startIndex] using async write().
833851
asyncfunctionwriteBatchAsyncFallback(batch,startIndex){
834852
for(leti=startIndex;i<batch.length;i++){
835853
constchunk=batch[i];
836854
if(hasWriteSync&&writer.writeSync(chunk)){
837855
// Sync retry succeeded
856+
}elseif(syncFalseWasAccepted()){
857+
totalBytes+=TypedArrayPrototypeGetByteLength(chunk);
858+
awaitwaitForSyncBackpressure();
859+
continue;
838860
}else{
839861
constresult=writer.write(
840862
chunk,signal ? {__proto__: null, signal } : undefined);
@@ -852,6 +874,12 @@ async function pipeTo(source, ...args) {
852874
functionwriteBatch(batch){
853875
if(hasWritev&&batch.length>1){
854876
if(!hasWritevSync||!writer.writevSync(batch)){
877+
if(hasWritevSync&&syncFalseWasAccepted()){
878+
for(leti=0;i<batch.length;i++){
879+
totalBytes+=TypedArrayPrototypeGetByteLength(batch[i]);
880+
}
881+
returnwaitForSyncBackpressure();
882+
}
855883
constopts=signal ? {__proto__: null, signal } : undefined;
856884
returnPromisePrototypeThen(writer.writev(batch,opts),()=>{
857885
for(leti=0;i<batch.length;i++){
@@ -867,6 +895,10 @@ async function pipeTo(source, ...args) {
867895
for(leti=0;i<batch.length;i++){
868896
constchunk=batch[i];
869897
if(!hasWriteSync||!writer.writeSync(chunk)){
898+
if(hasWriteSync&&syncFalseWasAccepted()){
899+
totalBytes+=TypedArrayPrototypeGetByteLength(chunk);
900+
returnwriteBatchAfterAcceptedBackpressure(batch,i+1);
901+
}
870902
// Sync path failed at index i - fall back to async for the rest.
871903
// Count bytes for chunks already written synchronously (0..i-1).
872904
returnwriteBatchAsyncFallback(batch,i);

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

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@ const {
3232

3333
const{
3434
drainableProtocol,
35+
kSyncWriteAcceptedOnFalse,
3536
}=require('internal/streams/iter/types');
3637

3738
const{
@@ -560,6 +561,10 @@ class PushWriter {
560561
returnthis.#queue.desiredSize;
561562
}
562563

564+
get[kSyncWriteAcceptedOnFalse](){
565+
returnthis.#queue.backpressurePolicy==='block';
566+
}
567+
563568
write(chunk,options){
564569
if(!options?.signal&&this.#queue.canWriteSync()){
565570
constbytes=toUint8Array(chunk);
@@ -586,7 +591,8 @@ class PushWriter {
586591
writeSync(chunk){
587592
constbytes=toUint8Array(chunk);
588593
constresult=this.#queue.writeSync([bytes]);
589-
if(!result&&this.#queue.backpressurePolicy==='block'){
594+
if(!result&&this.#queue.backpressurePolicy==='block'&&
595+
this.#queue.desiredSize===0){
590596
// Block policy: force-enqueue and return false as backpressure signal.
591597
// Data IS accepted; false tells caller to slow down.
592598
this.#queue.forceEnqueue([bytes]);
@@ -601,7 +607,8 @@ class PushWriter {
601607
}
602608
constbytes=convertChunks(chunks);
603609
constresult=this.#queue.writeSync(bytes);
604-
if(!result&&this.#queue.backpressurePolicy==='block'){
610+
if(!result&&this.#queue.backpressurePolicy==='block'&&
611+
this.#queue.desiredSize===0){
605612
this.#queue.forceEnqueue(bytes);
606613
returnfalse;
607614
}

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,9 +64,12 @@ const kValidatedTransform = Symbol('kValidatedTransform');
6464
*/
6565
constkValidatedSource=Symbol('kValidatedSource');
6666

67+
constkSyncWriteAcceptedOnFalse=Symbol('kSyncWriteAcceptedOnFalse');
68+
6769
module.exports={
6870
broadcastProtocol,
6971
drainableProtocol,
72+
kSyncWriteAcceptedOnFalse,
7073
kValidatedSource,
7174
kValidatedTransform,
7275
shareProtocol,

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

Lines changed: 32 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,8 @@
55

66
constcommon=require('../common');
77
constassert=require('assert');
8-
const{ pipeTo, pipeToSync }=require('stream/iter');
8+
const{setImmediate: setImmediatePromise}=require('timers/promises');
9+
const{ pipeTo, pipeToSync, push, text }=require('stream/iter');
910

1011
// Multi-chunk batch with writevSync (sync success path)
1112
asyncfunctiontestWritevSyncSuccess(){
@@ -104,6 +105,35 @@ async function testWriteSyncAlwaysFails() {
104105
assert.strictEqual(total,2);
105106
}
106107

108+
// PushWriter block mode accepts sync writes even when returning false for
109+
// backpressure. pipeTo must wait for drain, not retry the same write.
110+
asyncfunctionassertPushWriterBlockPipeTo(source,expected,expectedTotal){
111+
const{ writer, readable }=push({
112+
highWaterMark: 1,
113+
backpressure: 'block',
114+
});
115+
116+
constpipe=pipeTo(source,writer);
117+
awaitsetImmediatePromise();
118+
constdata=awaittext(readable);
119+
consttotal=awaitpipe;
120+
121+
assert.strictEqual(data,expected);
122+
assert.strictEqual(total,expectedTotal);
123+
}
124+
125+
asyncfunctiontestPushWriterBlockSyncFalseAccepted(){
126+
awaitassertPushWriterBlockPipeTo((asyncfunction*(){
127+
yield[newUint8Array([97])];
128+
yield[newUint8Array([98])];
129+
})(),'ab',2);
130+
131+
awaitassertPushWriterBlockPipeTo((asyncfunction*(){
132+
yield[newUint8Array([97,98])];
133+
yield[newUint8Array([99]),newUint8Array([100])];
134+
})(),'abcd',4);
135+
}
136+
107137
// pipeToSync with writevSync
108138
asyncfunctiontestPipeToSyncWritev(){
109139
constbatches=[];
@@ -142,6 +172,7 @@ Promise.all([
142172
testWritevSyncFails(),
143173
testWriteSyncFailsMidBatch(),
144174
testWriteSyncAlwaysFails(),
175+
testPushWriterBlockSyncFalseAccepted(),
145176
testPipeToSyncWritev(),
146177
testPipeToSyncWriteFallback(),
147178
]).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 ec2666e

Browse files
trivikraduh95
authored andcommitted
stream: avoid retrying accepted pipeTo writes
PushWriter in block backpressure mode can return false from writeSync() and writevSync() after accepting data. Treat that false return as backpressure and wait for drain instead of retrying the same chunks asynchronously. Fixes: #63296 Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: openai:gpt-5.5 PR-URL: #63297 Backport-PR-URL: #64675Fixes: #63296 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Ethan Arrowood <ethan@arrowood.dev>
1 parent f63143f commit ec2666e

4 files changed

Lines changed: 76 additions & 3 deletions

File tree

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

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,8 @@ const {
5151
}=require('internal/streams/iter/utils');
5252

5353
const{
54+
drainableProtocol,
55+
kSyncWriteAcceptedOnFalse,
5456
kValidatedTransform,
5557
}=require('internal/streams/iter/types');
5658

@@ -828,13 +830,33 @@ async function pipeTo(source, ...args) {
828830
consthasWriteSync=typeofwriter.writeSync==='function';
829831
consthasWritevSync=typeofwriter.writevSync==='function';
830832
consthasEndSync=typeofwriter.endSync==='function';
833+
constsyncFalseCanBeAccepted=writer[kSyncWriteAcceptedOnFalse]===true;
834+
835+
functionsyncFalseWasAccepted(){
836+
returnsyncFalseCanBeAccepted&&writer.desiredSize===0;
837+
}
838+
839+
functionwaitForSyncBackpressure(){
840+
constondrain=writer[drainableProtocol];
841+
returnondrain?.call(writer);
842+
}
843+
844+
asyncfunctionwriteBatchAfterAcceptedBackpressure(batch,startIndex){
845+
awaitwaitForSyncBackpressure();
846+
awaitwriteBatchAsyncFallback(batch,startIndex);
847+
}
848+
831849
// Async fallback for writeBatch when sync write fails partway through.
832850
// Continues writing from batch[startIndex] using async write().
833851
asyncfunctionwriteBatchAsyncFallback(batch,startIndex){
834852
for(leti=startIndex;i<batch.length;i++){
835853
constchunk=batch[i];
836854
if(hasWriteSync&&writer.writeSync(chunk)){
837855
// Sync retry succeeded
856+
}elseif(syncFalseWasAccepted()){
857+
totalBytes+=TypedArrayPrototypeGetByteLength(chunk);
858+
awaitwaitForSyncBackpressure();
859+
continue;
838860
}else{
839861
constresult=writer.write(
840862
chunk,signal ? {__proto__: null, signal } : undefined);
@@ -852,6 +874,12 @@ async function pipeTo(source, ...args) {
852874
functionwriteBatch(batch){
853875
if(hasWritev&&batch.length>1){
854876
if(!hasWritevSync||!writer.writevSync(batch)){
877+
if(hasWritevSync&&syncFalseWasAccepted()){
878+
for(leti=0;i<batch.length;i++){
879+
totalBytes+=TypedArrayPrototypeGetByteLength(batch[i]);
880+
}
881+
returnwaitForSyncBackpressure();
882+
}
855883
constopts=signal ? {__proto__: null, signal } : undefined;
856884
returnPromisePrototypeThen(writer.writev(batch,opts),()=>{
857885
for(leti=0;i<batch.length;i++){
@@ -867,6 +895,10 @@ async function pipeTo(source, ...args) {
867895
for(leti=0;i<batch.length;i++){
868896
constchunk=batch[i];
869897
if(!hasWriteSync||!writer.writeSync(chunk)){
898+
if(hasWriteSync&&syncFalseWasAccepted()){
899+
totalBytes+=TypedArrayPrototypeGetByteLength(chunk);
900+
returnwriteBatchAfterAcceptedBackpressure(batch,i+1);
901+
}
870902
// Sync path failed at index i - fall back to async for the rest.
871903
// Count bytes for chunks already written synchronously (0..i-1).
872904
returnwriteBatchAsyncFallback(batch,i);

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

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@ const {
3232

3333
const{
3434
drainableProtocol,
35+
kSyncWriteAcceptedOnFalse,
3536
}=require('internal/streams/iter/types');
3637

3738
const{
@@ -560,6 +561,10 @@ class PushWriter {
560561
returnthis.#queue.desiredSize;
561562
}
562563

564+
get[kSyncWriteAcceptedOnFalse](){
565+
returnthis.#queue.backpressurePolicy==='block';
566+
}
567+
563568
write(chunk,options){
564569
if(!options?.signal&&this.#queue.canWriteSync()){
565570
constbytes=toUint8Array(chunk);
@@ -586,7 +591,8 @@ class PushWriter {
586591
writeSync(chunk){
587592
constbytes=toUint8Array(chunk);
588593
constresult=this.#queue.writeSync([bytes]);
589-
if(!result&&this.#queue.backpressurePolicy==='block'){
594+
if(!result&&this.#queue.backpressurePolicy==='block'&&
595+
this.#queue.desiredSize===0){
590596
// Block policy: force-enqueue and return false as backpressure signal.
591597
// Data IS accepted; false tells caller to slow down.
592598
this.#queue.forceEnqueue([bytes]);
@@ -601,7 +607,8 @@ class PushWriter {
601607
}
602608
constbytes=convertChunks(chunks);
603609
constresult=this.#queue.writeSync(bytes);
604-
if(!result&&this.#queue.backpressurePolicy==='block'){
610+
if(!result&&this.#queue.backpressurePolicy==='block'&&
611+
this.#queue.desiredSize===0){
605612
this.#queue.forceEnqueue(bytes);
606613
returnfalse;
607614
}

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,9 +64,12 @@ const kValidatedTransform = Symbol('kValidatedTransform');
6464
*/
6565
constkValidatedSource=Symbol('kValidatedSource');
6666

67+
constkSyncWriteAcceptedOnFalse=Symbol('kSyncWriteAcceptedOnFalse');
68+
6769
module.exports={
6870
broadcastProtocol,
6971
drainableProtocol,
72+
kSyncWriteAcceptedOnFalse,
7073
kValidatedSource,
7174
kValidatedTransform,
7275
shareProtocol,

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

Lines changed: 32 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,8 @@
55

66
constcommon=require('../common');
77
constassert=require('assert');
8-
const{ pipeTo, pipeToSync }=require('stream/iter');
8+
const{setImmediate: setImmediatePromise}=require('timers/promises');
9+
const{ pipeTo, pipeToSync, push, text }=require('stream/iter');
910

1011
// Multi-chunk batch with writevSync (sync success path)
1112
asyncfunctiontestWritevSyncSuccess(){
@@ -104,6 +105,35 @@ async function testWriteSyncAlwaysFails() {
104105
assert.strictEqual(total,2);
105106
}
106107

108+
// PushWriter block mode accepts sync writes even when returning false for
109+
// backpressure. pipeTo must wait for drain, not retry the same write.
110+
asyncfunctionassertPushWriterBlockPipeTo(source,expected,expectedTotal){
111+
const{ writer, readable }=push({
112+
highWaterMark: 1,
113+
backpressure: 'block',
114+
});
115+
116+
constpipe=pipeTo(source,writer);
117+
awaitsetImmediatePromise();
118+
constdata=awaittext(readable);
119+
consttotal=awaitpipe;
120+
121+
assert.strictEqual(data,expected);
122+
assert.strictEqual(total,expectedTotal);
123+
}
124+
125+
asyncfunctiontestPushWriterBlockSyncFalseAccepted(){
126+
awaitassertPushWriterBlockPipeTo((asyncfunction*(){
127+
yield[newUint8Array([97])];
128+
yield[newUint8Array([98])];
129+
})(),'ab',2);
130+
131+
awaitassertPushWriterBlockPipeTo((asyncfunction*(){
132+
yield[newUint8Array([97,98])];
133+
yield[newUint8Array([99]),newUint8Array([100])];
134+
})(),'abcd',4);
135+
}
136+
107137
// pipeToSync with writevSync
108138
asyncfunctiontestPipeToSyncWritev(){
109139
constbatches=[];
@@ -142,6 +172,7 @@ Promise.all([
142172
testWritevSyncFails(),
143173
testWriteSyncFailsMidBatch(),
144174
testWriteSyncAlwaysFails(),
175+
testPushWriterBlockSyncFalseAccepted(),
145176
testPipeToSyncWritev(),
146177
testPipeToSyncWriteFallback(),
147178
]).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 ec2666e

Browse files
trivikraduh95
authored andcommitted
stream: avoid retrying accepted pipeTo writes
PushWriter in block backpressure mode can return false from writeSync() and writevSync() after accepting data. Treat that false return as backpressure and wait for drain instead of retrying the same chunks asynchronously. Fixes: #63296 Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: openai:gpt-5.5 PR-URL: #63297 Backport-PR-URL: #64675Fixes: #63296 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Ethan Arrowood <ethan@arrowood.dev>
1 parent f63143f commit ec2666e

4 files changed

Lines changed: 76 additions & 3 deletions

File tree

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

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,8 @@ const {
5151
}=require('internal/streams/iter/utils');
5252

5353
const{
54+
drainableProtocol,
55+
kSyncWriteAcceptedOnFalse,
5456
kValidatedTransform,
5557
}=require('internal/streams/iter/types');
5658

@@ -828,13 +830,33 @@ async function pipeTo(source, ...args) {
828830
consthasWriteSync=typeofwriter.writeSync==='function';
829831
consthasWritevSync=typeofwriter.writevSync==='function';
830832
consthasEndSync=typeofwriter.endSync==='function';
833+
constsyncFalseCanBeAccepted=writer[kSyncWriteAcceptedOnFalse]===true;
834+
835+
functionsyncFalseWasAccepted(){
836+
returnsyncFalseCanBeAccepted&&writer.desiredSize===0;
837+
}
838+
839+
functionwaitForSyncBackpressure(){
840+
constondrain=writer[drainableProtocol];
841+
returnondrain?.call(writer);
842+
}
843+
844+
asyncfunctionwriteBatchAfterAcceptedBackpressure(batch,startIndex){
845+
awaitwaitForSyncBackpressure();
846+
awaitwriteBatchAsyncFallback(batch,startIndex);
847+
}
848+
831849
// Async fallback for writeBatch when sync write fails partway through.
832850
// Continues writing from batch[startIndex] using async write().
833851
asyncfunctionwriteBatchAsyncFallback(batch,startIndex){
834852
for(leti=startIndex;i<batch.length;i++){
835853
constchunk=batch[i];
836854
if(hasWriteSync&&writer.writeSync(chunk)){
837855
// Sync retry succeeded
856+
}elseif(syncFalseWasAccepted()){
857+
totalBytes+=TypedArrayPrototypeGetByteLength(chunk);
858+
awaitwaitForSyncBackpressure();
859+
continue;
838860
}else{
839861
constresult=writer.write(
840862
chunk,signal ? {__proto__: null, signal } : undefined);
@@ -852,6 +874,12 @@ async function pipeTo(source, ...args) {
852874
functionwriteBatch(batch){
853875
if(hasWritev&&batch.length>1){
854876
if(!hasWritevSync||!writer.writevSync(batch)){
877+
if(hasWritevSync&&syncFalseWasAccepted()){
878+
for(leti=0;i<batch.length;i++){
879+
totalBytes+=TypedArrayPrototypeGetByteLength(batch[i]);
880+
}
881+
returnwaitForSyncBackpressure();
882+
}
855883
constopts=signal ? {__proto__: null, signal } : undefined;
856884
returnPromisePrototypeThen(writer.writev(batch,opts),()=>{
857885
for(leti=0;i<batch.length;i++){
@@ -867,6 +895,10 @@ async function pipeTo(source, ...args) {
867895
for(leti=0;i<batch.length;i++){
868896
constchunk=batch[i];
869897
if(!hasWriteSync||!writer.writeSync(chunk)){
898+
if(hasWriteSync&&syncFalseWasAccepted()){
899+
totalBytes+=TypedArrayPrototypeGetByteLength(chunk);
900+
returnwriteBatchAfterAcceptedBackpressure(batch,i+1);
901+
}
870902
// Sync path failed at index i - fall back to async for the rest.
871903
// Count bytes for chunks already written synchronously (0..i-1).
872904
returnwriteBatchAsyncFallback(batch,i);

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

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@ const {
3232

3333
const{
3434
drainableProtocol,
35+
kSyncWriteAcceptedOnFalse,
3536
}=require('internal/streams/iter/types');
3637

3738
const{
@@ -560,6 +561,10 @@ class PushWriter {
560561
returnthis.#queue.desiredSize;
561562
}
562563

564+
get[kSyncWriteAcceptedOnFalse](){
565+
returnthis.#queue.backpressurePolicy==='block';
566+
}
567+
563568
write(chunk,options){
564569
if(!options?.signal&&this.#queue.canWriteSync()){
565570
constbytes=toUint8Array(chunk);
@@ -586,7 +591,8 @@ class PushWriter {
586591
writeSync(chunk){
587592
constbytes=toUint8Array(chunk);
588593
constresult=this.#queue.writeSync([bytes]);
589-
if(!result&&this.#queue.backpressurePolicy==='block'){
594+
if(!result&&this.#queue.backpressurePolicy==='block'&&
595+
this.#queue.desiredSize===0){
590596
// Block policy: force-enqueue and return false as backpressure signal.
591597
// Data IS accepted; false tells caller to slow down.
592598
this.#queue.forceEnqueue([bytes]);
@@ -601,7 +607,8 @@ class PushWriter {
601607
}
602608
constbytes=convertChunks(chunks);
603609
constresult=this.#queue.writeSync(bytes);
604-
if(!result&&this.#queue.backpressurePolicy==='block'){
610+
if(!result&&this.#queue.backpressurePolicy==='block'&&
611+
this.#queue.desiredSize===0){
605612
this.#queue.forceEnqueue(bytes);
606613
returnfalse;
607614
}

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,9 +64,12 @@ const kValidatedTransform = Symbol('kValidatedTransform');
6464
*/
6565
constkValidatedSource=Symbol('kValidatedSource');
6666

67+
constkSyncWriteAcceptedOnFalse=Symbol('kSyncWriteAcceptedOnFalse');
68+
6769
module.exports={
6870
broadcastProtocol,
6971
drainableProtocol,
72+
kSyncWriteAcceptedOnFalse,
7073
kValidatedSource,
7174
kValidatedTransform,
7275
shareProtocol,

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

Lines changed: 32 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,8 @@
55

66
constcommon=require('../common');
77
constassert=require('assert');
8-
const{ pipeTo, pipeToSync }=require('stream/iter');
8+
const{setImmediate: setImmediatePromise}=require('timers/promises');
9+
const{ pipeTo, pipeToSync, push, text }=require('stream/iter');
910

1011
// Multi-chunk batch with writevSync (sync success path)
1112
asyncfunctiontestWritevSyncSuccess(){
@@ -104,6 +105,35 @@ async function testWriteSyncAlwaysFails() {
104105
assert.strictEqual(total,2);
105106
}
106107

108+
// PushWriter block mode accepts sync writes even when returning false for
109+
// backpressure. pipeTo must wait for drain, not retry the same write.
110+
asyncfunctionassertPushWriterBlockPipeTo(source,expected,expectedTotal){
111+
const{ writer, readable }=push({
112+
highWaterMark: 1,
113+
backpressure: 'block',
114+
});
115+
116+
constpipe=pipeTo(source,writer);
117+
awaitsetImmediatePromise();
118+
constdata=awaittext(readable);
119+
consttotal=awaitpipe;
120+
121+
assert.strictEqual(data,expected);
122+
assert.strictEqual(total,expectedTotal);
123+
}
124+
125+
asyncfunctiontestPushWriterBlockSyncFalseAccepted(){
126+
awaitassertPushWriterBlockPipeTo((asyncfunction*(){
127+
yield[newUint8Array([97])];
128+
yield[newUint8Array([98])];
129+
})(),'ab',2);
130+
131+
awaitassertPushWriterBlockPipeTo((asyncfunction*(){
132+
yield[newUint8Array([97,98])];
133+
yield[newUint8Array([99]),newUint8Array([100])];
134+
})(),'abcd',4);
135+
}
136+
107137
// pipeToSync with writevSync
108138
asyncfunctiontestPipeToSyncWritev(){
109139
constbatches=[];
@@ -142,6 +172,7 @@ Promise.all([
142172
testWritevSyncFails(),
143173
testWriteSyncFailsMidBatch(),
144174
testWriteSyncAlwaysFails(),
175+
testPushWriterBlockSyncFalseAccepted(),
145176
testPipeToSyncWritev(),
146177
testPipeToSyncWriteFallback(),
147178
]).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 ec2666e

Browse files
trivikraduh95
authored andcommitted
stream: avoid retrying accepted pipeTo writes
PushWriter in block backpressure mode can return false from writeSync() and writevSync() after accepting data. Treat that false return as backpressure and wait for drain instead of retrying the same chunks asynchronously. Fixes: #63296 Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: openai:gpt-5.5 PR-URL: #63297 Backport-PR-URL: #64675Fixes: #63296 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Ethan Arrowood <ethan@arrowood.dev>
1 parent f63143f commit ec2666e

4 files changed

Lines changed: 76 additions & 3 deletions

File tree

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

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,8 @@ const {
5151
}=require('internal/streams/iter/utils');
5252

5353
const{
54+
drainableProtocol,
55+
kSyncWriteAcceptedOnFalse,
5456
kValidatedTransform,
5557
}=require('internal/streams/iter/types');
5658

@@ -828,13 +830,33 @@ async function pipeTo(source, ...args) {
828830
consthasWriteSync=typeofwriter.writeSync==='function';
829831
consthasWritevSync=typeofwriter.writevSync==='function';
830832
consthasEndSync=typeofwriter.endSync==='function';
833+
constsyncFalseCanBeAccepted=writer[kSyncWriteAcceptedOnFalse]===true;
834+
835+
functionsyncFalseWasAccepted(){
836+
returnsyncFalseCanBeAccepted&&writer.desiredSize===0;
837+
}
838+
839+
functionwaitForSyncBackpressure(){
840+
constondrain=writer[drainableProtocol];
841+
returnondrain?.call(writer);
842+
}
843+
844+
asyncfunctionwriteBatchAfterAcceptedBackpressure(batch,startIndex){
845+
awaitwaitForSyncBackpressure();
846+
awaitwriteBatchAsyncFallback(batch,startIndex);
847+
}
848+
831849
// Async fallback for writeBatch when sync write fails partway through.
832850
// Continues writing from batch[startIndex] using async write().
833851
asyncfunctionwriteBatchAsyncFallback(batch,startIndex){
834852
for(leti=startIndex;i<batch.length;i++){
835853
constchunk=batch[i];
836854
if(hasWriteSync&&writer.writeSync(chunk)){
837855
// Sync retry succeeded
856+
}elseif(syncFalseWasAccepted()){
857+
totalBytes+=TypedArrayPrototypeGetByteLength(chunk);
858+
awaitwaitForSyncBackpressure();
859+
continue;
838860
}else{
839861
constresult=writer.write(
840862
chunk,signal ? {__proto__: null, signal } : undefined);
@@ -852,6 +874,12 @@ async function pipeTo(source, ...args) {
852874
functionwriteBatch(batch){
853875
if(hasWritev&&batch.length>1){
854876
if(!hasWritevSync||!writer.writevSync(batch)){
877+
if(hasWritevSync&&syncFalseWasAccepted()){
878+
for(leti=0;i<batch.length;i++){
879+
totalBytes+=TypedArrayPrototypeGetByteLength(batch[i]);
880+
}
881+
returnwaitForSyncBackpressure();
882+
}
855883
constopts=signal ? {__proto__: null, signal } : undefined;
856884
returnPromisePrototypeThen(writer.writev(batch,opts),()=>{
857885
for(leti=0;i<batch.length;i++){
@@ -867,6 +895,10 @@ async function pipeTo(source, ...args) {
867895
for(leti=0;i<batch.length;i++){
868896
constchunk=batch[i];
869897
if(!hasWriteSync||!writer.writeSync(chunk)){
898+
if(hasWriteSync&&syncFalseWasAccepted()){
899+
totalBytes+=TypedArrayPrototypeGetByteLength(chunk);
900+
returnwriteBatchAfterAcceptedBackpressure(batch,i+1);
901+
}
870902
// Sync path failed at index i - fall back to async for the rest.
871903
// Count bytes for chunks already written synchronously (0..i-1).
872904
returnwriteBatchAsyncFallback(batch,i);

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

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@ const {
3232

3333
const{
3434
drainableProtocol,
35+
kSyncWriteAcceptedOnFalse,
3536
}=require('internal/streams/iter/types');
3637

3738
const{
@@ -560,6 +561,10 @@ class PushWriter {
560561
returnthis.#queue.desiredSize;
561562
}
562563

564+
get[kSyncWriteAcceptedOnFalse](){
565+
returnthis.#queue.backpressurePolicy==='block';
566+
}
567+
563568
write(chunk,options){
564569
if(!options?.signal&&this.#queue.canWriteSync()){
565570
constbytes=toUint8Array(chunk);
@@ -586,7 +591,8 @@ class PushWriter {
586591
writeSync(chunk){
587592
constbytes=toUint8Array(chunk);
588593
constresult=this.#queue.writeSync([bytes]);
589-
if(!result&&this.#queue.backpressurePolicy==='block'){
594+
if(!result&&this.#queue.backpressurePolicy==='block'&&
595+
this.#queue.desiredSize===0){
590596
// Block policy: force-enqueue and return false as backpressure signal.
591597
// Data IS accepted; false tells caller to slow down.
592598
this.#queue.forceEnqueue([bytes]);
@@ -601,7 +607,8 @@ class PushWriter {
601607
}
602608
constbytes=convertChunks(chunks);
603609
constresult=this.#queue.writeSync(bytes);
604-
if(!result&&this.#queue.backpressurePolicy==='block'){
610+
if(!result&&this.#queue.backpressurePolicy==='block'&&
611+
this.#queue.desiredSize===0){
605612
this.#queue.forceEnqueue(bytes);
606613
returnfalse;
607614
}

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,9 +64,12 @@ const kValidatedTransform = Symbol('kValidatedTransform');
6464
*/
6565
constkValidatedSource=Symbol('kValidatedSource');
6666

67+
constkSyncWriteAcceptedOnFalse=Symbol('kSyncWriteAcceptedOnFalse');
68+
6769
module.exports={
6870
broadcastProtocol,
6971
drainableProtocol,
72+
kSyncWriteAcceptedOnFalse,
7073
kValidatedSource,
7174
kValidatedTransform,
7275
shareProtocol,

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

Lines changed: 32 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,8 @@
55

66
constcommon=require('../common');
77
constassert=require('assert');
8-
const{ pipeTo, pipeToSync }=require('stream/iter');
8+
const{setImmediate: setImmediatePromise}=require('timers/promises');
9+
const{ pipeTo, pipeToSync, push, text }=require('stream/iter');
910

1011
// Multi-chunk batch with writevSync (sync success path)
1112
asyncfunctiontestWritevSyncSuccess(){
@@ -104,6 +105,35 @@ async function testWriteSyncAlwaysFails() {
104105
assert.strictEqual(total,2);
105106
}
106107

108+
// PushWriter block mode accepts sync writes even when returning false for
109+
// backpressure. pipeTo must wait for drain, not retry the same write.
110+
asyncfunctionassertPushWriterBlockPipeTo(source,expected,expectedTotal){
111+
const{ writer, readable }=push({
112+
highWaterMark: 1,
113+
backpressure: 'block',
114+
});
115+
116+
constpipe=pipeTo(source,writer);
117+
awaitsetImmediatePromise();
118+
constdata=awaittext(readable);
119+
consttotal=awaitpipe;
120+
121+
assert.strictEqual(data,expected);
122+
assert.strictEqual(total,expectedTotal);
123+
}
124+
125+
asyncfunctiontestPushWriterBlockSyncFalseAccepted(){
126+
awaitassertPushWriterBlockPipeTo((asyncfunction*(){
127+
yield[newUint8Array([97])];
128+
yield[newUint8Array([98])];
129+
})(),'ab',2);
130+
131+
awaitassertPushWriterBlockPipeTo((asyncfunction*(){
132+
yield[newUint8Array([97,98])];
133+
yield[newUint8Array([99]),newUint8Array([100])];
134+
})(),'abcd',4);
135+
}
136+
107137
// pipeToSync with writevSync
108138
asyncfunctiontestPipeToSyncWritev(){
109139
constbatches=[];
@@ -142,6 +172,7 @@ Promise.all([
142172
testWritevSyncFails(),
143173
testWriteSyncFailsMidBatch(),
144174
testWriteSyncAlwaysFails(),
175+
testPushWriterBlockSyncFalseAccepted(),
145176
testPipeToSyncWritev(),
146177
testPipeToSyncWriteFallback(),
147178
]).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 ec2666e

Browse files
trivikraduh95
authored andcommitted
stream: avoid retrying accepted pipeTo writes
PushWriter in block backpressure mode can return false from writeSync() and writevSync() after accepting data. Treat that false return as backpressure and wait for drain instead of retrying the same chunks asynchronously. Fixes: #63296 Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: openai:gpt-5.5 PR-URL: #63297 Backport-PR-URL: #64675Fixes: #63296 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Ethan Arrowood <ethan@arrowood.dev>
1 parent f63143f commit ec2666e

4 files changed

Lines changed: 76 additions & 3 deletions

File tree

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

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,8 @@ const {
5151
}=require('internal/streams/iter/utils');
5252

5353
const{
54+
drainableProtocol,
55+
kSyncWriteAcceptedOnFalse,
5456
kValidatedTransform,
5557
}=require('internal/streams/iter/types');
5658

@@ -828,13 +830,33 @@ async function pipeTo(source, ...args) {
828830
consthasWriteSync=typeofwriter.writeSync==='function';
829831
consthasWritevSync=typeofwriter.writevSync==='function';
830832
consthasEndSync=typeofwriter.endSync==='function';
833+
constsyncFalseCanBeAccepted=writer[kSyncWriteAcceptedOnFalse]===true;
834+
835+
functionsyncFalseWasAccepted(){
836+
returnsyncFalseCanBeAccepted&&writer.desiredSize===0;
837+
}
838+
839+
functionwaitForSyncBackpressure(){
840+
constondrain=writer[drainableProtocol];
841+
returnondrain?.call(writer);
842+
}
843+
844+
asyncfunctionwriteBatchAfterAcceptedBackpressure(batch,startIndex){
845+
awaitwaitForSyncBackpressure();
846+
awaitwriteBatchAsyncFallback(batch,startIndex);
847+
}
848+
831849
// Async fallback for writeBatch when sync write fails partway through.
832850
// Continues writing from batch[startIndex] using async write().
833851
asyncfunctionwriteBatchAsyncFallback(batch,startIndex){
834852
for(leti=startIndex;i<batch.length;i++){
835853
constchunk=batch[i];
836854
if(hasWriteSync&&writer.writeSync(chunk)){
837855
// Sync retry succeeded
856+
}elseif(syncFalseWasAccepted()){
857+
totalBytes+=TypedArrayPrototypeGetByteLength(chunk);
858+
awaitwaitForSyncBackpressure();
859+
continue;
838860
}else{
839861
constresult=writer.write(
840862
chunk,signal ? {__proto__: null, signal } : undefined);
@@ -852,6 +874,12 @@ async function pipeTo(source, ...args) {
852874
functionwriteBatch(batch){
853875
if(hasWritev&&batch.length>1){
854876
if(!hasWritevSync||!writer.writevSync(batch)){
877+
if(hasWritevSync&&syncFalseWasAccepted()){
878+
for(leti=0;i<batch.length;i++){
879+
totalBytes+=TypedArrayPrototypeGetByteLength(batch[i]);
880+
}
881+
returnwaitForSyncBackpressure();
882+
}
855883
constopts=signal ? {__proto__: null, signal } : undefined;
856884
returnPromisePrototypeThen(writer.writev(batch,opts),()=>{
857885
for(leti=0;i<batch.length;i++){
@@ -867,6 +895,10 @@ async function pipeTo(source, ...args) {
867895
for(leti=0;i<batch.length;i++){
868896
constchunk=batch[i];
869897
if(!hasWriteSync||!writer.writeSync(chunk)){
898+
if(hasWriteSync&&syncFalseWasAccepted()){
899+
totalBytes+=TypedArrayPrototypeGetByteLength(chunk);
900+
returnwriteBatchAfterAcceptedBackpressure(batch,i+1);
901+
}
870902
// Sync path failed at index i - fall back to async for the rest.
871903
// Count bytes for chunks already written synchronously (0..i-1).
872904
returnwriteBatchAsyncFallback(batch,i);

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

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@ const {
3232

3333
const{
3434
drainableProtocol,
35+
kSyncWriteAcceptedOnFalse,
3536
}=require('internal/streams/iter/types');
3637

3738
const{
@@ -560,6 +561,10 @@ class PushWriter {
560561
returnthis.#queue.desiredSize;
561562
}
562563

564+
get[kSyncWriteAcceptedOnFalse](){
565+
returnthis.#queue.backpressurePolicy==='block';
566+
}
567+
563568
write(chunk,options){
564569
if(!options?.signal&&this.#queue.canWriteSync()){
565570
constbytes=toUint8Array(chunk);
@@ -586,7 +591,8 @@ class PushWriter {
586591
writeSync(chunk){
587592
constbytes=toUint8Array(chunk);
588593
constresult=this.#queue.writeSync([bytes]);
589-
if(!result&&this.#queue.backpressurePolicy==='block'){
594+
if(!result&&this.#queue.backpressurePolicy==='block'&&
595+
this.#queue.desiredSize===0){
590596
// Block policy: force-enqueue and return false as backpressure signal.
591597
// Data IS accepted; false tells caller to slow down.
592598
this.#queue.forceEnqueue([bytes]);
@@ -601,7 +607,8 @@ class PushWriter {
601607
}
602608
constbytes=convertChunks(chunks);
603609
constresult=this.#queue.writeSync(bytes);
604-
if(!result&&this.#queue.backpressurePolicy==='block'){
610+
if(!result&&this.#queue.backpressurePolicy==='block'&&
611+
this.#queue.desiredSize===0){
605612
this.#queue.forceEnqueue(bytes);
606613
returnfalse;
607614
}

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,9 +64,12 @@ const kValidatedTransform = Symbol('kValidatedTransform');
6464
*/
6565
constkValidatedSource=Symbol('kValidatedSource');
6666

67+
constkSyncWriteAcceptedOnFalse=Symbol('kSyncWriteAcceptedOnFalse');
68+
6769
module.exports={
6870
broadcastProtocol,
6971
drainableProtocol,
72+
kSyncWriteAcceptedOnFalse,
7073
kValidatedSource,
7174
kValidatedTransform,
7275
shareProtocol,

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

Lines changed: 32 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,8 @@
55

66
constcommon=require('../common');
77
constassert=require('assert');
8-
const{ pipeTo, pipeToSync }=require('stream/iter');
8+
const{setImmediate: setImmediatePromise}=require('timers/promises');
9+
const{ pipeTo, pipeToSync, push, text }=require('stream/iter');
910

1011
// Multi-chunk batch with writevSync (sync success path)
1112
asyncfunctiontestWritevSyncSuccess(){
@@ -104,6 +105,35 @@ async function testWriteSyncAlwaysFails() {
104105
assert.strictEqual(total,2);
105106
}
106107

108+
// PushWriter block mode accepts sync writes even when returning false for
109+
// backpressure. pipeTo must wait for drain, not retry the same write.
110+
asyncfunctionassertPushWriterBlockPipeTo(source,expected,expectedTotal){
111+
const{ writer, readable }=push({
112+
highWaterMark: 1,
113+
backpressure: 'block',
114+
});
115+
116+
constpipe=pipeTo(source,writer);
117+
awaitsetImmediatePromise();
118+
constdata=awaittext(readable);
119+
consttotal=awaitpipe;
120+
121+
assert.strictEqual(data,expected);
122+
assert.strictEqual(total,expectedTotal);
123+
}
124+
125+
asyncfunctiontestPushWriterBlockSyncFalseAccepted(){
126+
awaitassertPushWriterBlockPipeTo((asyncfunction*(){
127+
yield[newUint8Array([97])];
128+
yield[newUint8Array([98])];
129+
})(),'ab',2);
130+
131+
awaitassertPushWriterBlockPipeTo((asyncfunction*(){
132+
yield[newUint8Array([97,98])];
133+
yield[newUint8Array([99]),newUint8Array([100])];
134+
})(),'abcd',4);
135+
}
136+
107137
// pipeToSync with writevSync
108138
asyncfunctiontestPipeToSyncWritev(){
109139
constbatches=[];
@@ -142,6 +172,7 @@ Promise.all([
142172
testWritevSyncFails(),
143173
testWriteSyncFailsMidBatch(),
144174
testWriteSyncAlwaysFails(),
175+
testPushWriterBlockSyncFalseAccepted(),
145176
testPipeToSyncWritev(),
146177
testPipeToSyncWriteFallback(),
147178
]).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 ec2666e

Browse files
trivikraduh95
authored andcommitted
stream: avoid retrying accepted pipeTo writes
PushWriter in block backpressure mode can return false from writeSync() and writevSync() after accepting data. Treat that false return as backpressure and wait for drain instead of retrying the same chunks asynchronously. Fixes: #63296 Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: openai:gpt-5.5 PR-URL: #63297 Backport-PR-URL: #64675Fixes: #63296 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Ethan Arrowood <ethan@arrowood.dev>
1 parent f63143f commit ec2666e

4 files changed

Lines changed: 76 additions & 3 deletions

File tree

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

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,8 @@ const {
5151
}=require('internal/streams/iter/utils');
5252

5353
const{
54+
drainableProtocol,
55+
kSyncWriteAcceptedOnFalse,
5456
kValidatedTransform,
5557
}=require('internal/streams/iter/types');
5658

@@ -828,13 +830,33 @@ async function pipeTo(source, ...args) {
828830
consthasWriteSync=typeofwriter.writeSync==='function';
829831
consthasWritevSync=typeofwriter.writevSync==='function';
830832
consthasEndSync=typeofwriter.endSync==='function';
833+
constsyncFalseCanBeAccepted=writer[kSyncWriteAcceptedOnFalse]===true;
834+
835+
functionsyncFalseWasAccepted(){
836+
returnsyncFalseCanBeAccepted&&writer.desiredSize===0;
837+
}
838+
839+
functionwaitForSyncBackpressure(){
840+
constondrain=writer[drainableProtocol];
841+
returnondrain?.call(writer);
842+
}
843+
844+
asyncfunctionwriteBatchAfterAcceptedBackpressure(batch,startIndex){
845+
awaitwaitForSyncBackpressure();
846+
awaitwriteBatchAsyncFallback(batch,startIndex);
847+
}
848+
831849
// Async fallback for writeBatch when sync write fails partway through.
832850
// Continues writing from batch[startIndex] using async write().
833851
asyncfunctionwriteBatchAsyncFallback(batch,startIndex){
834852
for(leti=startIndex;i<batch.length;i++){
835853
constchunk=batch[i];
836854
if(hasWriteSync&&writer.writeSync(chunk)){
837855
// Sync retry succeeded
856+
}elseif(syncFalseWasAccepted()){
857+
totalBytes+=TypedArrayPrototypeGetByteLength(chunk);
858+
awaitwaitForSyncBackpressure();
859+
continue;
838860
}else{
839861
constresult=writer.write(
840862
chunk,signal ? {__proto__: null, signal } : undefined);
@@ -852,6 +874,12 @@ async function pipeTo(source, ...args) {
852874
functionwriteBatch(batch){
853875
if(hasWritev&&batch.length>1){
854876
if(!hasWritevSync||!writer.writevSync(batch)){
877+
if(hasWritevSync&&syncFalseWasAccepted()){
878+
for(leti=0;i<batch.length;i++){
879+
totalBytes+=TypedArrayPrototypeGetByteLength(batch[i]);
880+
}
881+
returnwaitForSyncBackpressure();
882+
}
855883
constopts=signal ? {__proto__: null, signal } : undefined;
856884
returnPromisePrototypeThen(writer.writev(batch,opts),()=>{
857885
for(leti=0;i<batch.length;i++){
@@ -867,6 +895,10 @@ async function pipeTo(source, ...args) {
867895
for(leti=0;i<batch.length;i++){
868896
constchunk=batch[i];
869897
if(!hasWriteSync||!writer.writeSync(chunk)){
898+
if(hasWriteSync&&syncFalseWasAccepted()){
899+
totalBytes+=TypedArrayPrototypeGetByteLength(chunk);
900+
returnwriteBatchAfterAcceptedBackpressure(batch,i+1);
901+
}
870902
// Sync path failed at index i - fall back to async for the rest.
871903
// Count bytes for chunks already written synchronously (0..i-1).
872904
returnwriteBatchAsyncFallback(batch,i);

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

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@ const {
3232

3333
const{
3434
drainableProtocol,
35+
kSyncWriteAcceptedOnFalse,
3536
}=require('internal/streams/iter/types');
3637

3738
const{
@@ -560,6 +561,10 @@ class PushWriter {
560561
returnthis.#queue.desiredSize;
561562
}
562563

564+
get[kSyncWriteAcceptedOnFalse](){
565+
returnthis.#queue.backpressurePolicy==='block';
566+
}
567+
563568
write(chunk,options){
564569
if(!options?.signal&&this.#queue.canWriteSync()){
565570
constbytes=toUint8Array(chunk);
@@ -586,7 +591,8 @@ class PushWriter {
586591
writeSync(chunk){
587592
constbytes=toUint8Array(chunk);
588593
constresult=this.#queue.writeSync([bytes]);
589-
if(!result&&this.#queue.backpressurePolicy==='block'){
594+
if(!result&&this.#queue.backpressurePolicy==='block'&&
595+
this.#queue.desiredSize===0){
590596
// Block policy: force-enqueue and return false as backpressure signal.
591597
// Data IS accepted; false tells caller to slow down.
592598
this.#queue.forceEnqueue([bytes]);
@@ -601,7 +607,8 @@ class PushWriter {
601607
}
602608
constbytes=convertChunks(chunks);
603609
constresult=this.#queue.writeSync(bytes);
604-
if(!result&&this.#queue.backpressurePolicy==='block'){
610+
if(!result&&this.#queue.backpressurePolicy==='block'&&
611+
this.#queue.desiredSize===0){
605612
this.#queue.forceEnqueue(bytes);
606613
returnfalse;
607614
}

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,9 +64,12 @@ const kValidatedTransform = Symbol('kValidatedTransform');
6464
*/
6565
constkValidatedSource=Symbol('kValidatedSource');
6666

67+
constkSyncWriteAcceptedOnFalse=Symbol('kSyncWriteAcceptedOnFalse');
68+
6769
module.exports={
6870
broadcastProtocol,
6971
drainableProtocol,
72+
kSyncWriteAcceptedOnFalse,
7073
kValidatedSource,
7174
kValidatedTransform,
7275
shareProtocol,

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

Lines changed: 32 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,8 @@
55

66
constcommon=require('../common');
77
constassert=require('assert');
8-
const{ pipeTo, pipeToSync }=require('stream/iter');
8+
const{setImmediate: setImmediatePromise}=require('timers/promises');
9+
const{ pipeTo, pipeToSync, push, text }=require('stream/iter');
910

1011
// Multi-chunk batch with writevSync (sync success path)
1112
asyncfunctiontestWritevSyncSuccess(){
@@ -104,6 +105,35 @@ async function testWriteSyncAlwaysFails() {
104105
assert.strictEqual(total,2);
105106
}
106107

108+
// PushWriter block mode accepts sync writes even when returning false for
109+
// backpressure. pipeTo must wait for drain, not retry the same write.
110+
asyncfunctionassertPushWriterBlockPipeTo(source,expected,expectedTotal){
111+
const{ writer, readable }=push({
112+
highWaterMark: 1,
113+
backpressure: 'block',
114+
});
115+
116+
constpipe=pipeTo(source,writer);
117+
awaitsetImmediatePromise();
118+
constdata=awaittext(readable);
119+
consttotal=awaitpipe;
120+
121+
assert.strictEqual(data,expected);
122+
assert.strictEqual(total,expectedTotal);
123+
}
124+
125+
asyncfunctiontestPushWriterBlockSyncFalseAccepted(){
126+
awaitassertPushWriterBlockPipeTo((asyncfunction*(){
127+
yield[newUint8Array([97])];
128+
yield[newUint8Array([98])];
129+
})(),'ab',2);
130+
131+
awaitassertPushWriterBlockPipeTo((asyncfunction*(){
132+
yield[newUint8Array([97,98])];
133+
yield[newUint8Array([99]),newUint8Array([100])];
134+
})(),'abcd',4);
135+
}
136+
107137
// pipeToSync with writevSync
108138
asyncfunctiontestPipeToSyncWritev(){
109139
constbatches=[];
@@ -142,6 +172,7 @@ Promise.all([
142172
testWritevSyncFails(),
143173
testWriteSyncFailsMidBatch(),
144174
testWriteSyncAlwaysFails(),
175+
testPushWriterBlockSyncFalseAccepted(),
145176
testPipeToSyncWritev(),
146177
testPipeToSyncWriteFallback(),
147178
]).then(common.mustCall());

0 commit comments

Comments
Β (0)