Commit 3169ab0ccfa for nodejs
commit 3169ab0ccfa734677ac87aaa74ccd65123406df7
Author: James M Snell <jasnell@gmail.com>
Date: Sat Oct 3 14:47:19 2026 +0000
stream: guard shareSync() against re-entrant source reads
If a shareSync() source read from a consumer of the same share while
producing a value, the nested read re-entered the source iterator.
For generators this threw "Generator is already running" from inside
the nested read, which recorded that error as the share's source
error while the outer read was still in progress, leaving every
consumer in an error state.
Fail the nested read with ERR_INVALID_STATE before touching the
share's state. The source sees the error and may handle it; if it lets
it escape, it becomes the source error as with any other source
failure.
Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
PR-URL: https://github.com/nodejs/node/pull/66483
Reviewed-By: Trivikram Kamat <trivikr.dev@gmail.com>
diff --git a/lib/internal/streams/iter/share.js b/lib/internal/streams/iter/share.js
index ac287d9472a..17dabf5da03 100644
--- a/lib/internal/streams/iter/share.js
+++ b/lib/internal/streams/iter/share.js
@@ -60,6 +60,7 @@ const {
ERR_INVALID_ARG_TYPE,
ERR_INVALID_ARG_VALUE,
ERR_INVALID_RETURN_VALUE,
+ ERR_INVALID_STATE,
ERR_OUT_OF_RANGE,
},
} = require('internal/errors');
@@ -548,6 +549,7 @@ class SyncShareImpl {
#sourceExhausted = false;
#sourceError = kNoShareError;
#cancelled = false;
+ #pulling = false;
#cachedMinCursor = 0;
#cachedMinCursorConsumers = 0;
/** Cumulative byte size of buffered entries */
@@ -747,6 +749,15 @@ class SyncShareImpl {
#pullFromSource() {
if (this.#sourceExhausted || this.#cancelled) return;
+ // A source that reads from its own share would re-enter its own
+ // iterator. Fail that read without touching the share's state; the
+ // source sees the error and may handle it.
+ if (this.#pulling) {
+ throw new ERR_INVALID_STATE(
+ 'shareSync() source cannot be read while it is producing a value');
+ }
+
+ this.#pulling = true;
try {
this.#sourceIterator ||= this.#source[SymbolIterator]();
@@ -760,6 +771,8 @@ class SyncShareImpl {
} catch (error) {
this.#sourceError = error;
this.#sourceExhausted = true;
+ } finally {
+ this.#pulling = false;
}
}
diff --git a/test/parallel/test-stream-iter-share-sync.js b/test/parallel/test-stream-iter-share-sync.js
index dd392d9c992..a843e7311dd 100644
--- a/test/parallel/test-stream-iter-share-sync.js
+++ b/test/parallel/test-stream-iter-share-sync.js
@@ -264,6 +264,50 @@ function testShareSyncDropOldestSplitsOversizedBatches() {
Array.from({ length: slowSeen.length }, (_, i) => 50 - slowSeen.length + i));
}
+function testShareSyncReentrantSourceRead() {
+ const enc = new TextEncoder();
+ let sibling;
+ let reentrantError;
+ function* source() {
+ yield [enc.encode('a')];
+ try {
+ sibling.next();
+ } catch (err) {
+ reentrantError = err;
+ }
+ yield [enc.encode('b')];
+ }
+ const shared = shareSync(source());
+ const c1 = shared.pull()[Symbol.iterator]();
+ sibling = shared.pull()[Symbol.iterator]();
+
+ assert.deepStrictEqual(c1.next().value, [enc.encode('a')]);
+ assert.deepStrictEqual(sibling.next().value, [enc.encode('a')]);
+ // Pulling 'b' runs the source, which tries to read its own share.
+ assert.deepStrictEqual(c1.next().value, [enc.encode('b')]);
+ assert.strictEqual(reentrantError?.code, 'ERR_INVALID_STATE');
+ // The failed re-entrant read left the share intact.
+ assert.deepStrictEqual(sibling.next().value, [enc.encode('b')]);
+ assert.strictEqual(c1.next().done, true);
+ assert.strictEqual(sibling.next().done, true);
+}
+
+function testShareSyncReentrantSourceReadUncaught() {
+ let sibling;
+ function* source() {
+ yield [new Uint8Array(1)];
+ sibling.next();
+ }
+ const shared = shareSync(source());
+ const c1 = shared.pull()[Symbol.iterator]();
+ sibling = shared.pull()[Symbol.iterator]();
+ c1.next();
+ sibling.next();
+ // The error escapes the source, so it becomes the share's source error.
+ assert.throws(() => c1.next(), { code: 'ERR_INVALID_STATE' });
+ assert.throws(() => sibling.next(), { code: 'ERR_INVALID_STATE' });
+}
+
function testShareSyncStringSource() {
const shared = shareSync('hello-sync-share');
const result = textSync(shared.pull());
@@ -281,6 +325,8 @@ Promise.all([
testShareSyncRejectsUnbounded(),
testShareSyncRejectsDropNewest(),
testShareSyncDropOldestSplitsOversizedBatches(),
+ testShareSyncReentrantSourceRead(),
+ testShareSyncReentrantSourceReadUncaught(),
testShareSyncStringSource(),
testShareSyncRetainsBufferWhenAllConsumersDetach(),
testShareSyncStrictBackpressureDetaches(),