Commit 486cff4

Browse files
trivikraduh95
authored andcommitted
stream: honor AbortSignal in Writer.end()
Reject Writer.end() when its signal is already aborted without closing the writer. For push writers, reject the pending operation if the signal aborts while buffered data drains, while allowing the graceful close to continue. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: codex:gpt-5.6-sol PR-URL: #64727Fixes: #64726 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Gürgün Dayıoğlu <hey@gurgun.day>
1 parent dc30379 commit 486cff4

5 files changed

Lines changed: 108 additions & 5 deletions

File tree

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

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -630,9 +630,10 @@ class BroadcastWriter {
630630
returnfalse;
631631
}
632632

633-
// end() is synchronous internally - signal accepted for interface compliance.
634633
end(options){
635-
getWriterSignal(options);
634+
constsignal=getWriterSignal(options);
635+
if(signal?.aborted)returnPromiseReject(signal.reason);
636+
636637
if(this.#isClosed())returnthis.#closed;
637638
this.#closed =PromiseResolve(this.#totalBytes);
638639
this.#broadcast[kEnd]();

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

Lines changed: 35 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@
88
const{
99
ArrayIsArray,
1010
ArrayPrototypePush,
11+
PromisePrototypeThen,
1112
PromiseReject,
1213
PromiseResolve,
1314
PromiseWithResolvers,
@@ -55,6 +56,35 @@ const {
5556

5657
constkNoFailReason=Symbol('kNoFailReason');
5758

59+
functionraceEndWithSignal(promise,signal){
60+
if(!signal)returnpromise;
61+
if(signal.aborted)returnPromiseReject(signal.reason);
62+
63+
const{
64+
promise: signaledPromise,
65+
resolve,
66+
reject,
67+
}=PromiseWithResolvers();
68+
constonAbort=()=>reject(signal.reason);
69+
70+
signal.addEventListener('abort',onAbort,{
71+
__proto__: null,
72+
once: true,
73+
});
74+
PromisePrototypeThen(
75+
promise,
76+
(value)=>{
77+
signal.removeEventListener('abort',onAbort);
78+
resolve(value);
79+
},
80+
(reason)=>{
81+
signal.removeEventListener('abort',onAbort);
82+
reject(reason);
83+
},
84+
);
85+
returnsignaledPromise;
86+
}
87+
5888
// =============================================================================
5989
// PushQueue - Internal Queue with Chunk-Based Backpressure
6090
// =============================================================================
@@ -628,7 +658,9 @@ class PushWriter {
628658
}
629659

630660
end(options){
631-
getWriterSignal(options);
661+
constsignal=getWriterSignal(options);
662+
if(signal?.aborted)returnPromiseReject(signal.reason);
663+
632664
constresult=this.#queue.end();
633665
if(result===-2){
634666
// Errored: reject with stored error
@@ -639,11 +671,11 @@ class PushWriter {
639671
// when consumer drains past the end sentinel
640672
constpendingEndPromise=this.#queue.pendingEndPromise;
641673
if(pendingEndPromise!==null){
642-
returnpendingEndPromise;
674+
returnraceEndWithSignal(pendingEndPromise,signal);
643675
}
644676
const{ promise, resolve, reject }=PromiseWithResolvers();
645677
this.#queue.setPendingEnd({__proto__: null, promise, resolve, reject });
646-
returnpromise;
678+
returnraceEndWithSignal(promise,signal);
647679
}
648680
// >= 0: byte count (immediate close or idempotent)
649681
returnPromiseResolve(result);

β€Žtest/parallel/test-stream-iter-broadcast-basic.jsβ€Ž

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -111,6 +111,22 @@ async function testWriterEnd() {
111111
assert.strictEqual(data,'data');
112112
}
113113

114+
asyncfunctiontestWriterEndWithPreAbortedSignal(){
115+
const{ writer,broadcast: bc}=broadcast();
116+
constconsumer=bc.push();
117+
constreason=newError('end aborted');
118+
119+
awaitassert.rejects(
120+
writer.end({signal: AbortSignal.abort(reason)}),
121+
(error)=>error===reason,
122+
);
123+
124+
// A rejected end must leave the writer open.
125+
awaitwriter.write('data');
126+
assert.strictEqual(awaitwriter.end(),4);
127+
assert.strictEqual(awaittext(consumer),'data');
128+
}
129+
114130
asyncfunctiontestWriterFail(){
115131
const{ writer,broadcast: bc}=broadcast();
116132
constconsumer=bc.push();
@@ -308,6 +324,7 @@ Promise.all([
308324
testWriteSync(),
309325
testWritevSync(),
310326
testWriterEnd(),
327+
testWriterEndWithPreAbortedSignal(),
311328
testWriterFail(),
312329
testCancelWithoutReason(),
313330
testCancelWithReason(),

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

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -127,6 +127,22 @@ async function testAbortSignal() {
127127
);
128128
}
129129

130+
asyncfunctiontestWriterEndWithPreAbortedSignal(){
131+
const[channelA,channelB]=duplex();
132+
constreason=newError('end aborted');
133+
134+
awaitassert.rejects(
135+
channelA.writer.end({signal: AbortSignal.abort(reason)}),
136+
(error)=>error===reason,
137+
);
138+
139+
awaitchannelA.writer.write('still open');
140+
constcompletedEnd=channelA.writer.end();
141+
assert.strictEqual(awaittext(channelB.readable),'still open');
142+
assert.strictEqual(awaitcompletedEnd,10);
143+
awaitchannelB.close();
144+
}
145+
130146
asyncfunctiontestEmptyDuplex(){
131147
const[channelA,channelB]=duplex();
132148

@@ -182,6 +198,7 @@ Promise.all([
182198
testWithOptions(),
183199
testPerChannelOptions(),
184200
testAbortSignal(),
201+
testWriterEndWithPreAbortedSignal(),
185202
testEmptyDuplex(),
186203
testChannelFail(),
187204
testAbortSignalBothChannels(),

β€Žtest/parallel/test-stream-iter-push-writer.jsβ€Ž

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -260,6 +260,40 @@ async function testEndAsyncReturnValue() {
260260
awaitconsume;
261261
}
262262

263+
asyncfunctiontestEndWithPreAbortedSignal(){
264+
const{ writer, readable }=push();
265+
constreason=newError('end aborted');
266+
267+
writer.writeSync('hello');
268+
awaitassert.rejects(
269+
writer.end({signal: AbortSignal.abort(reason)}),
270+
(error)=>error===reason,
271+
);
272+
273+
// A rejected end must leave the writer open.
274+
writer.writeSync(' world');
275+
constconsume=text(readable);
276+
assert.strictEqual(awaitwriter.end(),11);
277+
assert.strictEqual(awaitconsume,'hello world');
278+
}
279+
280+
asyncfunctiontestEndSignalAbortWhileDraining(){
281+
const{ writer, readable }=push();
282+
constcontroller=newAbortController();
283+
constreason=newError('end aborted while draining');
284+
285+
writer.writeSync('hello');
286+
constabortedEnd=writer.end({signal: controller.signal});
287+
controller.abort(reason);
288+
289+
awaitassert.rejects(abortedEnd,(error)=>error===reason);
290+
291+
// Aborting the operation does not undo the end-of-stream signal.
292+
constcompletedEnd=writer.end();
293+
assert.strictEqual(awaittext(readable),'hello');
294+
assert.strictEqual(awaitcompletedEnd,5);
295+
}
296+
263297
asyncfunctiontestEndAfterEndSyncWaitsForDrain(){
264298
const{ writer, readable }=push();
265299
writer.writeSync('hello');
@@ -553,6 +587,8 @@ Promise.all([
553587
testOndrainProtocolErrorPropagates(),
554588
testFail(),
555589
testEndAsyncReturnValue(),
590+
testEndWithPreAbortedSignal(),
591+
testEndSignalAbortWhileDraining(),
556592
testEndAfterEndSyncWaitsForDrain(),
557593
testWriteUint8Array(),
558594
testOndrainWaitsForDrain(),

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 486cff4

Browse files
trivikraduh95
authored andcommitted
stream: honor AbortSignal in Writer.end()
Reject Writer.end() when its signal is already aborted without closing the writer. For push writers, reject the pending operation if the signal aborts while buffered data drains, while allowing the graceful close to continue. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: codex:gpt-5.6-sol PR-URL: #64727Fixes: #64726 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Gürgün Dayıoğlu <hey@gurgun.day>
1 parent dc30379 commit 486cff4

5 files changed

Lines changed: 108 additions & 5 deletions

File tree

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

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -630,9 +630,10 @@ class BroadcastWriter {
630630
returnfalse;
631631
}
632632

633-
// end() is synchronous internally - signal accepted for interface compliance.
634633
end(options){
635-
getWriterSignal(options);
634+
constsignal=getWriterSignal(options);
635+
if(signal?.aborted)returnPromiseReject(signal.reason);
636+
636637
if(this.#isClosed())returnthis.#closed;
637638
this.#closed =PromiseResolve(this.#totalBytes);
638639
this.#broadcast[kEnd]();

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

Lines changed: 35 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@
88
const{
99
ArrayIsArray,
1010
ArrayPrototypePush,
11+
PromisePrototypeThen,
1112
PromiseReject,
1213
PromiseResolve,
1314
PromiseWithResolvers,
@@ -55,6 +56,35 @@ const {
5556

5657
constkNoFailReason=Symbol('kNoFailReason');
5758

59+
functionraceEndWithSignal(promise,signal){
60+
if(!signal)returnpromise;
61+
if(signal.aborted)returnPromiseReject(signal.reason);
62+
63+
const{
64+
promise: signaledPromise,
65+
resolve,
66+
reject,
67+
}=PromiseWithResolvers();
68+
constonAbort=()=>reject(signal.reason);
69+
70+
signal.addEventListener('abort',onAbort,{
71+
__proto__: null,
72+
once: true,
73+
});
74+
PromisePrototypeThen(
75+
promise,
76+
(value)=>{
77+
signal.removeEventListener('abort',onAbort);
78+
resolve(value);
79+
},
80+
(reason)=>{
81+
signal.removeEventListener('abort',onAbort);
82+
reject(reason);
83+
},
84+
);
85+
returnsignaledPromise;
86+
}
87+
5888
// =============================================================================
5989
// PushQueue - Internal Queue with Chunk-Based Backpressure
6090
// =============================================================================
@@ -628,7 +658,9 @@ class PushWriter {
628658
}
629659

630660
end(options){
631-
getWriterSignal(options);
661+
constsignal=getWriterSignal(options);
662+
if(signal?.aborted)returnPromiseReject(signal.reason);
663+
632664
constresult=this.#queue.end();
633665
if(result===-2){
634666
// Errored: reject with stored error
@@ -639,11 +671,11 @@ class PushWriter {
639671
// when consumer drains past the end sentinel
640672
constpendingEndPromise=this.#queue.pendingEndPromise;
641673
if(pendingEndPromise!==null){
642-
returnpendingEndPromise;
674+
returnraceEndWithSignal(pendingEndPromise,signal);
643675
}
644676
const{ promise, resolve, reject }=PromiseWithResolvers();
645677
this.#queue.setPendingEnd({__proto__: null, promise, resolve, reject });
646-
returnpromise;
678+
returnraceEndWithSignal(promise,signal);
647679
}
648680
// >= 0: byte count (immediate close or idempotent)
649681
returnPromiseResolve(result);

β€Žtest/parallel/test-stream-iter-broadcast-basic.jsβ€Ž

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -111,6 +111,22 @@ async function testWriterEnd() {
111111
assert.strictEqual(data,'data');
112112
}
113113

114+
asyncfunctiontestWriterEndWithPreAbortedSignal(){
115+
const{ writer,broadcast: bc}=broadcast();
116+
constconsumer=bc.push();
117+
constreason=newError('end aborted');
118+
119+
awaitassert.rejects(
120+
writer.end({signal: AbortSignal.abort(reason)}),
121+
(error)=>error===reason,
122+
);
123+
124+
// A rejected end must leave the writer open.
125+
awaitwriter.write('data');
126+
assert.strictEqual(awaitwriter.end(),4);
127+
assert.strictEqual(awaittext(consumer),'data');
128+
}
129+
114130
asyncfunctiontestWriterFail(){
115131
const{ writer,broadcast: bc}=broadcast();
116132
constconsumer=bc.push();
@@ -308,6 +324,7 @@ Promise.all([
308324
testWriteSync(),
309325
testWritevSync(),
310326
testWriterEnd(),
327+
testWriterEndWithPreAbortedSignal(),
311328
testWriterFail(),
312329
testCancelWithoutReason(),
313330
testCancelWithReason(),

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

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -127,6 +127,22 @@ async function testAbortSignal() {
127127
);
128128
}
129129

130+
asyncfunctiontestWriterEndWithPreAbortedSignal(){
131+
const[channelA,channelB]=duplex();
132+
constreason=newError('end aborted');
133+
134+
awaitassert.rejects(
135+
channelA.writer.end({signal: AbortSignal.abort(reason)}),
136+
(error)=>error===reason,
137+
);
138+
139+
awaitchannelA.writer.write('still open');
140+
constcompletedEnd=channelA.writer.end();
141+
assert.strictEqual(awaittext(channelB.readable),'still open');
142+
assert.strictEqual(awaitcompletedEnd,10);
143+
awaitchannelB.close();
144+
}
145+
130146
asyncfunctiontestEmptyDuplex(){
131147
const[channelA,channelB]=duplex();
132148

@@ -182,6 +198,7 @@ Promise.all([
182198
testWithOptions(),
183199
testPerChannelOptions(),
184200
testAbortSignal(),
201+
testWriterEndWithPreAbortedSignal(),
185202
testEmptyDuplex(),
186203
testChannelFail(),
187204
testAbortSignalBothChannels(),

β€Žtest/parallel/test-stream-iter-push-writer.jsβ€Ž

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -260,6 +260,40 @@ async function testEndAsyncReturnValue() {
260260
awaitconsume;
261261
}
262262

263+
asyncfunctiontestEndWithPreAbortedSignal(){
264+
const{ writer, readable }=push();
265+
constreason=newError('end aborted');
266+
267+
writer.writeSync('hello');
268+
awaitassert.rejects(
269+
writer.end({signal: AbortSignal.abort(reason)}),
270+
(error)=>error===reason,
271+
);
272+
273+
// A rejected end must leave the writer open.
274+
writer.writeSync(' world');
275+
constconsume=text(readable);
276+
assert.strictEqual(awaitwriter.end(),11);
277+
assert.strictEqual(awaitconsume,'hello world');
278+
}
279+
280+
asyncfunctiontestEndSignalAbortWhileDraining(){
281+
const{ writer, readable }=push();
282+
constcontroller=newAbortController();
283+
constreason=newError('end aborted while draining');
284+
285+
writer.writeSync('hello');
286+
constabortedEnd=writer.end({signal: controller.signal});
287+
controller.abort(reason);
288+
289+
awaitassert.rejects(abortedEnd,(error)=>error===reason);
290+
291+
// Aborting the operation does not undo the end-of-stream signal.
292+
constcompletedEnd=writer.end();
293+
assert.strictEqual(awaittext(readable),'hello');
294+
assert.strictEqual(awaitcompletedEnd,5);
295+
}
296+
263297
asyncfunctiontestEndAfterEndSyncWaitsForDrain(){
264298
const{ writer, readable }=push();
265299
writer.writeSync('hello');
@@ -553,6 +587,8 @@ Promise.all([
553587
testOndrainProtocolErrorPropagates(),
554588
testFail(),
555589
testEndAsyncReturnValue(),
590+
testEndWithPreAbortedSignal(),
591+
testEndSignalAbortWhileDraining(),
556592
testEndAfterEndSyncWaitsForDrain(),
557593
testWriteUint8Array(),
558594
testOndrainWaitsForDrain(),

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 486cff4

Browse files
trivikraduh95
authored andcommitted
stream: honor AbortSignal in Writer.end()
Reject Writer.end() when its signal is already aborted without closing the writer. For push writers, reject the pending operation if the signal aborts while buffered data drains, while allowing the graceful close to continue. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: codex:gpt-5.6-sol PR-URL: #64727Fixes: #64726 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Gürgün Dayıoğlu <hey@gurgun.day>
1 parent dc30379 commit 486cff4

5 files changed

Lines changed: 108 additions & 5 deletions

File tree

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

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -630,9 +630,10 @@ class BroadcastWriter {
630630
returnfalse;
631631
}
632632

633-
// end() is synchronous internally - signal accepted for interface compliance.
634633
end(options){
635-
getWriterSignal(options);
634+
constsignal=getWriterSignal(options);
635+
if(signal?.aborted)returnPromiseReject(signal.reason);
636+
636637
if(this.#isClosed())returnthis.#closed;
637638
this.#closed =PromiseResolve(this.#totalBytes);
638639
this.#broadcast[kEnd]();

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

Lines changed: 35 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@
88
const{
99
ArrayIsArray,
1010
ArrayPrototypePush,
11+
PromisePrototypeThen,
1112
PromiseReject,
1213
PromiseResolve,
1314
PromiseWithResolvers,
@@ -55,6 +56,35 @@ const {
5556

5657
constkNoFailReason=Symbol('kNoFailReason');
5758

59+
functionraceEndWithSignal(promise,signal){
60+
if(!signal)returnpromise;
61+
if(signal.aborted)returnPromiseReject(signal.reason);
62+
63+
const{
64+
promise: signaledPromise,
65+
resolve,
66+
reject,
67+
}=PromiseWithResolvers();
68+
constonAbort=()=>reject(signal.reason);
69+
70+
signal.addEventListener('abort',onAbort,{
71+
__proto__: null,
72+
once: true,
73+
});
74+
PromisePrototypeThen(
75+
promise,
76+
(value)=>{
77+
signal.removeEventListener('abort',onAbort);
78+
resolve(value);
79+
},
80+
(reason)=>{
81+
signal.removeEventListener('abort',onAbort);
82+
reject(reason);
83+
},
84+
);
85+
returnsignaledPromise;
86+
}
87+
5888
// =============================================================================
5989
// PushQueue - Internal Queue with Chunk-Based Backpressure
6090
// =============================================================================
@@ -628,7 +658,9 @@ class PushWriter {
628658
}
629659

630660
end(options){
631-
getWriterSignal(options);
661+
constsignal=getWriterSignal(options);
662+
if(signal?.aborted)returnPromiseReject(signal.reason);
663+
632664
constresult=this.#queue.end();
633665
if(result===-2){
634666
// Errored: reject with stored error
@@ -639,11 +671,11 @@ class PushWriter {
639671
// when consumer drains past the end sentinel
640672
constpendingEndPromise=this.#queue.pendingEndPromise;
641673
if(pendingEndPromise!==null){
642-
returnpendingEndPromise;
674+
returnraceEndWithSignal(pendingEndPromise,signal);
643675
}
644676
const{ promise, resolve, reject }=PromiseWithResolvers();
645677
this.#queue.setPendingEnd({__proto__: null, promise, resolve, reject });
646-
returnpromise;
678+
returnraceEndWithSignal(promise,signal);
647679
}
648680
// >= 0: byte count (immediate close or idempotent)
649681
returnPromiseResolve(result);

β€Žtest/parallel/test-stream-iter-broadcast-basic.jsβ€Ž

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -111,6 +111,22 @@ async function testWriterEnd() {
111111
assert.strictEqual(data,'data');
112112
}
113113

114+
asyncfunctiontestWriterEndWithPreAbortedSignal(){
115+
const{ writer,broadcast: bc}=broadcast();
116+
constconsumer=bc.push();
117+
constreason=newError('end aborted');
118+
119+
awaitassert.rejects(
120+
writer.end({signal: AbortSignal.abort(reason)}),
121+
(error)=>error===reason,
122+
);
123+
124+
// A rejected end must leave the writer open.
125+
awaitwriter.write('data');
126+
assert.strictEqual(awaitwriter.end(),4);
127+
assert.strictEqual(awaittext(consumer),'data');
128+
}
129+
114130
asyncfunctiontestWriterFail(){
115131
const{ writer,broadcast: bc}=broadcast();
116132
constconsumer=bc.push();
@@ -308,6 +324,7 @@ Promise.all([
308324
testWriteSync(),
309325
testWritevSync(),
310326
testWriterEnd(),
327+
testWriterEndWithPreAbortedSignal(),
311328
testWriterFail(),
312329
testCancelWithoutReason(),
313330
testCancelWithReason(),

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

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -127,6 +127,22 @@ async function testAbortSignal() {
127127
);
128128
}
129129

130+
asyncfunctiontestWriterEndWithPreAbortedSignal(){
131+
const[channelA,channelB]=duplex();
132+
constreason=newError('end aborted');
133+
134+
awaitassert.rejects(
135+
channelA.writer.end({signal: AbortSignal.abort(reason)}),
136+
(error)=>error===reason,
137+
);
138+
139+
awaitchannelA.writer.write('still open');
140+
constcompletedEnd=channelA.writer.end();
141+
assert.strictEqual(awaittext(channelB.readable),'still open');
142+
assert.strictEqual(awaitcompletedEnd,10);
143+
awaitchannelB.close();
144+
}
145+
130146
asyncfunctiontestEmptyDuplex(){
131147
const[channelA,channelB]=duplex();
132148

@@ -182,6 +198,7 @@ Promise.all([
182198
testWithOptions(),
183199
testPerChannelOptions(),
184200
testAbortSignal(),
201+
testWriterEndWithPreAbortedSignal(),
185202
testEmptyDuplex(),
186203
testChannelFail(),
187204
testAbortSignalBothChannels(),

β€Žtest/parallel/test-stream-iter-push-writer.jsβ€Ž

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -260,6 +260,40 @@ async function testEndAsyncReturnValue() {
260260
awaitconsume;
261261
}
262262

263+
asyncfunctiontestEndWithPreAbortedSignal(){
264+
const{ writer, readable }=push();
265+
constreason=newError('end aborted');
266+
267+
writer.writeSync('hello');
268+
awaitassert.rejects(
269+
writer.end({signal: AbortSignal.abort(reason)}),
270+
(error)=>error===reason,
271+
);
272+
273+
// A rejected end must leave the writer open.
274+
writer.writeSync(' world');
275+
constconsume=text(readable);
276+
assert.strictEqual(awaitwriter.end(),11);
277+
assert.strictEqual(awaitconsume,'hello world');
278+
}
279+
280+
asyncfunctiontestEndSignalAbortWhileDraining(){
281+
const{ writer, readable }=push();
282+
constcontroller=newAbortController();
283+
constreason=newError('end aborted while draining');
284+
285+
writer.writeSync('hello');
286+
constabortedEnd=writer.end({signal: controller.signal});
287+
controller.abort(reason);
288+
289+
awaitassert.rejects(abortedEnd,(error)=>error===reason);
290+
291+
// Aborting the operation does not undo the end-of-stream signal.
292+
constcompletedEnd=writer.end();
293+
assert.strictEqual(awaittext(readable),'hello');
294+
assert.strictEqual(awaitcompletedEnd,5);
295+
}
296+
263297
asyncfunctiontestEndAfterEndSyncWaitsForDrain(){
264298
const{ writer, readable }=push();
265299
writer.writeSync('hello');
@@ -553,6 +587,8 @@ Promise.all([
553587
testOndrainProtocolErrorPropagates(),
554588
testFail(),
555589
testEndAsyncReturnValue(),
590+
testEndWithPreAbortedSignal(),
591+
testEndSignalAbortWhileDraining(),
556592
testEndAfterEndSyncWaitsForDrain(),
557593
testWriteUint8Array(),
558594
testOndrainWaitsForDrain(),

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 486cff4

Browse files
trivikraduh95
authored andcommitted
stream: honor AbortSignal in Writer.end()
Reject Writer.end() when its signal is already aborted without closing the writer. For push writers, reject the pending operation if the signal aborts while buffered data drains, while allowing the graceful close to continue. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: codex:gpt-5.6-sol PR-URL: #64727Fixes: #64726 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Gürgün Dayıoğlu <hey@gurgun.day>
1 parent dc30379 commit 486cff4

5 files changed

Lines changed: 108 additions & 5 deletions

File tree

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

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -630,9 +630,10 @@ class BroadcastWriter {
630630
returnfalse;
631631
}
632632

633-
// end() is synchronous internally - signal accepted for interface compliance.
634633
end(options){
635-
getWriterSignal(options);
634+
constsignal=getWriterSignal(options);
635+
if(signal?.aborted)returnPromiseReject(signal.reason);
636+
636637
if(this.#isClosed())returnthis.#closed;
637638
this.#closed =PromiseResolve(this.#totalBytes);
638639
this.#broadcast[kEnd]();

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

Lines changed: 35 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@
88
const{
99
ArrayIsArray,
1010
ArrayPrototypePush,
11+
PromisePrototypeThen,
1112
PromiseReject,
1213
PromiseResolve,
1314
PromiseWithResolvers,
@@ -55,6 +56,35 @@ const {
5556

5657
constkNoFailReason=Symbol('kNoFailReason');
5758

59+
functionraceEndWithSignal(promise,signal){
60+
if(!signal)returnpromise;
61+
if(signal.aborted)returnPromiseReject(signal.reason);
62+
63+
const{
64+
promise: signaledPromise,
65+
resolve,
66+
reject,
67+
}=PromiseWithResolvers();
68+
constonAbort=()=>reject(signal.reason);
69+
70+
signal.addEventListener('abort',onAbort,{
71+
__proto__: null,
72+
once: true,
73+
});
74+
PromisePrototypeThen(
75+
promise,
76+
(value)=>{
77+
signal.removeEventListener('abort',onAbort);
78+
resolve(value);
79+
},
80+
(reason)=>{
81+
signal.removeEventListener('abort',onAbort);
82+
reject(reason);
83+
},
84+
);
85+
returnsignaledPromise;
86+
}
87+
5888
// =============================================================================
5989
// PushQueue - Internal Queue with Chunk-Based Backpressure
6090
// =============================================================================
@@ -628,7 +658,9 @@ class PushWriter {
628658
}
629659

630660
end(options){
631-
getWriterSignal(options);
661+
constsignal=getWriterSignal(options);
662+
if(signal?.aborted)returnPromiseReject(signal.reason);
663+
632664
constresult=this.#queue.end();
633665
if(result===-2){
634666
// Errored: reject with stored error
@@ -639,11 +671,11 @@ class PushWriter {
639671
// when consumer drains past the end sentinel
640672
constpendingEndPromise=this.#queue.pendingEndPromise;
641673
if(pendingEndPromise!==null){
642-
returnpendingEndPromise;
674+
returnraceEndWithSignal(pendingEndPromise,signal);
643675
}
644676
const{ promise, resolve, reject }=PromiseWithResolvers();
645677
this.#queue.setPendingEnd({__proto__: null, promise, resolve, reject });
646-
returnpromise;
678+
returnraceEndWithSignal(promise,signal);
647679
}
648680
// >= 0: byte count (immediate close or idempotent)
649681
returnPromiseResolve(result);

β€Žtest/parallel/test-stream-iter-broadcast-basic.jsβ€Ž

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -111,6 +111,22 @@ async function testWriterEnd() {
111111
assert.strictEqual(data,'data');
112112
}
113113

114+
asyncfunctiontestWriterEndWithPreAbortedSignal(){
115+
const{ writer,broadcast: bc}=broadcast();
116+
constconsumer=bc.push();
117+
constreason=newError('end aborted');
118+
119+
awaitassert.rejects(
120+
writer.end({signal: AbortSignal.abort(reason)}),
121+
(error)=>error===reason,
122+
);
123+
124+
// A rejected end must leave the writer open.
125+
awaitwriter.write('data');
126+
assert.strictEqual(awaitwriter.end(),4);
127+
assert.strictEqual(awaittext(consumer),'data');
128+
}
129+
114130
asyncfunctiontestWriterFail(){
115131
const{ writer,broadcast: bc}=broadcast();
116132
constconsumer=bc.push();
@@ -308,6 +324,7 @@ Promise.all([
308324
testWriteSync(),
309325
testWritevSync(),
310326
testWriterEnd(),
327+
testWriterEndWithPreAbortedSignal(),
311328
testWriterFail(),
312329
testCancelWithoutReason(),
313330
testCancelWithReason(),

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

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -127,6 +127,22 @@ async function testAbortSignal() {
127127
);
128128
}
129129

130+
asyncfunctiontestWriterEndWithPreAbortedSignal(){
131+
const[channelA,channelB]=duplex();
132+
constreason=newError('end aborted');
133+
134+
awaitassert.rejects(
135+
channelA.writer.end({signal: AbortSignal.abort(reason)}),
136+
(error)=>error===reason,
137+
);
138+
139+
awaitchannelA.writer.write('still open');
140+
constcompletedEnd=channelA.writer.end();
141+
assert.strictEqual(awaittext(channelB.readable),'still open');
142+
assert.strictEqual(awaitcompletedEnd,10);
143+
awaitchannelB.close();
144+
}
145+
130146
asyncfunctiontestEmptyDuplex(){
131147
const[channelA,channelB]=duplex();
132148

@@ -182,6 +198,7 @@ Promise.all([
182198
testWithOptions(),
183199
testPerChannelOptions(),
184200
testAbortSignal(),
201+
testWriterEndWithPreAbortedSignal(),
185202
testEmptyDuplex(),
186203
testChannelFail(),
187204
testAbortSignalBothChannels(),

β€Žtest/parallel/test-stream-iter-push-writer.jsβ€Ž

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -260,6 +260,40 @@ async function testEndAsyncReturnValue() {
260260
awaitconsume;
261261
}
262262

263+
asyncfunctiontestEndWithPreAbortedSignal(){
264+
const{ writer, readable }=push();
265+
constreason=newError('end aborted');
266+
267+
writer.writeSync('hello');
268+
awaitassert.rejects(
269+
writer.end({signal: AbortSignal.abort(reason)}),
270+
(error)=>error===reason,
271+
);
272+
273+
// A rejected end must leave the writer open.
274+
writer.writeSync(' world');
275+
constconsume=text(readable);
276+
assert.strictEqual(awaitwriter.end(),11);
277+
assert.strictEqual(awaitconsume,'hello world');
278+
}
279+
280+
asyncfunctiontestEndSignalAbortWhileDraining(){
281+
const{ writer, readable }=push();
282+
constcontroller=newAbortController();
283+
constreason=newError('end aborted while draining');
284+
285+
writer.writeSync('hello');
286+
constabortedEnd=writer.end({signal: controller.signal});
287+
controller.abort(reason);
288+
289+
awaitassert.rejects(abortedEnd,(error)=>error===reason);
290+
291+
// Aborting the operation does not undo the end-of-stream signal.
292+
constcompletedEnd=writer.end();
293+
assert.strictEqual(awaittext(readable),'hello');
294+
assert.strictEqual(awaitcompletedEnd,5);
295+
}
296+
263297
asyncfunctiontestEndAfterEndSyncWaitsForDrain(){
264298
const{ writer, readable }=push();
265299
writer.writeSync('hello');
@@ -553,6 +587,8 @@ Promise.all([
553587
testOndrainProtocolErrorPropagates(),
554588
testFail(),
555589
testEndAsyncReturnValue(),
590+
testEndWithPreAbortedSignal(),
591+
testEndSignalAbortWhileDraining(),
556592
testEndAfterEndSyncWaitsForDrain(),
557593
testWriteUint8Array(),
558594
testOndrainWaitsForDrain(),

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 486cff4

Browse files
trivikraduh95
authored andcommitted
stream: honor AbortSignal in Writer.end()
Reject Writer.end() when its signal is already aborted without closing the writer. For push writers, reject the pending operation if the signal aborts while buffered data drains, while allowing the graceful close to continue. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: codex:gpt-5.6-sol PR-URL: #64727Fixes: #64726 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Gürgün Dayıoğlu <hey@gurgun.day>
1 parent dc30379 commit 486cff4

5 files changed

Lines changed: 108 additions & 5 deletions

File tree

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

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -630,9 +630,10 @@ class BroadcastWriter {
630630
returnfalse;
631631
}
632632

633-
// end() is synchronous internally - signal accepted for interface compliance.
634633
end(options){
635-
getWriterSignal(options);
634+
constsignal=getWriterSignal(options);
635+
if(signal?.aborted)returnPromiseReject(signal.reason);
636+
636637
if(this.#isClosed())returnthis.#closed;
637638
this.#closed =PromiseResolve(this.#totalBytes);
638639
this.#broadcast[kEnd]();

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

Lines changed: 35 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@
88
const{
99
ArrayIsArray,
1010
ArrayPrototypePush,
11+
PromisePrototypeThen,
1112
PromiseReject,
1213
PromiseResolve,
1314
PromiseWithResolvers,
@@ -55,6 +56,35 @@ const {
5556

5657
constkNoFailReason=Symbol('kNoFailReason');
5758

59+
functionraceEndWithSignal(promise,signal){
60+
if(!signal)returnpromise;
61+
if(signal.aborted)returnPromiseReject(signal.reason);
62+
63+
const{
64+
promise: signaledPromise,
65+
resolve,
66+
reject,
67+
}=PromiseWithResolvers();
68+
constonAbort=()=>reject(signal.reason);
69+
70+
signal.addEventListener('abort',onAbort,{
71+
__proto__: null,
72+
once: true,
73+
});
74+
PromisePrototypeThen(
75+
promise,
76+
(value)=>{
77+
signal.removeEventListener('abort',onAbort);
78+
resolve(value);
79+
},
80+
(reason)=>{
81+
signal.removeEventListener('abort',onAbort);
82+
reject(reason);
83+
},
84+
);
85+
returnsignaledPromise;
86+
}
87+
5888
// =============================================================================
5989
// PushQueue - Internal Queue with Chunk-Based Backpressure
6090
// =============================================================================
@@ -628,7 +658,9 @@ class PushWriter {
628658
}
629659

630660
end(options){
631-
getWriterSignal(options);
661+
constsignal=getWriterSignal(options);
662+
if(signal?.aborted)returnPromiseReject(signal.reason);
663+
632664
constresult=this.#queue.end();
633665
if(result===-2){
634666
// Errored: reject with stored error
@@ -639,11 +671,11 @@ class PushWriter {
639671
// when consumer drains past the end sentinel
640672
constpendingEndPromise=this.#queue.pendingEndPromise;
641673
if(pendingEndPromise!==null){
642-
returnpendingEndPromise;
674+
returnraceEndWithSignal(pendingEndPromise,signal);
643675
}
644676
const{ promise, resolve, reject }=PromiseWithResolvers();
645677
this.#queue.setPendingEnd({__proto__: null, promise, resolve, reject });
646-
returnpromise;
678+
returnraceEndWithSignal(promise,signal);
647679
}
648680
// >= 0: byte count (immediate close or idempotent)
649681
returnPromiseResolve(result);

β€Žtest/parallel/test-stream-iter-broadcast-basic.jsβ€Ž

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -111,6 +111,22 @@ async function testWriterEnd() {
111111
assert.strictEqual(data,'data');
112112
}
113113

114+
asyncfunctiontestWriterEndWithPreAbortedSignal(){
115+
const{ writer,broadcast: bc}=broadcast();
116+
constconsumer=bc.push();
117+
constreason=newError('end aborted');
118+
119+
awaitassert.rejects(
120+
writer.end({signal: AbortSignal.abort(reason)}),
121+
(error)=>error===reason,
122+
);
123+
124+
// A rejected end must leave the writer open.
125+
awaitwriter.write('data');
126+
assert.strictEqual(awaitwriter.end(),4);
127+
assert.strictEqual(awaittext(consumer),'data');
128+
}
129+
114130
asyncfunctiontestWriterFail(){
115131
const{ writer,broadcast: bc}=broadcast();
116132
constconsumer=bc.push();
@@ -308,6 +324,7 @@ Promise.all([
308324
testWriteSync(),
309325
testWritevSync(),
310326
testWriterEnd(),
327+
testWriterEndWithPreAbortedSignal(),
311328
testWriterFail(),
312329
testCancelWithoutReason(),
313330
testCancelWithReason(),

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

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -127,6 +127,22 @@ async function testAbortSignal() {
127127
);
128128
}
129129

130+
asyncfunctiontestWriterEndWithPreAbortedSignal(){
131+
const[channelA,channelB]=duplex();
132+
constreason=newError('end aborted');
133+
134+
awaitassert.rejects(
135+
channelA.writer.end({signal: AbortSignal.abort(reason)}),
136+
(error)=>error===reason,
137+
);
138+
139+
awaitchannelA.writer.write('still open');
140+
constcompletedEnd=channelA.writer.end();
141+
assert.strictEqual(awaittext(channelB.readable),'still open');
142+
assert.strictEqual(awaitcompletedEnd,10);
143+
awaitchannelB.close();
144+
}
145+
130146
asyncfunctiontestEmptyDuplex(){
131147
const[channelA,channelB]=duplex();
132148

@@ -182,6 +198,7 @@ Promise.all([
182198
testWithOptions(),
183199
testPerChannelOptions(),
184200
testAbortSignal(),
201+
testWriterEndWithPreAbortedSignal(),
185202
testEmptyDuplex(),
186203
testChannelFail(),
187204
testAbortSignalBothChannels(),

β€Žtest/parallel/test-stream-iter-push-writer.jsβ€Ž

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -260,6 +260,40 @@ async function testEndAsyncReturnValue() {
260260
awaitconsume;
261261
}
262262

263+
asyncfunctiontestEndWithPreAbortedSignal(){
264+
const{ writer, readable }=push();
265+
constreason=newError('end aborted');
266+
267+
writer.writeSync('hello');
268+
awaitassert.rejects(
269+
writer.end({signal: AbortSignal.abort(reason)}),
270+
(error)=>error===reason,
271+
);
272+
273+
// A rejected end must leave the writer open.
274+
writer.writeSync(' world');
275+
constconsume=text(readable);
276+
assert.strictEqual(awaitwriter.end(),11);
277+
assert.strictEqual(awaitconsume,'hello world');
278+
}
279+
280+
asyncfunctiontestEndSignalAbortWhileDraining(){
281+
const{ writer, readable }=push();
282+
constcontroller=newAbortController();
283+
constreason=newError('end aborted while draining');
284+
285+
writer.writeSync('hello');
286+
constabortedEnd=writer.end({signal: controller.signal});
287+
controller.abort(reason);
288+
289+
awaitassert.rejects(abortedEnd,(error)=>error===reason);
290+
291+
// Aborting the operation does not undo the end-of-stream signal.
292+
constcompletedEnd=writer.end();
293+
assert.strictEqual(awaittext(readable),'hello');
294+
assert.strictEqual(awaitcompletedEnd,5);
295+
}
296+
263297
asyncfunctiontestEndAfterEndSyncWaitsForDrain(){
264298
const{ writer, readable }=push();
265299
writer.writeSync('hello');
@@ -553,6 +587,8 @@ Promise.all([
553587
testOndrainProtocolErrorPropagates(),
554588
testFail(),
555589
testEndAsyncReturnValue(),
590+
testEndWithPreAbortedSignal(),
591+
testEndSignalAbortWhileDraining(),
556592
testEndAfterEndSyncWaitsForDrain(),
557593
testWriteUint8Array(),
558594
testOndrainWaitsForDrain(),

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 486cff4

Browse files
trivikraduh95
authored andcommitted
stream: honor AbortSignal in Writer.end()
Reject Writer.end() when its signal is already aborted without closing the writer. For push writers, reject the pending operation if the signal aborts while buffered data drains, while allowing the graceful close to continue. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: codex:gpt-5.6-sol PR-URL: #64727Fixes: #64726 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Gürgün Dayıoğlu <hey@gurgun.day>
1 parent dc30379 commit 486cff4

5 files changed

Lines changed: 108 additions & 5 deletions

File tree

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

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -630,9 +630,10 @@ class BroadcastWriter {
630630
returnfalse;
631631
}
632632

633-
// end() is synchronous internally - signal accepted for interface compliance.
634633
end(options){
635-
getWriterSignal(options);
634+
constsignal=getWriterSignal(options);
635+
if(signal?.aborted)returnPromiseReject(signal.reason);
636+
636637
if(this.#isClosed())returnthis.#closed;
637638
this.#closed =PromiseResolve(this.#totalBytes);
638639
this.#broadcast[kEnd]();

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

Lines changed: 35 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@
88
const{
99
ArrayIsArray,
1010
ArrayPrototypePush,
11+
PromisePrototypeThen,
1112
PromiseReject,
1213
PromiseResolve,
1314
PromiseWithResolvers,
@@ -55,6 +56,35 @@ const {
5556

5657
constkNoFailReason=Symbol('kNoFailReason');
5758

59+
functionraceEndWithSignal(promise,signal){
60+
if(!signal)returnpromise;
61+
if(signal.aborted)returnPromiseReject(signal.reason);
62+
63+
const{
64+
promise: signaledPromise,
65+
resolve,
66+
reject,
67+
}=PromiseWithResolvers();
68+
constonAbort=()=>reject(signal.reason);
69+
70+
signal.addEventListener('abort',onAbort,{
71+
__proto__: null,
72+
once: true,
73+
});
74+
PromisePrototypeThen(
75+
promise,
76+
(value)=>{
77+
signal.removeEventListener('abort',onAbort);
78+
resolve(value);
79+
},
80+
(reason)=>{
81+
signal.removeEventListener('abort',onAbort);
82+
reject(reason);
83+
},
84+
);
85+
returnsignaledPromise;
86+
}
87+
5888
// =============================================================================
5989
// PushQueue - Internal Queue with Chunk-Based Backpressure
6090
// =============================================================================
@@ -628,7 +658,9 @@ class PushWriter {
628658
}
629659

630660
end(options){
631-
getWriterSignal(options);
661+
constsignal=getWriterSignal(options);
662+
if(signal?.aborted)returnPromiseReject(signal.reason);
663+
632664
constresult=this.#queue.end();
633665
if(result===-2){
634666
// Errored: reject with stored error
@@ -639,11 +671,11 @@ class PushWriter {
639671
// when consumer drains past the end sentinel
640672
constpendingEndPromise=this.#queue.pendingEndPromise;
641673
if(pendingEndPromise!==null){
642-
returnpendingEndPromise;
674+
returnraceEndWithSignal(pendingEndPromise,signal);
643675
}
644676
const{ promise, resolve, reject }=PromiseWithResolvers();
645677
this.#queue.setPendingEnd({__proto__: null, promise, resolve, reject });
646-
returnpromise;
678+
returnraceEndWithSignal(promise,signal);
647679
}
648680
// >= 0: byte count (immediate close or idempotent)
649681
returnPromiseResolve(result);

β€Žtest/parallel/test-stream-iter-broadcast-basic.jsβ€Ž

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -111,6 +111,22 @@ async function testWriterEnd() {
111111
assert.strictEqual(data,'data');
112112
}
113113

114+
asyncfunctiontestWriterEndWithPreAbortedSignal(){
115+
const{ writer,broadcast: bc}=broadcast();
116+
constconsumer=bc.push();
117+
constreason=newError('end aborted');
118+
119+
awaitassert.rejects(
120+
writer.end({signal: AbortSignal.abort(reason)}),
121+
(error)=>error===reason,
122+
);
123+
124+
// A rejected end must leave the writer open.
125+
awaitwriter.write('data');
126+
assert.strictEqual(awaitwriter.end(),4);
127+
assert.strictEqual(awaittext(consumer),'data');
128+
}
129+
114130
asyncfunctiontestWriterFail(){
115131
const{ writer,broadcast: bc}=broadcast();
116132
constconsumer=bc.push();
@@ -308,6 +324,7 @@ Promise.all([
308324
testWriteSync(),
309325
testWritevSync(),
310326
testWriterEnd(),
327+
testWriterEndWithPreAbortedSignal(),
311328
testWriterFail(),
312329
testCancelWithoutReason(),
313330
testCancelWithReason(),

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

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -127,6 +127,22 @@ async function testAbortSignal() {
127127
);
128128
}
129129

130+
asyncfunctiontestWriterEndWithPreAbortedSignal(){
131+
const[channelA,channelB]=duplex();
132+
constreason=newError('end aborted');
133+
134+
awaitassert.rejects(
135+
channelA.writer.end({signal: AbortSignal.abort(reason)}),
136+
(error)=>error===reason,
137+
);
138+
139+
awaitchannelA.writer.write('still open');
140+
constcompletedEnd=channelA.writer.end();
141+
assert.strictEqual(awaittext(channelB.readable),'still open');
142+
assert.strictEqual(awaitcompletedEnd,10);
143+
awaitchannelB.close();
144+
}
145+
130146
asyncfunctiontestEmptyDuplex(){
131147
const[channelA,channelB]=duplex();
132148

@@ -182,6 +198,7 @@ Promise.all([
182198
testWithOptions(),
183199
testPerChannelOptions(),
184200
testAbortSignal(),
201+
testWriterEndWithPreAbortedSignal(),
185202
testEmptyDuplex(),
186203
testChannelFail(),
187204
testAbortSignalBothChannels(),

β€Žtest/parallel/test-stream-iter-push-writer.jsβ€Ž

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -260,6 +260,40 @@ async function testEndAsyncReturnValue() {
260260
awaitconsume;
261261
}
262262

263+
asyncfunctiontestEndWithPreAbortedSignal(){
264+
const{ writer, readable }=push();
265+
constreason=newError('end aborted');
266+
267+
writer.writeSync('hello');
268+
awaitassert.rejects(
269+
writer.end({signal: AbortSignal.abort(reason)}),
270+
(error)=>error===reason,
271+
);
272+
273+
// A rejected end must leave the writer open.
274+
writer.writeSync(' world');
275+
constconsume=text(readable);
276+
assert.strictEqual(awaitwriter.end(),11);
277+
assert.strictEqual(awaitconsume,'hello world');
278+
}
279+
280+
asyncfunctiontestEndSignalAbortWhileDraining(){
281+
const{ writer, readable }=push();
282+
constcontroller=newAbortController();
283+
constreason=newError('end aborted while draining');
284+
285+
writer.writeSync('hello');
286+
constabortedEnd=writer.end({signal: controller.signal});
287+
controller.abort(reason);
288+
289+
awaitassert.rejects(abortedEnd,(error)=>error===reason);
290+
291+
// Aborting the operation does not undo the end-of-stream signal.
292+
constcompletedEnd=writer.end();
293+
assert.strictEqual(awaittext(readable),'hello');
294+
assert.strictEqual(awaitcompletedEnd,5);
295+
}
296+
263297
asyncfunctiontestEndAfterEndSyncWaitsForDrain(){
264298
const{ writer, readable }=push();
265299
writer.writeSync('hello');
@@ -553,6 +587,8 @@ Promise.all([
553587
testOndrainProtocolErrorPropagates(),
554588
testFail(),
555589
testEndAsyncReturnValue(),
590+
testEndWithPreAbortedSignal(),
591+
testEndSignalAbortWhileDraining(),
556592
testEndAfterEndSyncWaitsForDrain(),
557593
testWriteUint8Array(),
558594
testOndrainWaitsForDrain(),

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 486cff4

Browse files
trivikraduh95
authored andcommitted
stream: honor AbortSignal in Writer.end()
Reject Writer.end() when its signal is already aborted without closing the writer. For push writers, reject the pending operation if the signal aborts while buffered data drains, while allowing the graceful close to continue. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: codex:gpt-5.6-sol PR-URL: #64727Fixes: #64726 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Gürgün Dayıoğlu <hey@gurgun.day>
1 parent dc30379 commit 486cff4

5 files changed

Lines changed: 108 additions & 5 deletions

File tree

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

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -630,9 +630,10 @@ class BroadcastWriter {
630630
returnfalse;
631631
}
632632

633-
// end() is synchronous internally - signal accepted for interface compliance.
634633
end(options){
635-
getWriterSignal(options);
634+
constsignal=getWriterSignal(options);
635+
if(signal?.aborted)returnPromiseReject(signal.reason);
636+
636637
if(this.#isClosed())returnthis.#closed;
637638
this.#closed =PromiseResolve(this.#totalBytes);
638639
this.#broadcast[kEnd]();

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

Lines changed: 35 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@
88
const{
99
ArrayIsArray,
1010
ArrayPrototypePush,
11+
PromisePrototypeThen,
1112
PromiseReject,
1213
PromiseResolve,
1314
PromiseWithResolvers,
@@ -55,6 +56,35 @@ const {
5556

5657
constkNoFailReason=Symbol('kNoFailReason');
5758

59+
functionraceEndWithSignal(promise,signal){
60+
if(!signal)returnpromise;
61+
if(signal.aborted)returnPromiseReject(signal.reason);
62+
63+
const{
64+
promise: signaledPromise,
65+
resolve,
66+
reject,
67+
}=PromiseWithResolvers();
68+
constonAbort=()=>reject(signal.reason);
69+
70+
signal.addEventListener('abort',onAbort,{
71+
__proto__: null,
72+
once: true,
73+
});
74+
PromisePrototypeThen(
75+
promise,
76+
(value)=>{
77+
signal.removeEventListener('abort',onAbort);
78+
resolve(value);
79+
},
80+
(reason)=>{
81+
signal.removeEventListener('abort',onAbort);
82+
reject(reason);
83+
},
84+
);
85+
returnsignaledPromise;
86+
}
87+
5888
// =============================================================================
5989
// PushQueue - Internal Queue with Chunk-Based Backpressure
6090
// =============================================================================
@@ -628,7 +658,9 @@ class PushWriter {
628658
}
629659

630660
end(options){
631-
getWriterSignal(options);
661+
constsignal=getWriterSignal(options);
662+
if(signal?.aborted)returnPromiseReject(signal.reason);
663+
632664
constresult=this.#queue.end();
633665
if(result===-2){
634666
// Errored: reject with stored error
@@ -639,11 +671,11 @@ class PushWriter {
639671
// when consumer drains past the end sentinel
640672
constpendingEndPromise=this.#queue.pendingEndPromise;
641673
if(pendingEndPromise!==null){
642-
returnpendingEndPromise;
674+
returnraceEndWithSignal(pendingEndPromise,signal);
643675
}
644676
const{ promise, resolve, reject }=PromiseWithResolvers();
645677
this.#queue.setPendingEnd({__proto__: null, promise, resolve, reject });
646-
returnpromise;
678+
returnraceEndWithSignal(promise,signal);
647679
}
648680
// >= 0: byte count (immediate close or idempotent)
649681
returnPromiseResolve(result);

β€Žtest/parallel/test-stream-iter-broadcast-basic.jsβ€Ž

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -111,6 +111,22 @@ async function testWriterEnd() {
111111
assert.strictEqual(data,'data');
112112
}
113113

114+
asyncfunctiontestWriterEndWithPreAbortedSignal(){
115+
const{ writer,broadcast: bc}=broadcast();
116+
constconsumer=bc.push();
117+
constreason=newError('end aborted');
118+
119+
awaitassert.rejects(
120+
writer.end({signal: AbortSignal.abort(reason)}),
121+
(error)=>error===reason,
122+
);
123+
124+
// A rejected end must leave the writer open.
125+
awaitwriter.write('data');
126+
assert.strictEqual(awaitwriter.end(),4);
127+
assert.strictEqual(awaittext(consumer),'data');
128+
}
129+
114130
asyncfunctiontestWriterFail(){
115131
const{ writer,broadcast: bc}=broadcast();
116132
constconsumer=bc.push();
@@ -308,6 +324,7 @@ Promise.all([
308324
testWriteSync(),
309325
testWritevSync(),
310326
testWriterEnd(),
327+
testWriterEndWithPreAbortedSignal(),
311328
testWriterFail(),
312329
testCancelWithoutReason(),
313330
testCancelWithReason(),

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

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -127,6 +127,22 @@ async function testAbortSignal() {
127127
);
128128
}
129129

130+
asyncfunctiontestWriterEndWithPreAbortedSignal(){
131+
const[channelA,channelB]=duplex();
132+
constreason=newError('end aborted');
133+
134+
awaitassert.rejects(
135+
channelA.writer.end({signal: AbortSignal.abort(reason)}),
136+
(error)=>error===reason,
137+
);
138+
139+
awaitchannelA.writer.write('still open');
140+
constcompletedEnd=channelA.writer.end();
141+
assert.strictEqual(awaittext(channelB.readable),'still open');
142+
assert.strictEqual(awaitcompletedEnd,10);
143+
awaitchannelB.close();
144+
}
145+
130146
asyncfunctiontestEmptyDuplex(){
131147
const[channelA,channelB]=duplex();
132148

@@ -182,6 +198,7 @@ Promise.all([
182198
testWithOptions(),
183199
testPerChannelOptions(),
184200
testAbortSignal(),
201+
testWriterEndWithPreAbortedSignal(),
185202
testEmptyDuplex(),
186203
testChannelFail(),
187204
testAbortSignalBothChannels(),

β€Žtest/parallel/test-stream-iter-push-writer.jsβ€Ž

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -260,6 +260,40 @@ async function testEndAsyncReturnValue() {
260260
awaitconsume;
261261
}
262262

263+
asyncfunctiontestEndWithPreAbortedSignal(){
264+
const{ writer, readable }=push();
265+
constreason=newError('end aborted');
266+
267+
writer.writeSync('hello');
268+
awaitassert.rejects(
269+
writer.end({signal: AbortSignal.abort(reason)}),
270+
(error)=>error===reason,
271+
);
272+
273+
// A rejected end must leave the writer open.
274+
writer.writeSync(' world');
275+
constconsume=text(readable);
276+
assert.strictEqual(awaitwriter.end(),11);
277+
assert.strictEqual(awaitconsume,'hello world');
278+
}
279+
280+
asyncfunctiontestEndSignalAbortWhileDraining(){
281+
const{ writer, readable }=push();
282+
constcontroller=newAbortController();
283+
constreason=newError('end aborted while draining');
284+
285+
writer.writeSync('hello');
286+
constabortedEnd=writer.end({signal: controller.signal});
287+
controller.abort(reason);
288+
289+
awaitassert.rejects(abortedEnd,(error)=>error===reason);
290+
291+
// Aborting the operation does not undo the end-of-stream signal.
292+
constcompletedEnd=writer.end();
293+
assert.strictEqual(awaittext(readable),'hello');
294+
assert.strictEqual(awaitcompletedEnd,5);
295+
}
296+
263297
asyncfunctiontestEndAfterEndSyncWaitsForDrain(){
264298
const{ writer, readable }=push();
265299
writer.writeSync('hello');
@@ -553,6 +587,8 @@ Promise.all([
553587
testOndrainProtocolErrorPropagates(),
554588
testFail(),
555589
testEndAsyncReturnValue(),
590+
testEndWithPreAbortedSignal(),
591+
testEndSignalAbortWhileDraining(),
556592
testEndAfterEndSyncWaitsForDrain(),
557593
testWriteUint8Array(),
558594
testOndrainWaitsForDrain(),

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 486cff4

Browse files
trivikraduh95
authored andcommitted
stream: honor AbortSignal in Writer.end()
Reject Writer.end() when its signal is already aborted without closing the writer. For push writers, reject the pending operation if the signal aborts while buffered data drains, while allowing the graceful close to continue. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: codex:gpt-5.6-sol PR-URL: #64727Fixes: #64726 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Gürgün Dayıoğlu <hey@gurgun.day>
1 parent dc30379 commit 486cff4

5 files changed

Lines changed: 108 additions & 5 deletions

File tree

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

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -630,9 +630,10 @@ class BroadcastWriter {
630630
returnfalse;
631631
}
632632

633-
// end() is synchronous internally - signal accepted for interface compliance.
634633
end(options){
635-
getWriterSignal(options);
634+
constsignal=getWriterSignal(options);
635+
if(signal?.aborted)returnPromiseReject(signal.reason);
636+
636637
if(this.#isClosed())returnthis.#closed;
637638
this.#closed =PromiseResolve(this.#totalBytes);
638639
this.#broadcast[kEnd]();

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

Lines changed: 35 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@
88
const{
99
ArrayIsArray,
1010
ArrayPrototypePush,
11+
PromisePrototypeThen,
1112
PromiseReject,
1213
PromiseResolve,
1314
PromiseWithResolvers,
@@ -55,6 +56,35 @@ const {
5556

5657
constkNoFailReason=Symbol('kNoFailReason');
5758

59+
functionraceEndWithSignal(promise,signal){
60+
if(!signal)returnpromise;
61+
if(signal.aborted)returnPromiseReject(signal.reason);
62+
63+
const{
64+
promise: signaledPromise,
65+
resolve,
66+
reject,
67+
}=PromiseWithResolvers();
68+
constonAbort=()=>reject(signal.reason);
69+
70+
signal.addEventListener('abort',onAbort,{
71+
__proto__: null,
72+
once: true,
73+
});
74+
PromisePrototypeThen(
75+
promise,
76+
(value)=>{
77+
signal.removeEventListener('abort',onAbort);
78+
resolve(value);
79+
},
80+
(reason)=>{
81+
signal.removeEventListener('abort',onAbort);
82+
reject(reason);
83+
},
84+
);
85+
returnsignaledPromise;
86+
}
87+
5888
// =============================================================================
5989
// PushQueue - Internal Queue with Chunk-Based Backpressure
6090
// =============================================================================
@@ -628,7 +658,9 @@ class PushWriter {
628658
}
629659

630660
end(options){
631-
getWriterSignal(options);
661+
constsignal=getWriterSignal(options);
662+
if(signal?.aborted)returnPromiseReject(signal.reason);
663+
632664
constresult=this.#queue.end();
633665
if(result===-2){
634666
// Errored: reject with stored error
@@ -639,11 +671,11 @@ class PushWriter {
639671
// when consumer drains past the end sentinel
640672
constpendingEndPromise=this.#queue.pendingEndPromise;
641673
if(pendingEndPromise!==null){
642-
returnpendingEndPromise;
674+
returnraceEndWithSignal(pendingEndPromise,signal);
643675
}
644676
const{ promise, resolve, reject }=PromiseWithResolvers();
645677
this.#queue.setPendingEnd({__proto__: null, promise, resolve, reject });
646-
returnpromise;
678+
returnraceEndWithSignal(promise,signal);
647679
}
648680
// >= 0: byte count (immediate close or idempotent)
649681
returnPromiseResolve(result);

β€Žtest/parallel/test-stream-iter-broadcast-basic.jsβ€Ž

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -111,6 +111,22 @@ async function testWriterEnd() {
111111
assert.strictEqual(data,'data');
112112
}
113113

114+
asyncfunctiontestWriterEndWithPreAbortedSignal(){
115+
const{ writer,broadcast: bc}=broadcast();
116+
constconsumer=bc.push();
117+
constreason=newError('end aborted');
118+
119+
awaitassert.rejects(
120+
writer.end({signal: AbortSignal.abort(reason)}),
121+
(error)=>error===reason,
122+
);
123+
124+
// A rejected end must leave the writer open.
125+
awaitwriter.write('data');
126+
assert.strictEqual(awaitwriter.end(),4);
127+
assert.strictEqual(awaittext(consumer),'data');
128+
}
129+
114130
asyncfunctiontestWriterFail(){
115131
const{ writer,broadcast: bc}=broadcast();
116132
constconsumer=bc.push();
@@ -308,6 +324,7 @@ Promise.all([
308324
testWriteSync(),
309325
testWritevSync(),
310326
testWriterEnd(),
327+
testWriterEndWithPreAbortedSignal(),
311328
testWriterFail(),
312329
testCancelWithoutReason(),
313330
testCancelWithReason(),

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

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -127,6 +127,22 @@ async function testAbortSignal() {
127127
);
128128
}
129129

130+
asyncfunctiontestWriterEndWithPreAbortedSignal(){
131+
const[channelA,channelB]=duplex();
132+
constreason=newError('end aborted');
133+
134+
awaitassert.rejects(
135+
channelA.writer.end({signal: AbortSignal.abort(reason)}),
136+
(error)=>error===reason,
137+
);
138+
139+
awaitchannelA.writer.write('still open');
140+
constcompletedEnd=channelA.writer.end();
141+
assert.strictEqual(awaittext(channelB.readable),'still open');
142+
assert.strictEqual(awaitcompletedEnd,10);
143+
awaitchannelB.close();
144+
}
145+
130146
asyncfunctiontestEmptyDuplex(){
131147
const[channelA,channelB]=duplex();
132148

@@ -182,6 +198,7 @@ Promise.all([
182198
testWithOptions(),
183199
testPerChannelOptions(),
184200
testAbortSignal(),
201+
testWriterEndWithPreAbortedSignal(),
185202
testEmptyDuplex(),
186203
testChannelFail(),
187204
testAbortSignalBothChannels(),

β€Žtest/parallel/test-stream-iter-push-writer.jsβ€Ž

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -260,6 +260,40 @@ async function testEndAsyncReturnValue() {
260260
awaitconsume;
261261
}
262262

263+
asyncfunctiontestEndWithPreAbortedSignal(){
264+
const{ writer, readable }=push();
265+
constreason=newError('end aborted');
266+
267+
writer.writeSync('hello');
268+
awaitassert.rejects(
269+
writer.end({signal: AbortSignal.abort(reason)}),
270+
(error)=>error===reason,
271+
);
272+
273+
// A rejected end must leave the writer open.
274+
writer.writeSync(' world');
275+
constconsume=text(readable);
276+
assert.strictEqual(awaitwriter.end(),11);
277+
assert.strictEqual(awaitconsume,'hello world');
278+
}
279+
280+
asyncfunctiontestEndSignalAbortWhileDraining(){
281+
const{ writer, readable }=push();
282+
constcontroller=newAbortController();
283+
constreason=newError('end aborted while draining');
284+
285+
writer.writeSync('hello');
286+
constabortedEnd=writer.end({signal: controller.signal});
287+
controller.abort(reason);
288+
289+
awaitassert.rejects(abortedEnd,(error)=>error===reason);
290+
291+
// Aborting the operation does not undo the end-of-stream signal.
292+
constcompletedEnd=writer.end();
293+
assert.strictEqual(awaittext(readable),'hello');
294+
assert.strictEqual(awaitcompletedEnd,5);
295+
}
296+
263297
asyncfunctiontestEndAfterEndSyncWaitsForDrain(){
264298
const{ writer, readable }=push();
265299
writer.writeSync('hello');
@@ -553,6 +587,8 @@ Promise.all([
553587
testOndrainProtocolErrorPropagates(),
554588
testFail(),
555589
testEndAsyncReturnValue(),
590+
testEndWithPreAbortedSignal(),
591+
testEndSignalAbortWhileDraining(),
556592
testEndAfterEndSyncWaitsForDrain(),
557593
testWriteUint8Array(),
558594
testOndrainWaitsForDrain(),

0 commit comments

Comments
Β (0)