Commit 0876a29

Browse files
trivikraduh95
authored andcommitted
stream: preserve falsy cancellation reasons
Use undefined as the no-error sentinel when cancelling broadcast and share consumers. This ensures that 0, an empty string, false, and null are propagated instead of being converted to clean completion. Make sync share surface cancellation reasons before handling detached consumers, and add regression coverage for async and sync consumers. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: codex:gpt-5.6-sol PR-URL: #64705Fixes: #64704 Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent ab50ae6 commit 0876a29

5 files changed

Lines changed: 56 additions & 44 deletions

File tree

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

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -88,7 +88,7 @@ class BroadcastImpl {
8888
#consumers =newSafeSet();
8989
#waiters =[];// Consumers with pending resolve (subset of #consumers)
9090
#ended =false;
91-
#error=null;
91+
#error;
9292
#cancelled =false;
9393
#options;
9494
#writer =null;
@@ -177,7 +177,9 @@ class BroadcastImpl {
177177
__proto__: null,
178178
next(){
179179
if(state.detached){
180-
if(self.#error)returnPromiseReject(self.#error);
180+
if(self.#error !==undefined){
181+
returnPromiseReject(self.#error);
182+
}
181183
returnkDone;
182184
}
183185

@@ -194,7 +196,7 @@ class BroadcastImpl {
194196
{__proto__: null,done: false,value: chunk});
195197
}
196198

197-
if(self.#error){
199+
if(self.#error!==undefined){
198200
state.detached=true;
199201
self.#deleteConsumer(state);
200202
returnPromiseReject(self.#error);
@@ -344,7 +346,7 @@ class BroadcastImpl {
344346
}
345347

346348
[kAbort](reason){
347-
if(this.#ended ||this.#error)return;
349+
if(this.#ended ||this.#error!==undefined)return;
348350
this.#error =reason;
349351
this.#ended =true;
350352

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

Lines changed: 14 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -73,7 +73,7 @@ class ShareImpl {
7373
#consumers =newSafeSet();
7474
#sourceIterator =null;
7575
#sourceExhausted =false;
76-
#sourceError=null;
76+
#sourceError;
7777
#cancelled =false;
7878
#pulling =false;
7979
#pullWaiters =[];
@@ -129,7 +129,7 @@ class ShareImpl {
129129
__proto__: null,
130130
[SymbolAsyncIterator](){
131131
constgetNext=async()=>{
132-
if(self.#sourceError){
132+
if(self.#sourceError!==undefined){
133133
state.detached=true;
134134
self.#consumers.delete(state);
135135
throwself.#sourceError;
@@ -141,7 +141,7 @@ class ShareImpl {
141141
// cursor must re-pull rather than terminating prematurely.
142142
for(;;){
143143
if(state.detached){
144-
if(self.#sourceError)throwself.#sourceError;
144+
if(self.#sourceError!==undefined)throwself.#sourceError;
145145
return{__proto__: null,done: true,value: undefined};
146146
}
147147

@@ -167,7 +167,7 @@ class ShareImpl {
167167
if(self.#sourceExhausted){
168168
state.detached=true;
169169
self.#deleteConsumer(state);
170-
if(self.#sourceError)throwself.#sourceError;
170+
if(self.#sourceError!==undefined)throwself.#sourceError;
171171
return{__proto__: null,done: true,value: undefined};
172172
}
173173

@@ -176,7 +176,7 @@ class ShareImpl {
176176
if(shouldBuffer===null){
177177
state.detached=true;
178178
self.#deleteConsumer(state);
179-
if(self.#sourceError)throwself.#sourceError;
179+
if(self.#sourceError!==undefined)throwself.#sourceError;
180180
return{__proto__: null,done: true,value: undefined};
181181
}
182182

@@ -260,7 +260,9 @@ class ShareImpl {
260260

261261
async #waitForBufferSpace(){
262262
while(this.#bufferedBytes >=this.#options.budget){
263-
if(this.#cancelled ||this.#sourceError ||this.#sourceExhausted){
263+
if(this.#cancelled ||
264+
this.#sourceError !==undefined||
265+
this.#sourceExhausted){
264266
returnthis.#cancelled ? null : true;
265267
}
266268

@@ -418,7 +420,7 @@ class SyncShareImpl {
418420
#consumers =newSafeSet();
419421
#sourceIterator =null;
420422
#sourceExhausted =false;
421-
#sourceError=null;
423+
#sourceError;
422424
#cancelled =false;
423425
#cachedMinCursor =0;
424426
#cachedMinCursorConsumers =0;
@@ -467,14 +469,14 @@ class SyncShareImpl {
467469
return{
468470
__proto__: null,
469471
next(){
470-
if(state.detached){
471-
return{__proto__: null,done: true,value: undefined};
472-
}
473-
if(self.#sourceError){
472+
if(self.#sourceError !==undefined){
474473
state.detached=true;
475474
self.#deleteConsumer(state);
476475
throwself.#sourceError;
477476
}
477+
if(state.detached){
478+
return{__proto__: null,done: true,value: undefined};
479+
}
478480
if(self.#cancelled){
479481
state.detached=true;
480482
self.#deleteConsumer(state);
@@ -535,7 +537,7 @@ class SyncShareImpl {
535537

536538
self.#pullFromSource();
537539

538-
if(self.#sourceError){
540+
if(self.#sourceError!==undefined){
539541
state.detached=true;
540542
self.#deleteConsumer(state);
541543
throwself.#sourceError;

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

Lines changed: 8 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -260,15 +260,15 @@ async function testWriterFailIdempotent() {
260260
},{message: 'fail!'});
261261
}
262262

263-
// cancel() with falsy reason (0, "", false) should still treat as error
264263
asyncfunctiontestCancelWithFalsyReason(){
265-
const{broadcast: bc}=broadcast();
266-
constconsumer=bc.push();
267-
constresultPromise=text(consumer).catch((err)=>err);
268-
awaitnewPromise((resolve)=>setImmediate(resolve));
269-
bc.cancel(0);
270-
constresult=awaitresultPromise;
271-
assert.strictEqual(result,0);
264+
for(constreasonof[0,'',false,null]){
265+
const{broadcast: bc}=broadcast();
266+
constiterator=bc.push()[Symbol.asyncIterator]();
267+
268+
bc.cancel(reason);
269+
270+
awaitassert.rejects(iterator.next(),(error)=>error===reason);
271+
}
272272
}
273273

274274
// Late-joining consumer should read from oldest buffered entry

β€Žtest/parallel/test-stream-iter-share-async.jsβ€Ž

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -132,6 +132,17 @@ async function testShareCancelWithReason() {
132132
);
133133
}
134134

135+
asyncfunctiontestShareCancelWithFalsyReason(){
136+
for(constreasonof[0,'',false,null]){
137+
constshared=share(from('data'));
138+
constiterator=shared.pull()[Symbol.asyncIterator]();
139+
140+
shared.cancel(reason);
141+
142+
awaitassert.rejects(iterator.next(),(error)=>error===reason);
143+
}
144+
}
145+
135146
asyncfunctiontestShareAbortSignal(){
136147
constac=newAbortController();
137148
constreason=newError('share aborted');
@@ -357,6 +368,7 @@ Promise.all([
357368
testShareCancel(),
358369
testShareCancelMidIteration(),
359370
testShareCancelWithReason(),
371+
testShareCancelWithFalsyReason(),
360372
testShareAbortSignal(),
361373
testShareAbortSignalWhileSourcePullPending(),
362374
testSharePullAbortSignalRejectsPendingNext(),

β€Žtest/parallel/test-stream-iter-share-sync.jsβ€Ž

Lines changed: 16 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -87,37 +87,32 @@ function testShareSyncCancelMidIteration() {
8787
}
8888

8989
functiontestShareSyncCancelWithReason(){
90-
// When cancel(reason) is called, a consumer that hasn't started
91-
// iterating is already detached, so it sees done:true (not the error).
92-
// But a consumer that is mid-iteration when another consumer cancels
93-
// with a reason will see the error on the next pull after cancel.
9490
constenc=newTextEncoder();
9591
function*gen(){
9692
yield[enc.encode('a')];
9793
yield[enc.encode('b')];
98-
yield[enc.encode('c')];
9994
}
10095
constshared=shareSync(gen(),{budget: 16384});
101-
constc1=shared.pull();
102-
constc2=shared.pull();
96+
constiterator1=shared.pull()[Symbol.iterator]();
97+
constiterator2=shared.pull()[Symbol.iterator]();
98+
constreason=newError('sync cancel reason');
99+
100+
iterator1.next();
101+
shared.cancel(reason);
103102

104-
// c1 reads one item, then c2 cancels with a reason
105-
constiter1=c1[Symbol.iterator]();
106-
constfirst=iter1.next();
107-
assert.strictEqual(first.done,false);
103+
assert.throws(()=>iterator1.next(),(error)=>error===reason);
104+
assert.throws(()=>iterator2.next(),(error)=>error===reason);
105+
}
108106

109-
shared.cancel(newError('sync cancel reason'));
107+
functiontestShareSyncCancelWithFalsyReason(){
108+
for(constreasonof[0,'',false,null]){
109+
constshared=shareSync(fromSync('data'));
110+
constiterator=shared.pull()[Symbol.iterator]();
110111

111-
// c1 was already iterating, it's now detached β†’ done
112-
constnext=iter1.next();
113-
assert.strictEqual(next.done,true);
112+
shared.cancel(reason);
114113

115-
// c2 never started, also detached β†’ done (not error)
116-
constbatches=[];
117-
for(constbatchofc2){
118-
batches.push(batch);
114+
assert.throws(()=>iterator.next(),(error)=>error===reason);
119115
}
120-
assert.strictEqual(batches.length,0);
121116
}
122117

123118
// =============================================================================
@@ -157,6 +152,7 @@ Promise.all([
157152
testShareSyncCancel(),
158153
testShareSyncCancelMidIteration(),
159154
testShareSyncCancelWithReason(),
155+
testShareSyncCancelWithFalsyReason(),
160156
testShareSyncSourceError(),
161157
testShareSyncStringSource(),
162158
]).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 0876a29

Browse files
trivikraduh95
authored andcommitted
stream: preserve falsy cancellation reasons
Use undefined as the no-error sentinel when cancelling broadcast and share consumers. This ensures that 0, an empty string, false, and null are propagated instead of being converted to clean completion. Make sync share surface cancellation reasons before handling detached consumers, and add regression coverage for async and sync consumers. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: codex:gpt-5.6-sol PR-URL: #64705Fixes: #64704 Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent ab50ae6 commit 0876a29

5 files changed

Lines changed: 56 additions & 44 deletions

File tree

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

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -88,7 +88,7 @@ class BroadcastImpl {
8888
#consumers =newSafeSet();
8989
#waiters =[];// Consumers with pending resolve (subset of #consumers)
9090
#ended =false;
91-
#error=null;
91+
#error;
9292
#cancelled =false;
9393
#options;
9494
#writer =null;
@@ -177,7 +177,9 @@ class BroadcastImpl {
177177
__proto__: null,
178178
next(){
179179
if(state.detached){
180-
if(self.#error)returnPromiseReject(self.#error);
180+
if(self.#error !==undefined){
181+
returnPromiseReject(self.#error);
182+
}
181183
returnkDone;
182184
}
183185

@@ -194,7 +196,7 @@ class BroadcastImpl {
194196
{__proto__: null,done: false,value: chunk});
195197
}
196198

197-
if(self.#error){
199+
if(self.#error!==undefined){
198200
state.detached=true;
199201
self.#deleteConsumer(state);
200202
returnPromiseReject(self.#error);
@@ -344,7 +346,7 @@ class BroadcastImpl {
344346
}
345347

346348
[kAbort](reason){
347-
if(this.#ended ||this.#error)return;
349+
if(this.#ended ||this.#error!==undefined)return;
348350
this.#error =reason;
349351
this.#ended =true;
350352

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

Lines changed: 14 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -73,7 +73,7 @@ class ShareImpl {
7373
#consumers =newSafeSet();
7474
#sourceIterator =null;
7575
#sourceExhausted =false;
76-
#sourceError=null;
76+
#sourceError;
7777
#cancelled =false;
7878
#pulling =false;
7979
#pullWaiters =[];
@@ -129,7 +129,7 @@ class ShareImpl {
129129
__proto__: null,
130130
[SymbolAsyncIterator](){
131131
constgetNext=async()=>{
132-
if(self.#sourceError){
132+
if(self.#sourceError!==undefined){
133133
state.detached=true;
134134
self.#consumers.delete(state);
135135
throwself.#sourceError;
@@ -141,7 +141,7 @@ class ShareImpl {
141141
// cursor must re-pull rather than terminating prematurely.
142142
for(;;){
143143
if(state.detached){
144-
if(self.#sourceError)throwself.#sourceError;
144+
if(self.#sourceError!==undefined)throwself.#sourceError;
145145
return{__proto__: null,done: true,value: undefined};
146146
}
147147

@@ -167,7 +167,7 @@ class ShareImpl {
167167
if(self.#sourceExhausted){
168168
state.detached=true;
169169
self.#deleteConsumer(state);
170-
if(self.#sourceError)throwself.#sourceError;
170+
if(self.#sourceError!==undefined)throwself.#sourceError;
171171
return{__proto__: null,done: true,value: undefined};
172172
}
173173

@@ -176,7 +176,7 @@ class ShareImpl {
176176
if(shouldBuffer===null){
177177
state.detached=true;
178178
self.#deleteConsumer(state);
179-
if(self.#sourceError)throwself.#sourceError;
179+
if(self.#sourceError!==undefined)throwself.#sourceError;
180180
return{__proto__: null,done: true,value: undefined};
181181
}
182182

@@ -260,7 +260,9 @@ class ShareImpl {
260260

261261
async #waitForBufferSpace(){
262262
while(this.#bufferedBytes >=this.#options.budget){
263-
if(this.#cancelled ||this.#sourceError ||this.#sourceExhausted){
263+
if(this.#cancelled ||
264+
this.#sourceError !==undefined||
265+
this.#sourceExhausted){
264266
returnthis.#cancelled ? null : true;
265267
}
266268

@@ -418,7 +420,7 @@ class SyncShareImpl {
418420
#consumers =newSafeSet();
419421
#sourceIterator =null;
420422
#sourceExhausted =false;
421-
#sourceError=null;
423+
#sourceError;
422424
#cancelled =false;
423425
#cachedMinCursor =0;
424426
#cachedMinCursorConsumers =0;
@@ -467,14 +469,14 @@ class SyncShareImpl {
467469
return{
468470
__proto__: null,
469471
next(){
470-
if(state.detached){
471-
return{__proto__: null,done: true,value: undefined};
472-
}
473-
if(self.#sourceError){
472+
if(self.#sourceError !==undefined){
474473
state.detached=true;
475474
self.#deleteConsumer(state);
476475
throwself.#sourceError;
477476
}
477+
if(state.detached){
478+
return{__proto__: null,done: true,value: undefined};
479+
}
478480
if(self.#cancelled){
479481
state.detached=true;
480482
self.#deleteConsumer(state);
@@ -535,7 +537,7 @@ class SyncShareImpl {
535537

536538
self.#pullFromSource();
537539

538-
if(self.#sourceError){
540+
if(self.#sourceError!==undefined){
539541
state.detached=true;
540542
self.#deleteConsumer(state);
541543
throwself.#sourceError;

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

Lines changed: 8 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -260,15 +260,15 @@ async function testWriterFailIdempotent() {
260260
},{message: 'fail!'});
261261
}
262262

263-
// cancel() with falsy reason (0, "", false) should still treat as error
264263
asyncfunctiontestCancelWithFalsyReason(){
265-
const{broadcast: bc}=broadcast();
266-
constconsumer=bc.push();
267-
constresultPromise=text(consumer).catch((err)=>err);
268-
awaitnewPromise((resolve)=>setImmediate(resolve));
269-
bc.cancel(0);
270-
constresult=awaitresultPromise;
271-
assert.strictEqual(result,0);
264+
for(constreasonof[0,'',false,null]){
265+
const{broadcast: bc}=broadcast();
266+
constiterator=bc.push()[Symbol.asyncIterator]();
267+
268+
bc.cancel(reason);
269+
270+
awaitassert.rejects(iterator.next(),(error)=>error===reason);
271+
}
272272
}
273273

274274
// Late-joining consumer should read from oldest buffered entry

β€Žtest/parallel/test-stream-iter-share-async.jsβ€Ž

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -132,6 +132,17 @@ async function testShareCancelWithReason() {
132132
);
133133
}
134134

135+
asyncfunctiontestShareCancelWithFalsyReason(){
136+
for(constreasonof[0,'',false,null]){
137+
constshared=share(from('data'));
138+
constiterator=shared.pull()[Symbol.asyncIterator]();
139+
140+
shared.cancel(reason);
141+
142+
awaitassert.rejects(iterator.next(),(error)=>error===reason);
143+
}
144+
}
145+
135146
asyncfunctiontestShareAbortSignal(){
136147
constac=newAbortController();
137148
constreason=newError('share aborted');
@@ -357,6 +368,7 @@ Promise.all([
357368
testShareCancel(),
358369
testShareCancelMidIteration(),
359370
testShareCancelWithReason(),
371+
testShareCancelWithFalsyReason(),
360372
testShareAbortSignal(),
361373
testShareAbortSignalWhileSourcePullPending(),
362374
testSharePullAbortSignalRejectsPendingNext(),

β€Žtest/parallel/test-stream-iter-share-sync.jsβ€Ž

Lines changed: 16 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -87,37 +87,32 @@ function testShareSyncCancelMidIteration() {
8787
}
8888

8989
functiontestShareSyncCancelWithReason(){
90-
// When cancel(reason) is called, a consumer that hasn't started
91-
// iterating is already detached, so it sees done:true (not the error).
92-
// But a consumer that is mid-iteration when another consumer cancels
93-
// with a reason will see the error on the next pull after cancel.
9490
constenc=newTextEncoder();
9591
function*gen(){
9692
yield[enc.encode('a')];
9793
yield[enc.encode('b')];
98-
yield[enc.encode('c')];
9994
}
10095
constshared=shareSync(gen(),{budget: 16384});
101-
constc1=shared.pull();
102-
constc2=shared.pull();
96+
constiterator1=shared.pull()[Symbol.iterator]();
97+
constiterator2=shared.pull()[Symbol.iterator]();
98+
constreason=newError('sync cancel reason');
99+
100+
iterator1.next();
101+
shared.cancel(reason);
103102

104-
// c1 reads one item, then c2 cancels with a reason
105-
constiter1=c1[Symbol.iterator]();
106-
constfirst=iter1.next();
107-
assert.strictEqual(first.done,false);
103+
assert.throws(()=>iterator1.next(),(error)=>error===reason);
104+
assert.throws(()=>iterator2.next(),(error)=>error===reason);
105+
}
108106

109-
shared.cancel(newError('sync cancel reason'));
107+
functiontestShareSyncCancelWithFalsyReason(){
108+
for(constreasonof[0,'',false,null]){
109+
constshared=shareSync(fromSync('data'));
110+
constiterator=shared.pull()[Symbol.iterator]();
110111

111-
// c1 was already iterating, it's now detached β†’ done
112-
constnext=iter1.next();
113-
assert.strictEqual(next.done,true);
112+
shared.cancel(reason);
114113

115-
// c2 never started, also detached β†’ done (not error)
116-
constbatches=[];
117-
for(constbatchofc2){
118-
batches.push(batch);
114+
assert.throws(()=>iterator.next(),(error)=>error===reason);
119115
}
120-
assert.strictEqual(batches.length,0);
121116
}
122117

123118
// =============================================================================
@@ -157,6 +152,7 @@ Promise.all([
157152
testShareSyncCancel(),
158153
testShareSyncCancelMidIteration(),
159154
testShareSyncCancelWithReason(),
155+
testShareSyncCancelWithFalsyReason(),
160156
testShareSyncSourceError(),
161157
testShareSyncStringSource(),
162158
]).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 0876a29

Browse files
trivikraduh95
authored andcommitted
stream: preserve falsy cancellation reasons
Use undefined as the no-error sentinel when cancelling broadcast and share consumers. This ensures that 0, an empty string, false, and null are propagated instead of being converted to clean completion. Make sync share surface cancellation reasons before handling detached consumers, and add regression coverage for async and sync consumers. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: codex:gpt-5.6-sol PR-URL: #64705Fixes: #64704 Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent ab50ae6 commit 0876a29

5 files changed

Lines changed: 56 additions & 44 deletions

File tree

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

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -88,7 +88,7 @@ class BroadcastImpl {
8888
#consumers =newSafeSet();
8989
#waiters =[];// Consumers with pending resolve (subset of #consumers)
9090
#ended =false;
91-
#error=null;
91+
#error;
9292
#cancelled =false;
9393
#options;
9494
#writer =null;
@@ -177,7 +177,9 @@ class BroadcastImpl {
177177
__proto__: null,
178178
next(){
179179
if(state.detached){
180-
if(self.#error)returnPromiseReject(self.#error);
180+
if(self.#error !==undefined){
181+
returnPromiseReject(self.#error);
182+
}
181183
returnkDone;
182184
}
183185

@@ -194,7 +196,7 @@ class BroadcastImpl {
194196
{__proto__: null,done: false,value: chunk});
195197
}
196198

197-
if(self.#error){
199+
if(self.#error!==undefined){
198200
state.detached=true;
199201
self.#deleteConsumer(state);
200202
returnPromiseReject(self.#error);
@@ -344,7 +346,7 @@ class BroadcastImpl {
344346
}
345347

346348
[kAbort](reason){
347-
if(this.#ended ||this.#error)return;
349+
if(this.#ended ||this.#error!==undefined)return;
348350
this.#error =reason;
349351
this.#ended =true;
350352

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

Lines changed: 14 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -73,7 +73,7 @@ class ShareImpl {
7373
#consumers =newSafeSet();
7474
#sourceIterator =null;
7575
#sourceExhausted =false;
76-
#sourceError=null;
76+
#sourceError;
7777
#cancelled =false;
7878
#pulling =false;
7979
#pullWaiters =[];
@@ -129,7 +129,7 @@ class ShareImpl {
129129
__proto__: null,
130130
[SymbolAsyncIterator](){
131131
constgetNext=async()=>{
132-
if(self.#sourceError){
132+
if(self.#sourceError!==undefined){
133133
state.detached=true;
134134
self.#consumers.delete(state);
135135
throwself.#sourceError;
@@ -141,7 +141,7 @@ class ShareImpl {
141141
// cursor must re-pull rather than terminating prematurely.
142142
for(;;){
143143
if(state.detached){
144-
if(self.#sourceError)throwself.#sourceError;
144+
if(self.#sourceError!==undefined)throwself.#sourceError;
145145
return{__proto__: null,done: true,value: undefined};
146146
}
147147

@@ -167,7 +167,7 @@ class ShareImpl {
167167
if(self.#sourceExhausted){
168168
state.detached=true;
169169
self.#deleteConsumer(state);
170-
if(self.#sourceError)throwself.#sourceError;
170+
if(self.#sourceError!==undefined)throwself.#sourceError;
171171
return{__proto__: null,done: true,value: undefined};
172172
}
173173

@@ -176,7 +176,7 @@ class ShareImpl {
176176
if(shouldBuffer===null){
177177
state.detached=true;
178178
self.#deleteConsumer(state);
179-
if(self.#sourceError)throwself.#sourceError;
179+
if(self.#sourceError!==undefined)throwself.#sourceError;
180180
return{__proto__: null,done: true,value: undefined};
181181
}
182182

@@ -260,7 +260,9 @@ class ShareImpl {
260260

261261
async #waitForBufferSpace(){
262262
while(this.#bufferedBytes >=this.#options.budget){
263-
if(this.#cancelled ||this.#sourceError ||this.#sourceExhausted){
263+
if(this.#cancelled ||
264+
this.#sourceError !==undefined||
265+
this.#sourceExhausted){
264266
returnthis.#cancelled ? null : true;
265267
}
266268

@@ -418,7 +420,7 @@ class SyncShareImpl {
418420
#consumers =newSafeSet();
419421
#sourceIterator =null;
420422
#sourceExhausted =false;
421-
#sourceError=null;
423+
#sourceError;
422424
#cancelled =false;
423425
#cachedMinCursor =0;
424426
#cachedMinCursorConsumers =0;
@@ -467,14 +469,14 @@ class SyncShareImpl {
467469
return{
468470
__proto__: null,
469471
next(){
470-
if(state.detached){
471-
return{__proto__: null,done: true,value: undefined};
472-
}
473-
if(self.#sourceError){
472+
if(self.#sourceError !==undefined){
474473
state.detached=true;
475474
self.#deleteConsumer(state);
476475
throwself.#sourceError;
477476
}
477+
if(state.detached){
478+
return{__proto__: null,done: true,value: undefined};
479+
}
478480
if(self.#cancelled){
479481
state.detached=true;
480482
self.#deleteConsumer(state);
@@ -535,7 +537,7 @@ class SyncShareImpl {
535537

536538
self.#pullFromSource();
537539

538-
if(self.#sourceError){
540+
if(self.#sourceError!==undefined){
539541
state.detached=true;
540542
self.#deleteConsumer(state);
541543
throwself.#sourceError;

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

Lines changed: 8 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -260,15 +260,15 @@ async function testWriterFailIdempotent() {
260260
},{message: 'fail!'});
261261
}
262262

263-
// cancel() with falsy reason (0, "", false) should still treat as error
264263
asyncfunctiontestCancelWithFalsyReason(){
265-
const{broadcast: bc}=broadcast();
266-
constconsumer=bc.push();
267-
constresultPromise=text(consumer).catch((err)=>err);
268-
awaitnewPromise((resolve)=>setImmediate(resolve));
269-
bc.cancel(0);
270-
constresult=awaitresultPromise;
271-
assert.strictEqual(result,0);
264+
for(constreasonof[0,'',false,null]){
265+
const{broadcast: bc}=broadcast();
266+
constiterator=bc.push()[Symbol.asyncIterator]();
267+
268+
bc.cancel(reason);
269+
270+
awaitassert.rejects(iterator.next(),(error)=>error===reason);
271+
}
272272
}
273273

274274
// Late-joining consumer should read from oldest buffered entry

β€Žtest/parallel/test-stream-iter-share-async.jsβ€Ž

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -132,6 +132,17 @@ async function testShareCancelWithReason() {
132132
);
133133
}
134134

135+
asyncfunctiontestShareCancelWithFalsyReason(){
136+
for(constreasonof[0,'',false,null]){
137+
constshared=share(from('data'));
138+
constiterator=shared.pull()[Symbol.asyncIterator]();
139+
140+
shared.cancel(reason);
141+
142+
awaitassert.rejects(iterator.next(),(error)=>error===reason);
143+
}
144+
}
145+
135146
asyncfunctiontestShareAbortSignal(){
136147
constac=newAbortController();
137148
constreason=newError('share aborted');
@@ -357,6 +368,7 @@ Promise.all([
357368
testShareCancel(),
358369
testShareCancelMidIteration(),
359370
testShareCancelWithReason(),
371+
testShareCancelWithFalsyReason(),
360372
testShareAbortSignal(),
361373
testShareAbortSignalWhileSourcePullPending(),
362374
testSharePullAbortSignalRejectsPendingNext(),

β€Žtest/parallel/test-stream-iter-share-sync.jsβ€Ž

Lines changed: 16 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -87,37 +87,32 @@ function testShareSyncCancelMidIteration() {
8787
}
8888

8989
functiontestShareSyncCancelWithReason(){
90-
// When cancel(reason) is called, a consumer that hasn't started
91-
// iterating is already detached, so it sees done:true (not the error).
92-
// But a consumer that is mid-iteration when another consumer cancels
93-
// with a reason will see the error on the next pull after cancel.
9490
constenc=newTextEncoder();
9591
function*gen(){
9692
yield[enc.encode('a')];
9793
yield[enc.encode('b')];
98-
yield[enc.encode('c')];
9994
}
10095
constshared=shareSync(gen(),{budget: 16384});
101-
constc1=shared.pull();
102-
constc2=shared.pull();
96+
constiterator1=shared.pull()[Symbol.iterator]();
97+
constiterator2=shared.pull()[Symbol.iterator]();
98+
constreason=newError('sync cancel reason');
99+
100+
iterator1.next();
101+
shared.cancel(reason);
103102

104-
// c1 reads one item, then c2 cancels with a reason
105-
constiter1=c1[Symbol.iterator]();
106-
constfirst=iter1.next();
107-
assert.strictEqual(first.done,false);
103+
assert.throws(()=>iterator1.next(),(error)=>error===reason);
104+
assert.throws(()=>iterator2.next(),(error)=>error===reason);
105+
}
108106

109-
shared.cancel(newError('sync cancel reason'));
107+
functiontestShareSyncCancelWithFalsyReason(){
108+
for(constreasonof[0,'',false,null]){
109+
constshared=shareSync(fromSync('data'));
110+
constiterator=shared.pull()[Symbol.iterator]();
110111

111-
// c1 was already iterating, it's now detached β†’ done
112-
constnext=iter1.next();
113-
assert.strictEqual(next.done,true);
112+
shared.cancel(reason);
114113

115-
// c2 never started, also detached β†’ done (not error)
116-
constbatches=[];
117-
for(constbatchofc2){
118-
batches.push(batch);
114+
assert.throws(()=>iterator.next(),(error)=>error===reason);
119115
}
120-
assert.strictEqual(batches.length,0);
121116
}
122117

123118
// =============================================================================
@@ -157,6 +152,7 @@ Promise.all([
157152
testShareSyncCancel(),
158153
testShareSyncCancelMidIteration(),
159154
testShareSyncCancelWithReason(),
155+
testShareSyncCancelWithFalsyReason(),
160156
testShareSyncSourceError(),
161157
testShareSyncStringSource(),
162158
]).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 0876a29

Browse files
trivikraduh95
authored andcommitted
stream: preserve falsy cancellation reasons
Use undefined as the no-error sentinel when cancelling broadcast and share consumers. This ensures that 0, an empty string, false, and null are propagated instead of being converted to clean completion. Make sync share surface cancellation reasons before handling detached consumers, and add regression coverage for async and sync consumers. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: codex:gpt-5.6-sol PR-URL: #64705Fixes: #64704 Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent ab50ae6 commit 0876a29

5 files changed

Lines changed: 56 additions & 44 deletions

File tree

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

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -88,7 +88,7 @@ class BroadcastImpl {
8888
#consumers =newSafeSet();
8989
#waiters =[];// Consumers with pending resolve (subset of #consumers)
9090
#ended =false;
91-
#error=null;
91+
#error;
9292
#cancelled =false;
9393
#options;
9494
#writer =null;
@@ -177,7 +177,9 @@ class BroadcastImpl {
177177
__proto__: null,
178178
next(){
179179
if(state.detached){
180-
if(self.#error)returnPromiseReject(self.#error);
180+
if(self.#error !==undefined){
181+
returnPromiseReject(self.#error);
182+
}
181183
returnkDone;
182184
}
183185

@@ -194,7 +196,7 @@ class BroadcastImpl {
194196
{__proto__: null,done: false,value: chunk});
195197
}
196198

197-
if(self.#error){
199+
if(self.#error!==undefined){
198200
state.detached=true;
199201
self.#deleteConsumer(state);
200202
returnPromiseReject(self.#error);
@@ -344,7 +346,7 @@ class BroadcastImpl {
344346
}
345347

346348
[kAbort](reason){
347-
if(this.#ended ||this.#error)return;
349+
if(this.#ended ||this.#error!==undefined)return;
348350
this.#error =reason;
349351
this.#ended =true;
350352

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

Lines changed: 14 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -73,7 +73,7 @@ class ShareImpl {
7373
#consumers =newSafeSet();
7474
#sourceIterator =null;
7575
#sourceExhausted =false;
76-
#sourceError=null;
76+
#sourceError;
7777
#cancelled =false;
7878
#pulling =false;
7979
#pullWaiters =[];
@@ -129,7 +129,7 @@ class ShareImpl {
129129
__proto__: null,
130130
[SymbolAsyncIterator](){
131131
constgetNext=async()=>{
132-
if(self.#sourceError){
132+
if(self.#sourceError!==undefined){
133133
state.detached=true;
134134
self.#consumers.delete(state);
135135
throwself.#sourceError;
@@ -141,7 +141,7 @@ class ShareImpl {
141141
// cursor must re-pull rather than terminating prematurely.
142142
for(;;){
143143
if(state.detached){
144-
if(self.#sourceError)throwself.#sourceError;
144+
if(self.#sourceError!==undefined)throwself.#sourceError;
145145
return{__proto__: null,done: true,value: undefined};
146146
}
147147

@@ -167,7 +167,7 @@ class ShareImpl {
167167
if(self.#sourceExhausted){
168168
state.detached=true;
169169
self.#deleteConsumer(state);
170-
if(self.#sourceError)throwself.#sourceError;
170+
if(self.#sourceError!==undefined)throwself.#sourceError;
171171
return{__proto__: null,done: true,value: undefined};
172172
}
173173

@@ -176,7 +176,7 @@ class ShareImpl {
176176
if(shouldBuffer===null){
177177
state.detached=true;
178178
self.#deleteConsumer(state);
179-
if(self.#sourceError)throwself.#sourceError;
179+
if(self.#sourceError!==undefined)throwself.#sourceError;
180180
return{__proto__: null,done: true,value: undefined};
181181
}
182182

@@ -260,7 +260,9 @@ class ShareImpl {
260260

261261
async #waitForBufferSpace(){
262262
while(this.#bufferedBytes >=this.#options.budget){
263-
if(this.#cancelled ||this.#sourceError ||this.#sourceExhausted){
263+
if(this.#cancelled ||
264+
this.#sourceError !==undefined||
265+
this.#sourceExhausted){
264266
returnthis.#cancelled ? null : true;
265267
}
266268

@@ -418,7 +420,7 @@ class SyncShareImpl {
418420
#consumers =newSafeSet();
419421
#sourceIterator =null;
420422
#sourceExhausted =false;
421-
#sourceError=null;
423+
#sourceError;
422424
#cancelled =false;
423425
#cachedMinCursor =0;
424426
#cachedMinCursorConsumers =0;
@@ -467,14 +469,14 @@ class SyncShareImpl {
467469
return{
468470
__proto__: null,
469471
next(){
470-
if(state.detached){
471-
return{__proto__: null,done: true,value: undefined};
472-
}
473-
if(self.#sourceError){
472+
if(self.#sourceError !==undefined){
474473
state.detached=true;
475474
self.#deleteConsumer(state);
476475
throwself.#sourceError;
477476
}
477+
if(state.detached){
478+
return{__proto__: null,done: true,value: undefined};
479+
}
478480
if(self.#cancelled){
479481
state.detached=true;
480482
self.#deleteConsumer(state);
@@ -535,7 +537,7 @@ class SyncShareImpl {
535537

536538
self.#pullFromSource();
537539

538-
if(self.#sourceError){
540+
if(self.#sourceError!==undefined){
539541
state.detached=true;
540542
self.#deleteConsumer(state);
541543
throwself.#sourceError;

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

Lines changed: 8 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -260,15 +260,15 @@ async function testWriterFailIdempotent() {
260260
},{message: 'fail!'});
261261
}
262262

263-
// cancel() with falsy reason (0, "", false) should still treat as error
264263
asyncfunctiontestCancelWithFalsyReason(){
265-
const{broadcast: bc}=broadcast();
266-
constconsumer=bc.push();
267-
constresultPromise=text(consumer).catch((err)=>err);
268-
awaitnewPromise((resolve)=>setImmediate(resolve));
269-
bc.cancel(0);
270-
constresult=awaitresultPromise;
271-
assert.strictEqual(result,0);
264+
for(constreasonof[0,'',false,null]){
265+
const{broadcast: bc}=broadcast();
266+
constiterator=bc.push()[Symbol.asyncIterator]();
267+
268+
bc.cancel(reason);
269+
270+
awaitassert.rejects(iterator.next(),(error)=>error===reason);
271+
}
272272
}
273273

274274
// Late-joining consumer should read from oldest buffered entry

β€Žtest/parallel/test-stream-iter-share-async.jsβ€Ž

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -132,6 +132,17 @@ async function testShareCancelWithReason() {
132132
);
133133
}
134134

135+
asyncfunctiontestShareCancelWithFalsyReason(){
136+
for(constreasonof[0,'',false,null]){
137+
constshared=share(from('data'));
138+
constiterator=shared.pull()[Symbol.asyncIterator]();
139+
140+
shared.cancel(reason);
141+
142+
awaitassert.rejects(iterator.next(),(error)=>error===reason);
143+
}
144+
}
145+
135146
asyncfunctiontestShareAbortSignal(){
136147
constac=newAbortController();
137148
constreason=newError('share aborted');
@@ -357,6 +368,7 @@ Promise.all([
357368
testShareCancel(),
358369
testShareCancelMidIteration(),
359370
testShareCancelWithReason(),
371+
testShareCancelWithFalsyReason(),
360372
testShareAbortSignal(),
361373
testShareAbortSignalWhileSourcePullPending(),
362374
testSharePullAbortSignalRejectsPendingNext(),

β€Žtest/parallel/test-stream-iter-share-sync.jsβ€Ž

Lines changed: 16 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -87,37 +87,32 @@ function testShareSyncCancelMidIteration() {
8787
}
8888

8989
functiontestShareSyncCancelWithReason(){
90-
// When cancel(reason) is called, a consumer that hasn't started
91-
// iterating is already detached, so it sees done:true (not the error).
92-
// But a consumer that is mid-iteration when another consumer cancels
93-
// with a reason will see the error on the next pull after cancel.
9490
constenc=newTextEncoder();
9591
function*gen(){
9692
yield[enc.encode('a')];
9793
yield[enc.encode('b')];
98-
yield[enc.encode('c')];
9994
}
10095
constshared=shareSync(gen(),{budget: 16384});
101-
constc1=shared.pull();
102-
constc2=shared.pull();
96+
constiterator1=shared.pull()[Symbol.iterator]();
97+
constiterator2=shared.pull()[Symbol.iterator]();
98+
constreason=newError('sync cancel reason');
99+
100+
iterator1.next();
101+
shared.cancel(reason);
103102

104-
// c1 reads one item, then c2 cancels with a reason
105-
constiter1=c1[Symbol.iterator]();
106-
constfirst=iter1.next();
107-
assert.strictEqual(first.done,false);
103+
assert.throws(()=>iterator1.next(),(error)=>error===reason);
104+
assert.throws(()=>iterator2.next(),(error)=>error===reason);
105+
}
108106

109-
shared.cancel(newError('sync cancel reason'));
107+
functiontestShareSyncCancelWithFalsyReason(){
108+
for(constreasonof[0,'',false,null]){
109+
constshared=shareSync(fromSync('data'));
110+
constiterator=shared.pull()[Symbol.iterator]();
110111

111-
// c1 was already iterating, it's now detached β†’ done
112-
constnext=iter1.next();
113-
assert.strictEqual(next.done,true);
112+
shared.cancel(reason);
114113

115-
// c2 never started, also detached β†’ done (not error)
116-
constbatches=[];
117-
for(constbatchofc2){
118-
batches.push(batch);
114+
assert.throws(()=>iterator.next(),(error)=>error===reason);
119115
}
120-
assert.strictEqual(batches.length,0);
121116
}
122117

123118
// =============================================================================
@@ -157,6 +152,7 @@ Promise.all([
157152
testShareSyncCancel(),
158153
testShareSyncCancelMidIteration(),
159154
testShareSyncCancelWithReason(),
155+
testShareSyncCancelWithFalsyReason(),
160156
testShareSyncSourceError(),
161157
testShareSyncStringSource(),
162158
]).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 0876a29

Browse files
trivikraduh95
authored andcommitted
stream: preserve falsy cancellation reasons
Use undefined as the no-error sentinel when cancelling broadcast and share consumers. This ensures that 0, an empty string, false, and null are propagated instead of being converted to clean completion. Make sync share surface cancellation reasons before handling detached consumers, and add regression coverage for async and sync consumers. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: codex:gpt-5.6-sol PR-URL: #64705Fixes: #64704 Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent ab50ae6 commit 0876a29

5 files changed

Lines changed: 56 additions & 44 deletions

File tree

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

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -88,7 +88,7 @@ class BroadcastImpl {
8888
#consumers =newSafeSet();
8989
#waiters =[];// Consumers with pending resolve (subset of #consumers)
9090
#ended =false;
91-
#error=null;
91+
#error;
9292
#cancelled =false;
9393
#options;
9494
#writer =null;
@@ -177,7 +177,9 @@ class BroadcastImpl {
177177
__proto__: null,
178178
next(){
179179
if(state.detached){
180-
if(self.#error)returnPromiseReject(self.#error);
180+
if(self.#error !==undefined){
181+
returnPromiseReject(self.#error);
182+
}
181183
returnkDone;
182184
}
183185

@@ -194,7 +196,7 @@ class BroadcastImpl {
194196
{__proto__: null,done: false,value: chunk});
195197
}
196198

197-
if(self.#error){
199+
if(self.#error!==undefined){
198200
state.detached=true;
199201
self.#deleteConsumer(state);
200202
returnPromiseReject(self.#error);
@@ -344,7 +346,7 @@ class BroadcastImpl {
344346
}
345347

346348
[kAbort](reason){
347-
if(this.#ended ||this.#error)return;
349+
if(this.#ended ||this.#error!==undefined)return;
348350
this.#error =reason;
349351
this.#ended =true;
350352

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

Lines changed: 14 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -73,7 +73,7 @@ class ShareImpl {
7373
#consumers =newSafeSet();
7474
#sourceIterator =null;
7575
#sourceExhausted =false;
76-
#sourceError=null;
76+
#sourceError;
7777
#cancelled =false;
7878
#pulling =false;
7979
#pullWaiters =[];
@@ -129,7 +129,7 @@ class ShareImpl {
129129
__proto__: null,
130130
[SymbolAsyncIterator](){
131131
constgetNext=async()=>{
132-
if(self.#sourceError){
132+
if(self.#sourceError!==undefined){
133133
state.detached=true;
134134
self.#consumers.delete(state);
135135
throwself.#sourceError;
@@ -141,7 +141,7 @@ class ShareImpl {
141141
// cursor must re-pull rather than terminating prematurely.
142142
for(;;){
143143
if(state.detached){
144-
if(self.#sourceError)throwself.#sourceError;
144+
if(self.#sourceError!==undefined)throwself.#sourceError;
145145
return{__proto__: null,done: true,value: undefined};
146146
}
147147

@@ -167,7 +167,7 @@ class ShareImpl {
167167
if(self.#sourceExhausted){
168168
state.detached=true;
169169
self.#deleteConsumer(state);
170-
if(self.#sourceError)throwself.#sourceError;
170+
if(self.#sourceError!==undefined)throwself.#sourceError;
171171
return{__proto__: null,done: true,value: undefined};
172172
}
173173

@@ -176,7 +176,7 @@ class ShareImpl {
176176
if(shouldBuffer===null){
177177
state.detached=true;
178178
self.#deleteConsumer(state);
179-
if(self.#sourceError)throwself.#sourceError;
179+
if(self.#sourceError!==undefined)throwself.#sourceError;
180180
return{__proto__: null,done: true,value: undefined};
181181
}
182182

@@ -260,7 +260,9 @@ class ShareImpl {
260260

261261
async #waitForBufferSpace(){
262262
while(this.#bufferedBytes >=this.#options.budget){
263-
if(this.#cancelled ||this.#sourceError ||this.#sourceExhausted){
263+
if(this.#cancelled ||
264+
this.#sourceError !==undefined||
265+
this.#sourceExhausted){
264266
returnthis.#cancelled ? null : true;
265267
}
266268

@@ -418,7 +420,7 @@ class SyncShareImpl {
418420
#consumers =newSafeSet();
419421
#sourceIterator =null;
420422
#sourceExhausted =false;
421-
#sourceError=null;
423+
#sourceError;
422424
#cancelled =false;
423425
#cachedMinCursor =0;
424426
#cachedMinCursorConsumers =0;
@@ -467,14 +469,14 @@ class SyncShareImpl {
467469
return{
468470
__proto__: null,
469471
next(){
470-
if(state.detached){
471-
return{__proto__: null,done: true,value: undefined};
472-
}
473-
if(self.#sourceError){
472+
if(self.#sourceError !==undefined){
474473
state.detached=true;
475474
self.#deleteConsumer(state);
476475
throwself.#sourceError;
477476
}
477+
if(state.detached){
478+
return{__proto__: null,done: true,value: undefined};
479+
}
478480
if(self.#cancelled){
479481
state.detached=true;
480482
self.#deleteConsumer(state);
@@ -535,7 +537,7 @@ class SyncShareImpl {
535537

536538
self.#pullFromSource();
537539

538-
if(self.#sourceError){
540+
if(self.#sourceError!==undefined){
539541
state.detached=true;
540542
self.#deleteConsumer(state);
541543
throwself.#sourceError;

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

Lines changed: 8 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -260,15 +260,15 @@ async function testWriterFailIdempotent() {
260260
},{message: 'fail!'});
261261
}
262262

263-
// cancel() with falsy reason (0, "", false) should still treat as error
264263
asyncfunctiontestCancelWithFalsyReason(){
265-
const{broadcast: bc}=broadcast();
266-
constconsumer=bc.push();
267-
constresultPromise=text(consumer).catch((err)=>err);
268-
awaitnewPromise((resolve)=>setImmediate(resolve));
269-
bc.cancel(0);
270-
constresult=awaitresultPromise;
271-
assert.strictEqual(result,0);
264+
for(constreasonof[0,'',false,null]){
265+
const{broadcast: bc}=broadcast();
266+
constiterator=bc.push()[Symbol.asyncIterator]();
267+
268+
bc.cancel(reason);
269+
270+
awaitassert.rejects(iterator.next(),(error)=>error===reason);
271+
}
272272
}
273273

274274
// Late-joining consumer should read from oldest buffered entry

β€Žtest/parallel/test-stream-iter-share-async.jsβ€Ž

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -132,6 +132,17 @@ async function testShareCancelWithReason() {
132132
);
133133
}
134134

135+
asyncfunctiontestShareCancelWithFalsyReason(){
136+
for(constreasonof[0,'',false,null]){
137+
constshared=share(from('data'));
138+
constiterator=shared.pull()[Symbol.asyncIterator]();
139+
140+
shared.cancel(reason);
141+
142+
awaitassert.rejects(iterator.next(),(error)=>error===reason);
143+
}
144+
}
145+
135146
asyncfunctiontestShareAbortSignal(){
136147
constac=newAbortController();
137148
constreason=newError('share aborted');
@@ -357,6 +368,7 @@ Promise.all([
357368
testShareCancel(),
358369
testShareCancelMidIteration(),
359370
testShareCancelWithReason(),
371+
testShareCancelWithFalsyReason(),
360372
testShareAbortSignal(),
361373
testShareAbortSignalWhileSourcePullPending(),
362374
testSharePullAbortSignalRejectsPendingNext(),

β€Žtest/parallel/test-stream-iter-share-sync.jsβ€Ž

Lines changed: 16 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -87,37 +87,32 @@ function testShareSyncCancelMidIteration() {
8787
}
8888

8989
functiontestShareSyncCancelWithReason(){
90-
// When cancel(reason) is called, a consumer that hasn't started
91-
// iterating is already detached, so it sees done:true (not the error).
92-
// But a consumer that is mid-iteration when another consumer cancels
93-
// with a reason will see the error on the next pull after cancel.
9490
constenc=newTextEncoder();
9591
function*gen(){
9692
yield[enc.encode('a')];
9793
yield[enc.encode('b')];
98-
yield[enc.encode('c')];
9994
}
10095
constshared=shareSync(gen(),{budget: 16384});
101-
constc1=shared.pull();
102-
constc2=shared.pull();
96+
constiterator1=shared.pull()[Symbol.iterator]();
97+
constiterator2=shared.pull()[Symbol.iterator]();
98+
constreason=newError('sync cancel reason');
99+
100+
iterator1.next();
101+
shared.cancel(reason);
103102

104-
// c1 reads one item, then c2 cancels with a reason
105-
constiter1=c1[Symbol.iterator]();
106-
constfirst=iter1.next();
107-
assert.strictEqual(first.done,false);
103+
assert.throws(()=>iterator1.next(),(error)=>error===reason);
104+
assert.throws(()=>iterator2.next(),(error)=>error===reason);
105+
}
108106

109-
shared.cancel(newError('sync cancel reason'));
107+
functiontestShareSyncCancelWithFalsyReason(){
108+
for(constreasonof[0,'',false,null]){
109+
constshared=shareSync(fromSync('data'));
110+
constiterator=shared.pull()[Symbol.iterator]();
110111

111-
// c1 was already iterating, it's now detached β†’ done
112-
constnext=iter1.next();
113-
assert.strictEqual(next.done,true);
112+
shared.cancel(reason);
114113

115-
// c2 never started, also detached β†’ done (not error)
116-
constbatches=[];
117-
for(constbatchofc2){
118-
batches.push(batch);
114+
assert.throws(()=>iterator.next(),(error)=>error===reason);
119115
}
120-
assert.strictEqual(batches.length,0);
121116
}
122117

123118
// =============================================================================
@@ -157,6 +152,7 @@ Promise.all([
157152
testShareSyncCancel(),
158153
testShareSyncCancelMidIteration(),
159154
testShareSyncCancelWithReason(),
155+
testShareSyncCancelWithFalsyReason(),
160156
testShareSyncSourceError(),
161157
testShareSyncStringSource(),
162158
]).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 0876a29

Browse files
trivikraduh95
authored andcommitted
stream: preserve falsy cancellation reasons
Use undefined as the no-error sentinel when cancelling broadcast and share consumers. This ensures that 0, an empty string, false, and null are propagated instead of being converted to clean completion. Make sync share surface cancellation reasons before handling detached consumers, and add regression coverage for async and sync consumers. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: codex:gpt-5.6-sol PR-URL: #64705Fixes: #64704 Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent ab50ae6 commit 0876a29

5 files changed

Lines changed: 56 additions & 44 deletions

File tree

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

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -88,7 +88,7 @@ class BroadcastImpl {
8888
#consumers =newSafeSet();
8989
#waiters =[];// Consumers with pending resolve (subset of #consumers)
9090
#ended =false;
91-
#error=null;
91+
#error;
9292
#cancelled =false;
9393
#options;
9494
#writer =null;
@@ -177,7 +177,9 @@ class BroadcastImpl {
177177
__proto__: null,
178178
next(){
179179
if(state.detached){
180-
if(self.#error)returnPromiseReject(self.#error);
180+
if(self.#error !==undefined){
181+
returnPromiseReject(self.#error);
182+
}
181183
returnkDone;
182184
}
183185

@@ -194,7 +196,7 @@ class BroadcastImpl {
194196
{__proto__: null,done: false,value: chunk});
195197
}
196198

197-
if(self.#error){
199+
if(self.#error!==undefined){
198200
state.detached=true;
199201
self.#deleteConsumer(state);
200202
returnPromiseReject(self.#error);
@@ -344,7 +346,7 @@ class BroadcastImpl {
344346
}
345347

346348
[kAbort](reason){
347-
if(this.#ended ||this.#error)return;
349+
if(this.#ended ||this.#error!==undefined)return;
348350
this.#error =reason;
349351
this.#ended =true;
350352

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

Lines changed: 14 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -73,7 +73,7 @@ class ShareImpl {
7373
#consumers =newSafeSet();
7474
#sourceIterator =null;
7575
#sourceExhausted =false;
76-
#sourceError=null;
76+
#sourceError;
7777
#cancelled =false;
7878
#pulling =false;
7979
#pullWaiters =[];
@@ -129,7 +129,7 @@ class ShareImpl {
129129
__proto__: null,
130130
[SymbolAsyncIterator](){
131131
constgetNext=async()=>{
132-
if(self.#sourceError){
132+
if(self.#sourceError!==undefined){
133133
state.detached=true;
134134
self.#consumers.delete(state);
135135
throwself.#sourceError;
@@ -141,7 +141,7 @@ class ShareImpl {
141141
// cursor must re-pull rather than terminating prematurely.
142142
for(;;){
143143
if(state.detached){
144-
if(self.#sourceError)throwself.#sourceError;
144+
if(self.#sourceError!==undefined)throwself.#sourceError;
145145
return{__proto__: null,done: true,value: undefined};
146146
}
147147

@@ -167,7 +167,7 @@ class ShareImpl {
167167
if(self.#sourceExhausted){
168168
state.detached=true;
169169
self.#deleteConsumer(state);
170-
if(self.#sourceError)throwself.#sourceError;
170+
if(self.#sourceError!==undefined)throwself.#sourceError;
171171
return{__proto__: null,done: true,value: undefined};
172172
}
173173

@@ -176,7 +176,7 @@ class ShareImpl {
176176
if(shouldBuffer===null){
177177
state.detached=true;
178178
self.#deleteConsumer(state);
179-
if(self.#sourceError)throwself.#sourceError;
179+
if(self.#sourceError!==undefined)throwself.#sourceError;
180180
return{__proto__: null,done: true,value: undefined};
181181
}
182182

@@ -260,7 +260,9 @@ class ShareImpl {
260260

261261
async #waitForBufferSpace(){
262262
while(this.#bufferedBytes >=this.#options.budget){
263-
if(this.#cancelled ||this.#sourceError ||this.#sourceExhausted){
263+
if(this.#cancelled ||
264+
this.#sourceError !==undefined||
265+
this.#sourceExhausted){
264266
returnthis.#cancelled ? null : true;
265267
}
266268

@@ -418,7 +420,7 @@ class SyncShareImpl {
418420
#consumers =newSafeSet();
419421
#sourceIterator =null;
420422
#sourceExhausted =false;
421-
#sourceError=null;
423+
#sourceError;
422424
#cancelled =false;
423425
#cachedMinCursor =0;
424426
#cachedMinCursorConsumers =0;
@@ -467,14 +469,14 @@ class SyncShareImpl {
467469
return{
468470
__proto__: null,
469471
next(){
470-
if(state.detached){
471-
return{__proto__: null,done: true,value: undefined};
472-
}
473-
if(self.#sourceError){
472+
if(self.#sourceError !==undefined){
474473
state.detached=true;
475474
self.#deleteConsumer(state);
476475
throwself.#sourceError;
477476
}
477+
if(state.detached){
478+
return{__proto__: null,done: true,value: undefined};
479+
}
478480
if(self.#cancelled){
479481
state.detached=true;
480482
self.#deleteConsumer(state);
@@ -535,7 +537,7 @@ class SyncShareImpl {
535537

536538
self.#pullFromSource();
537539

538-
if(self.#sourceError){
540+
if(self.#sourceError!==undefined){
539541
state.detached=true;
540542
self.#deleteConsumer(state);
541543
throwself.#sourceError;

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

Lines changed: 8 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -260,15 +260,15 @@ async function testWriterFailIdempotent() {
260260
},{message: 'fail!'});
261261
}
262262

263-
// cancel() with falsy reason (0, "", false) should still treat as error
264263
asyncfunctiontestCancelWithFalsyReason(){
265-
const{broadcast: bc}=broadcast();
266-
constconsumer=bc.push();
267-
constresultPromise=text(consumer).catch((err)=>err);
268-
awaitnewPromise((resolve)=>setImmediate(resolve));
269-
bc.cancel(0);
270-
constresult=awaitresultPromise;
271-
assert.strictEqual(result,0);
264+
for(constreasonof[0,'',false,null]){
265+
const{broadcast: bc}=broadcast();
266+
constiterator=bc.push()[Symbol.asyncIterator]();
267+
268+
bc.cancel(reason);
269+
270+
awaitassert.rejects(iterator.next(),(error)=>error===reason);
271+
}
272272
}
273273

274274
// Late-joining consumer should read from oldest buffered entry

β€Žtest/parallel/test-stream-iter-share-async.jsβ€Ž

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -132,6 +132,17 @@ async function testShareCancelWithReason() {
132132
);
133133
}
134134

135+
asyncfunctiontestShareCancelWithFalsyReason(){
136+
for(constreasonof[0,'',false,null]){
137+
constshared=share(from('data'));
138+
constiterator=shared.pull()[Symbol.asyncIterator]();
139+
140+
shared.cancel(reason);
141+
142+
awaitassert.rejects(iterator.next(),(error)=>error===reason);
143+
}
144+
}
145+
135146
asyncfunctiontestShareAbortSignal(){
136147
constac=newAbortController();
137148
constreason=newError('share aborted');
@@ -357,6 +368,7 @@ Promise.all([
357368
testShareCancel(),
358369
testShareCancelMidIteration(),
359370
testShareCancelWithReason(),
371+
testShareCancelWithFalsyReason(),
360372
testShareAbortSignal(),
361373
testShareAbortSignalWhileSourcePullPending(),
362374
testSharePullAbortSignalRejectsPendingNext(),

β€Žtest/parallel/test-stream-iter-share-sync.jsβ€Ž

Lines changed: 16 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -87,37 +87,32 @@ function testShareSyncCancelMidIteration() {
8787
}
8888

8989
functiontestShareSyncCancelWithReason(){
90-
// When cancel(reason) is called, a consumer that hasn't started
91-
// iterating is already detached, so it sees done:true (not the error).
92-
// But a consumer that is mid-iteration when another consumer cancels
93-
// with a reason will see the error on the next pull after cancel.
9490
constenc=newTextEncoder();
9591
function*gen(){
9692
yield[enc.encode('a')];
9793
yield[enc.encode('b')];
98-
yield[enc.encode('c')];
9994
}
10095
constshared=shareSync(gen(),{budget: 16384});
101-
constc1=shared.pull();
102-
constc2=shared.pull();
96+
constiterator1=shared.pull()[Symbol.iterator]();
97+
constiterator2=shared.pull()[Symbol.iterator]();
98+
constreason=newError('sync cancel reason');
99+
100+
iterator1.next();
101+
shared.cancel(reason);
103102

104-
// c1 reads one item, then c2 cancels with a reason
105-
constiter1=c1[Symbol.iterator]();
106-
constfirst=iter1.next();
107-
assert.strictEqual(first.done,false);
103+
assert.throws(()=>iterator1.next(),(error)=>error===reason);
104+
assert.throws(()=>iterator2.next(),(error)=>error===reason);
105+
}
108106

109-
shared.cancel(newError('sync cancel reason'));
107+
functiontestShareSyncCancelWithFalsyReason(){
108+
for(constreasonof[0,'',false,null]){
109+
constshared=shareSync(fromSync('data'));
110+
constiterator=shared.pull()[Symbol.iterator]();
110111

111-
// c1 was already iterating, it's now detached β†’ done
112-
constnext=iter1.next();
113-
assert.strictEqual(next.done,true);
112+
shared.cancel(reason);
114113

115-
// c2 never started, also detached β†’ done (not error)
116-
constbatches=[];
117-
for(constbatchofc2){
118-
batches.push(batch);
114+
assert.throws(()=>iterator.next(),(error)=>error===reason);
119115
}
120-
assert.strictEqual(batches.length,0);
121116
}
122117

123118
// =============================================================================
@@ -157,6 +152,7 @@ Promise.all([
157152
testShareSyncCancel(),
158153
testShareSyncCancelMidIteration(),
159154
testShareSyncCancelWithReason(),
155+
testShareSyncCancelWithFalsyReason(),
160156
testShareSyncSourceError(),
161157
testShareSyncStringSource(),
162158
]).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 0876a29

Browse files
trivikraduh95
authored andcommitted
stream: preserve falsy cancellation reasons
Use undefined as the no-error sentinel when cancelling broadcast and share consumers. This ensures that 0, an empty string, false, and null are propagated instead of being converted to clean completion. Make sync share surface cancellation reasons before handling detached consumers, and add regression coverage for async and sync consumers. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: codex:gpt-5.6-sol PR-URL: #64705Fixes: #64704 Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent ab50ae6 commit 0876a29

5 files changed

Lines changed: 56 additions & 44 deletions

File tree

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

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -88,7 +88,7 @@ class BroadcastImpl {
8888
#consumers =newSafeSet();
8989
#waiters =[];// Consumers with pending resolve (subset of #consumers)
9090
#ended =false;
91-
#error=null;
91+
#error;
9292
#cancelled =false;
9393
#options;
9494
#writer =null;
@@ -177,7 +177,9 @@ class BroadcastImpl {
177177
__proto__: null,
178178
next(){
179179
if(state.detached){
180-
if(self.#error)returnPromiseReject(self.#error);
180+
if(self.#error !==undefined){
181+
returnPromiseReject(self.#error);
182+
}
181183
returnkDone;
182184
}
183185

@@ -194,7 +196,7 @@ class BroadcastImpl {
194196
{__proto__: null,done: false,value: chunk});
195197
}
196198

197-
if(self.#error){
199+
if(self.#error!==undefined){
198200
state.detached=true;
199201
self.#deleteConsumer(state);
200202
returnPromiseReject(self.#error);
@@ -344,7 +346,7 @@ class BroadcastImpl {
344346
}
345347

346348
[kAbort](reason){
347-
if(this.#ended ||this.#error)return;
349+
if(this.#ended ||this.#error!==undefined)return;
348350
this.#error =reason;
349351
this.#ended =true;
350352

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

Lines changed: 14 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -73,7 +73,7 @@ class ShareImpl {
7373
#consumers =newSafeSet();
7474
#sourceIterator =null;
7575
#sourceExhausted =false;
76-
#sourceError=null;
76+
#sourceError;
7777
#cancelled =false;
7878
#pulling =false;
7979
#pullWaiters =[];
@@ -129,7 +129,7 @@ class ShareImpl {
129129
__proto__: null,
130130
[SymbolAsyncIterator](){
131131
constgetNext=async()=>{
132-
if(self.#sourceError){
132+
if(self.#sourceError!==undefined){
133133
state.detached=true;
134134
self.#consumers.delete(state);
135135
throwself.#sourceError;
@@ -141,7 +141,7 @@ class ShareImpl {
141141
// cursor must re-pull rather than terminating prematurely.
142142
for(;;){
143143
if(state.detached){
144-
if(self.#sourceError)throwself.#sourceError;
144+
if(self.#sourceError!==undefined)throwself.#sourceError;
145145
return{__proto__: null,done: true,value: undefined};
146146
}
147147

@@ -167,7 +167,7 @@ class ShareImpl {
167167
if(self.#sourceExhausted){
168168
state.detached=true;
169169
self.#deleteConsumer(state);
170-
if(self.#sourceError)throwself.#sourceError;
170+
if(self.#sourceError!==undefined)throwself.#sourceError;
171171
return{__proto__: null,done: true,value: undefined};
172172
}
173173

@@ -176,7 +176,7 @@ class ShareImpl {
176176
if(shouldBuffer===null){
177177
state.detached=true;
178178
self.#deleteConsumer(state);
179-
if(self.#sourceError)throwself.#sourceError;
179+
if(self.#sourceError!==undefined)throwself.#sourceError;
180180
return{__proto__: null,done: true,value: undefined};
181181
}
182182

@@ -260,7 +260,9 @@ class ShareImpl {
260260

261261
async #waitForBufferSpace(){
262262
while(this.#bufferedBytes >=this.#options.budget){
263-
if(this.#cancelled ||this.#sourceError ||this.#sourceExhausted){
263+
if(this.#cancelled ||
264+
this.#sourceError !==undefined||
265+
this.#sourceExhausted){
264266
returnthis.#cancelled ? null : true;
265267
}
266268

@@ -418,7 +420,7 @@ class SyncShareImpl {
418420
#consumers =newSafeSet();
419421
#sourceIterator =null;
420422
#sourceExhausted =false;
421-
#sourceError=null;
423+
#sourceError;
422424
#cancelled =false;
423425
#cachedMinCursor =0;
424426
#cachedMinCursorConsumers =0;
@@ -467,14 +469,14 @@ class SyncShareImpl {
467469
return{
468470
__proto__: null,
469471
next(){
470-
if(state.detached){
471-
return{__proto__: null,done: true,value: undefined};
472-
}
473-
if(self.#sourceError){
472+
if(self.#sourceError !==undefined){
474473
state.detached=true;
475474
self.#deleteConsumer(state);
476475
throwself.#sourceError;
477476
}
477+
if(state.detached){
478+
return{__proto__: null,done: true,value: undefined};
479+
}
478480
if(self.#cancelled){
479481
state.detached=true;
480482
self.#deleteConsumer(state);
@@ -535,7 +537,7 @@ class SyncShareImpl {
535537

536538
self.#pullFromSource();
537539

538-
if(self.#sourceError){
540+
if(self.#sourceError!==undefined){
539541
state.detached=true;
540542
self.#deleteConsumer(state);
541543
throwself.#sourceError;

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

Lines changed: 8 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -260,15 +260,15 @@ async function testWriterFailIdempotent() {
260260
},{message: 'fail!'});
261261
}
262262

263-
// cancel() with falsy reason (0, "", false) should still treat as error
264263
asyncfunctiontestCancelWithFalsyReason(){
265-
const{broadcast: bc}=broadcast();
266-
constconsumer=bc.push();
267-
constresultPromise=text(consumer).catch((err)=>err);
268-
awaitnewPromise((resolve)=>setImmediate(resolve));
269-
bc.cancel(0);
270-
constresult=awaitresultPromise;
271-
assert.strictEqual(result,0);
264+
for(constreasonof[0,'',false,null]){
265+
const{broadcast: bc}=broadcast();
266+
constiterator=bc.push()[Symbol.asyncIterator]();
267+
268+
bc.cancel(reason);
269+
270+
awaitassert.rejects(iterator.next(),(error)=>error===reason);
271+
}
272272
}
273273

274274
// Late-joining consumer should read from oldest buffered entry

β€Žtest/parallel/test-stream-iter-share-async.jsβ€Ž

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -132,6 +132,17 @@ async function testShareCancelWithReason() {
132132
);
133133
}
134134

135+
asyncfunctiontestShareCancelWithFalsyReason(){
136+
for(constreasonof[0,'',false,null]){
137+
constshared=share(from('data'));
138+
constiterator=shared.pull()[Symbol.asyncIterator]();
139+
140+
shared.cancel(reason);
141+
142+
awaitassert.rejects(iterator.next(),(error)=>error===reason);
143+
}
144+
}
145+
135146
asyncfunctiontestShareAbortSignal(){
136147
constac=newAbortController();
137148
constreason=newError('share aborted');
@@ -357,6 +368,7 @@ Promise.all([
357368
testShareCancel(),
358369
testShareCancelMidIteration(),
359370
testShareCancelWithReason(),
371+
testShareCancelWithFalsyReason(),
360372
testShareAbortSignal(),
361373
testShareAbortSignalWhileSourcePullPending(),
362374
testSharePullAbortSignalRejectsPendingNext(),

β€Žtest/parallel/test-stream-iter-share-sync.jsβ€Ž

Lines changed: 16 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -87,37 +87,32 @@ function testShareSyncCancelMidIteration() {
8787
}
8888

8989
functiontestShareSyncCancelWithReason(){
90-
// When cancel(reason) is called, a consumer that hasn't started
91-
// iterating is already detached, so it sees done:true (not the error).
92-
// But a consumer that is mid-iteration when another consumer cancels
93-
// with a reason will see the error on the next pull after cancel.
9490
constenc=newTextEncoder();
9591
function*gen(){
9692
yield[enc.encode('a')];
9793
yield[enc.encode('b')];
98-
yield[enc.encode('c')];
9994
}
10095
constshared=shareSync(gen(),{budget: 16384});
101-
constc1=shared.pull();
102-
constc2=shared.pull();
96+
constiterator1=shared.pull()[Symbol.iterator]();
97+
constiterator2=shared.pull()[Symbol.iterator]();
98+
constreason=newError('sync cancel reason');
99+
100+
iterator1.next();
101+
shared.cancel(reason);
103102

104-
// c1 reads one item, then c2 cancels with a reason
105-
constiter1=c1[Symbol.iterator]();
106-
constfirst=iter1.next();
107-
assert.strictEqual(first.done,false);
103+
assert.throws(()=>iterator1.next(),(error)=>error===reason);
104+
assert.throws(()=>iterator2.next(),(error)=>error===reason);
105+
}
108106

109-
shared.cancel(newError('sync cancel reason'));
107+
functiontestShareSyncCancelWithFalsyReason(){
108+
for(constreasonof[0,'',false,null]){
109+
constshared=shareSync(fromSync('data'));
110+
constiterator=shared.pull()[Symbol.iterator]();
110111

111-
// c1 was already iterating, it's now detached β†’ done
112-
constnext=iter1.next();
113-
assert.strictEqual(next.done,true);
112+
shared.cancel(reason);
114113

115-
// c2 never started, also detached β†’ done (not error)
116-
constbatches=[];
117-
for(constbatchofc2){
118-
batches.push(batch);
114+
assert.throws(()=>iterator.next(),(error)=>error===reason);
119115
}
120-
assert.strictEqual(batches.length,0);
121116
}
122117

123118
// =============================================================================
@@ -157,6 +152,7 @@ Promise.all([
157152
testShareSyncCancel(),
158153
testShareSyncCancelMidIteration(),
159154
testShareSyncCancelWithReason(),
155+
testShareSyncCancelWithFalsyReason(),
160156
testShareSyncSourceError(),
161157
testShareSyncStringSource(),
162158
]).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 0876a29

Browse files
trivikraduh95
authored andcommitted
stream: preserve falsy cancellation reasons
Use undefined as the no-error sentinel when cancelling broadcast and share consumers. This ensures that 0, an empty string, false, and null are propagated instead of being converted to clean completion. Make sync share surface cancellation reasons before handling detached consumers, and add regression coverage for async and sync consumers. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: codex:gpt-5.6-sol PR-URL: #64705Fixes: #64704 Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent ab50ae6 commit 0876a29

5 files changed

Lines changed: 56 additions & 44 deletions

File tree

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

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -88,7 +88,7 @@ class BroadcastImpl {
8888
#consumers =newSafeSet();
8989
#waiters =[];// Consumers with pending resolve (subset of #consumers)
9090
#ended =false;
91-
#error=null;
91+
#error;
9292
#cancelled =false;
9393
#options;
9494
#writer =null;
@@ -177,7 +177,9 @@ class BroadcastImpl {
177177
__proto__: null,
178178
next(){
179179
if(state.detached){
180-
if(self.#error)returnPromiseReject(self.#error);
180+
if(self.#error !==undefined){
181+
returnPromiseReject(self.#error);
182+
}
181183
returnkDone;
182184
}
183185

@@ -194,7 +196,7 @@ class BroadcastImpl {
194196
{__proto__: null,done: false,value: chunk});
195197
}
196198

197-
if(self.#error){
199+
if(self.#error!==undefined){
198200
state.detached=true;
199201
self.#deleteConsumer(state);
200202
returnPromiseReject(self.#error);
@@ -344,7 +346,7 @@ class BroadcastImpl {
344346
}
345347

346348
[kAbort](reason){
347-
if(this.#ended ||this.#error)return;
349+
if(this.#ended ||this.#error!==undefined)return;
348350
this.#error =reason;
349351
this.#ended =true;
350352

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

Lines changed: 14 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -73,7 +73,7 @@ class ShareImpl {
7373
#consumers =newSafeSet();
7474
#sourceIterator =null;
7575
#sourceExhausted =false;
76-
#sourceError=null;
76+
#sourceError;
7777
#cancelled =false;
7878
#pulling =false;
7979
#pullWaiters =[];
@@ -129,7 +129,7 @@ class ShareImpl {
129129
__proto__: null,
130130
[SymbolAsyncIterator](){
131131
constgetNext=async()=>{
132-
if(self.#sourceError){
132+
if(self.#sourceError!==undefined){
133133
state.detached=true;
134134
self.#consumers.delete(state);
135135
throwself.#sourceError;
@@ -141,7 +141,7 @@ class ShareImpl {
141141
// cursor must re-pull rather than terminating prematurely.
142142
for(;;){
143143
if(state.detached){
144-
if(self.#sourceError)throwself.#sourceError;
144+
if(self.#sourceError!==undefined)throwself.#sourceError;
145145
return{__proto__: null,done: true,value: undefined};
146146
}
147147

@@ -167,7 +167,7 @@ class ShareImpl {
167167
if(self.#sourceExhausted){
168168
state.detached=true;
169169
self.#deleteConsumer(state);
170-
if(self.#sourceError)throwself.#sourceError;
170+
if(self.#sourceError!==undefined)throwself.#sourceError;
171171
return{__proto__: null,done: true,value: undefined};
172172
}
173173

@@ -176,7 +176,7 @@ class ShareImpl {
176176
if(shouldBuffer===null){
177177
state.detached=true;
178178
self.#deleteConsumer(state);
179-
if(self.#sourceError)throwself.#sourceError;
179+
if(self.#sourceError!==undefined)throwself.#sourceError;
180180
return{__proto__: null,done: true,value: undefined};
181181
}
182182

@@ -260,7 +260,9 @@ class ShareImpl {
260260

261261
async #waitForBufferSpace(){
262262
while(this.#bufferedBytes >=this.#options.budget){
263-
if(this.#cancelled ||this.#sourceError ||this.#sourceExhausted){
263+
if(this.#cancelled ||
264+
this.#sourceError !==undefined||
265+
this.#sourceExhausted){
264266
returnthis.#cancelled ? null : true;
265267
}
266268

@@ -418,7 +420,7 @@ class SyncShareImpl {
418420
#consumers =newSafeSet();
419421
#sourceIterator =null;
420422
#sourceExhausted =false;
421-
#sourceError=null;
423+
#sourceError;
422424
#cancelled =false;
423425
#cachedMinCursor =0;
424426
#cachedMinCursorConsumers =0;
@@ -467,14 +469,14 @@ class SyncShareImpl {
467469
return{
468470
__proto__: null,
469471
next(){
470-
if(state.detached){
471-
return{__proto__: null,done: true,value: undefined};
472-
}
473-
if(self.#sourceError){
472+
if(self.#sourceError !==undefined){
474473
state.detached=true;
475474
self.#deleteConsumer(state);
476475
throwself.#sourceError;
477476
}
477+
if(state.detached){
478+
return{__proto__: null,done: true,value: undefined};
479+
}
478480
if(self.#cancelled){
479481
state.detached=true;
480482
self.#deleteConsumer(state);
@@ -535,7 +537,7 @@ class SyncShareImpl {
535537

536538
self.#pullFromSource();
537539

538-
if(self.#sourceError){
540+
if(self.#sourceError!==undefined){
539541
state.detached=true;
540542
self.#deleteConsumer(state);
541543
throwself.#sourceError;

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

Lines changed: 8 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -260,15 +260,15 @@ async function testWriterFailIdempotent() {
260260
},{message: 'fail!'});
261261
}
262262

263-
// cancel() with falsy reason (0, "", false) should still treat as error
264263
asyncfunctiontestCancelWithFalsyReason(){
265-
const{broadcast: bc}=broadcast();
266-
constconsumer=bc.push();
267-
constresultPromise=text(consumer).catch((err)=>err);
268-
awaitnewPromise((resolve)=>setImmediate(resolve));
269-
bc.cancel(0);
270-
constresult=awaitresultPromise;
271-
assert.strictEqual(result,0);
264+
for(constreasonof[0,'',false,null]){
265+
const{broadcast: bc}=broadcast();
266+
constiterator=bc.push()[Symbol.asyncIterator]();
267+
268+
bc.cancel(reason);
269+
270+
awaitassert.rejects(iterator.next(),(error)=>error===reason);
271+
}
272272
}
273273

274274
// Late-joining consumer should read from oldest buffered entry

β€Žtest/parallel/test-stream-iter-share-async.jsβ€Ž

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -132,6 +132,17 @@ async function testShareCancelWithReason() {
132132
);
133133
}
134134

135+
asyncfunctiontestShareCancelWithFalsyReason(){
136+
for(constreasonof[0,'',false,null]){
137+
constshared=share(from('data'));
138+
constiterator=shared.pull()[Symbol.asyncIterator]();
139+
140+
shared.cancel(reason);
141+
142+
awaitassert.rejects(iterator.next(),(error)=>error===reason);
143+
}
144+
}
145+
135146
asyncfunctiontestShareAbortSignal(){
136147
constac=newAbortController();
137148
constreason=newError('share aborted');
@@ -357,6 +368,7 @@ Promise.all([
357368
testShareCancel(),
358369
testShareCancelMidIteration(),
359370
testShareCancelWithReason(),
371+
testShareCancelWithFalsyReason(),
360372
testShareAbortSignal(),
361373
testShareAbortSignalWhileSourcePullPending(),
362374
testSharePullAbortSignalRejectsPendingNext(),

β€Žtest/parallel/test-stream-iter-share-sync.jsβ€Ž

Lines changed: 16 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -87,37 +87,32 @@ function testShareSyncCancelMidIteration() {
8787
}
8888

8989
functiontestShareSyncCancelWithReason(){
90-
// When cancel(reason) is called, a consumer that hasn't started
91-
// iterating is already detached, so it sees done:true (not the error).
92-
// But a consumer that is mid-iteration when another consumer cancels
93-
// with a reason will see the error on the next pull after cancel.
9490
constenc=newTextEncoder();
9591
function*gen(){
9692
yield[enc.encode('a')];
9793
yield[enc.encode('b')];
98-
yield[enc.encode('c')];
9994
}
10095
constshared=shareSync(gen(),{budget: 16384});
101-
constc1=shared.pull();
102-
constc2=shared.pull();
96+
constiterator1=shared.pull()[Symbol.iterator]();
97+
constiterator2=shared.pull()[Symbol.iterator]();
98+
constreason=newError('sync cancel reason');
99+
100+
iterator1.next();
101+
shared.cancel(reason);
103102

104-
// c1 reads one item, then c2 cancels with a reason
105-
constiter1=c1[Symbol.iterator]();
106-
constfirst=iter1.next();
107-
assert.strictEqual(first.done,false);
103+
assert.throws(()=>iterator1.next(),(error)=>error===reason);
104+
assert.throws(()=>iterator2.next(),(error)=>error===reason);
105+
}
108106

109-
shared.cancel(newError('sync cancel reason'));
107+
functiontestShareSyncCancelWithFalsyReason(){
108+
for(constreasonof[0,'',false,null]){
109+
constshared=shareSync(fromSync('data'));
110+
constiterator=shared.pull()[Symbol.iterator]();
110111

111-
// c1 was already iterating, it's now detached β†’ done
112-
constnext=iter1.next();
113-
assert.strictEqual(next.done,true);
112+
shared.cancel(reason);
114113

115-
// c2 never started, also detached β†’ done (not error)
116-
constbatches=[];
117-
for(constbatchofc2){
118-
batches.push(batch);
114+
assert.throws(()=>iterator.next(),(error)=>error===reason);
119115
}
120-
assert.strictEqual(batches.length,0);
121116
}
122117

123118
// =============================================================================
@@ -157,6 +152,7 @@ Promise.all([
157152
testShareSyncCancel(),
158153
testShareSyncCancelMidIteration(),
159154
testShareSyncCancelWithReason(),
155+
testShareSyncCancelWithFalsyReason(),
160156
testShareSyncSourceError(),
161157
testShareSyncStringSource(),
162158
]).then(common.mustCall());

0 commit comments

Comments
Β (0)