FazBrowse GitHub Viewer | Trending |
URL:
| Home
Tools: [Download Repo ZIP]   [Original HTTPS Page]

stream: preserve falsy cancellation reasons · nodejs/node@0876a29 · GitHub

/ node Public

Commit 0876a29

Browse files
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: #64705 Fixes: #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 = new SafeSet();
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) return PromiseReject(self.#error);
180+
if (self.#error !== undefined) {
181+
return PromiseReject(self.#error);
182+
}
181183
return kDone;
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
return PromiseReject(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 = new SafeSet();
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
const getNext = async () => {
132-
if (self.#sourceError) {
132+
if (self.#sourceError !== undefined) {
133133
state.detached = true;
134134
self.#consumers.delete(state);
135135
throw self.#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) throw self.#sourceError;
144+
if (self.#sourceError !== undefined) throw self.#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) throw self.#sourceError;
170+
if (self.#sourceError !== undefined) throw self.#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) throw self.#sourceError;
179+
if (self.#sourceError !== undefined) throw self.#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
return this.#cancelled ? null : true;
265267
}
266268

@@ -418,7 +420,7 @@ class SyncShareImpl {
418420
#consumers = new SafeSet();
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
throw self.#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
throw self.#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
async function testCancelWithFalsyReason() {
265-
const { broadcast: bc } = broadcast();
266-
const consumer = bc.push();
267-
const resultPromise = text(consumer).catch((err) => err);
268-
await new Promise((resolve) => setImmediate(resolve));
269-
bc.cancel(0);
270-
const result = await resultPromise;
271-
assert.strictEqual(result, 0);
264+
for (const reason of [0, '', false, null]) {
265+
const { broadcast: bc } = broadcast();
266+
const iterator = bc.push()[Symbol.asyncIterator]();
267+
268+
bc.cancel(reason);
269+
270+
await assert.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+
async function testShareCancelWithFalsyReason() {
136+
for (const reason of [0, '', false, null]) {
137+
const shared = share(from('data'));
138+
const iterator = shared.pull()[Symbol.asyncIterator]();
139+
140+
shared.cancel(reason);
141+
142+
await assert.rejects(iterator.next(), (error) => error === reason);
143+
}
144+
}
145+
135146
async function testShareAbortSignal() {
136147
const ac = new AbortController();
137148
const reason = new Error('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
function testShareSyncCancelWithReason() {
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
const enc = new TextEncoder();
9591
function* gen() {
9692
yield [enc.encode('a')];
9793
yield [enc.encode('b')];
98-
yield [enc.encode('c')];
9994
}
10095
const shared = shareSync(gen(), { budget: 16384 });
101-
const c1 = shared.pull();
102-
const c2 = shared.pull();
96+
const iterator1 = shared.pull()[Symbol.iterator]();
97+
const iterator2 = shared.pull()[Symbol.iterator]();
98+
const reason = new Error('sync cancel reason');
99+
100+
iterator1.next();
101+
shared.cancel(reason);
103102

104-
// c1 reads one item, then c2 cancels with a reason
105-
const iter1 = c1[Symbol.iterator]();
106-
const first = 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(new Error('sync cancel reason'));
107+
function testShareSyncCancelWithFalsyReason() {
108+
for (const reason of [0, '', false, null]) {
109+
const shared = shareSync(fromSync('data'));
110+
const iterator = shared.pull()[Symbol.iterator]();
110111

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

115-
// c2 never started, also detached → done (not error)
116-
const batches = [];
117-
for (const batch of c2) {
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)

Back | FazBrowse Home | New Git URL