Commit a39c0fae57e for nodejs
commit a39c0fae57e34ddaed15f67ac16065910e451147
Author: James M Snell <jasnell@gmail.com>
Date: Sat Oct 3 14:45:55 2026 +0000
stream: split oversized share batches under drop-oldest
share() and shareSync() buffer each batch pulled from the source as a
single entry, and the 'drop-oldest' policy evicts whole entries until
the buffer is below the budget. from() and fromSync() combine up to 128
values of a sync source into one batch, so a single entry can be many
times larger than the budget. Evicting it discarded every chunk in it,
including chunks that slower consumers had not read yet: a consumer
could lose the entire stream while a faster consumer read all of it.
When the policy is 'drop-oldest', split batches that are larger than
the budget into consecutive entries that are each smaller than it, so
eviction keeps the newest chunks that fit within the budget.
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/doc/api/stream_iter.md b/doc/api/stream_iter.md
index 7124d5e8fcf..eeff3ed0d41 100644
--- a/doc/api/stream_iter.md
+++ b/doc/api/stream_iter.md
@@ -1504,7 +1504,9 @@ budget. With `'drop-newest'`, the entry pulled from the source is discarded
and the consumer then waits in the same way, so in both cases a stalled
consumer also stalls the consumers that are ahead of it. Only `'drop-oldest'`
lets consumers that are ahead continue, by discarding the oldest buffered
-entries that the slowest consumer has not read yet.
+entries that the slowest consumer has not read yet. A batch pulled from the
+source that is larger than `budget` is split into smaller entries first, so
+eviction keeps the newest chunks that fit within the budget.
```mjs
import { from, share, text } from 'node:stream/iter';
diff --git a/lib/internal/streams/iter/share.js b/lib/internal/streams/iter/share.js
index 5439dfeb3d6..ac287d9472a 100644
--- a/lib/internal/streams/iter/share.js
+++ b/lib/internal/streams/iter/share.js
@@ -44,6 +44,7 @@ const {
getMinCursor,
onSignalAbort,
parsePullArgs,
+ splitBatchEntry,
validateBatchEntry,
} = require('internal/streams/iter/utils');
const {
@@ -427,9 +428,7 @@ class ShareImpl {
if (result.done) {
this.#sourceExhausted = true;
} else if (!discard) {
- const entry = createBatchEntry(result.value);
- this.#buffer.push(entry);
- this.#bufferedBytes += entry.byteLength;
+ this.#bufferBatch(result.value);
}
} catch (error) {
this.#sourceError = error;
@@ -447,6 +446,25 @@ class ShareImpl {
})();
}
+ #bufferBatch(batch) {
+ const entry = createBatchEntry(batch);
+ // 'drop-oldest' evicts whole entries. A single pulled batch can be much
+ // larger than the budget (for example when from() combines many values
+ // of a sync source), and evicting it would discard every chunk in it,
+ // including ones that no consumer has read. Split such batches so that
+ // eviction keeps the newest chunks that fit the budget.
+ const entries = this.#options.backpressure === 'drop-oldest' ?
+ splitBatchEntry(entry, this.#options.budget) : undefined;
+ if (entries === undefined) {
+ this.#buffer.push(entry);
+ } else {
+ for (let i = 0; i < entries.length; i++) {
+ this.#buffer.push(entries[i]);
+ }
+ }
+ this.#bufferedBytes += entry.byteLength;
+ }
+
#tryTrimBuffer() {
// Retain buffered data for consumers that attach while none are active.
// Without this, the last consumer detaching would discard entries that
@@ -737,9 +755,7 @@ class SyncShareImpl {
if (result.done) {
this.#sourceExhausted = true;
} else {
- const entry = createBatchEntry(result.value);
- this.#buffer.push(entry);
- this.#bufferedBytes += entry.byteLength;
+ this.#bufferBatch(result.value);
}
} catch (error) {
this.#sourceError = error;
@@ -747,6 +763,25 @@ class SyncShareImpl {
}
}
+ #bufferBatch(batch) {
+ const entry = createBatchEntry(batch);
+ // 'drop-oldest' evicts whole entries. A single pulled batch can be much
+ // larger than the budget (for example when from() combines many values
+ // of a sync source), and evicting it would discard every chunk in it,
+ // including ones that no consumer has read. Split such batches so that
+ // eviction keeps the newest chunks that fit the budget.
+ const entries = this.#options.backpressure === 'drop-oldest' ?
+ splitBatchEntry(entry, this.#options.budget) : undefined;
+ if (entries === undefined) {
+ this.#buffer.push(entry);
+ } else {
+ for (let i = 0; i < entries.length; i++) {
+ this.#buffer.push(entries[i]);
+ }
+ }
+ this.#bufferedBytes += entry.byteLength;
+ }
+
#tryTrimBuffer() {
// Retain buffered data for consumers that attach while none are active.
// Without this, the last consumer detaching would discard entries that
diff --git a/lib/internal/streams/iter/utils.js b/lib/internal/streams/iter/utils.js
index a6069c9ed3b..34698dbdf0d 100644
--- a/lib/internal/streams/iter/utils.js
+++ b/lib/internal/streams/iter/utils.js
@@ -4,6 +4,7 @@ const {
Array,
ArrayBufferPrototypeGetByteLength,
ArrayBufferPrototypeGetDetached,
+ ArrayPrototypePush,
ArrayPrototypeSlice,
PromiseResolve,
PromiseWithResolvers,
@@ -245,6 +246,35 @@ function createBatchEntry(chunks) {
return { __proto__: null, views, byteLength };
}
+/**
+ * Split a batch entry into consecutive entries that are each smaller than
+ * `limit` bytes, preserving chunk order. A single chunk of `limit` bytes or
+ * more forms an entry on its own. Returns `undefined` if no split is needed.
+ * @param {{ views: object[], byteLength: number }} entry
+ * @param {number} limit
+ * @returns {Array<{ views: object[], byteLength: number }>|undefined}
+ */
+function splitBatchEntry(entry, limit) {
+ const { views } = entry;
+ if (entry.byteLength < limit || views.length < 2) return undefined;
+ const entries = [];
+ let current = [];
+ let byteLength = 0;
+ for (let i = 0; i < views.length; i++) {
+ const view = views[i];
+ if (current.length > 0 && byteLength + view.byteLength >= limit) {
+ ArrayPrototypePush(entries,
+ { __proto__: null, views: current, byteLength });
+ current = [];
+ byteLength = 0;
+ }
+ ArrayPrototypePush(current, view);
+ byteLength += view.byteLength;
+ }
+ ArrayPrototypePush(entries, { __proto__: null, views: current, byteLength });
+ return entries;
+}
+
function validateBatchEntry(entry) {
const chunks = new Array(entry.views.length);
for (let i = 0; i < entry.views.length; i++) {
@@ -456,6 +486,7 @@ module.exports = {
concatBytes,
convertChunks,
createBatchEntry,
+ splitBatchEntry,
getProtocolMethod,
getWriterSignal,
getMinCursor,
diff --git a/test/parallel/test-stream-iter-share-async.js b/test/parallel/test-stream-iter-share-async.js
index 29de8e8a141..8ea06427543 100644
--- a/test/parallel/test-stream-iter-share-async.js
+++ b/test/parallel/test-stream-iter-share-async.js
@@ -400,6 +400,43 @@ async function testShareRetainsBufferWhenAllConsumersDetach() {
assert.strictEqual(await text(shared.pull()), 'abc');
}
+async function testShareDropOldestSplitsOversizedBatches() {
+ // from() combines the values of this sync generator into a single batch
+ // that is much larger than the budget. Evicting that batch as a whole would
+ // leave the slower consumer with nothing at all.
+ function* source() {
+ for (let i = 0; i < 50; i++) {
+ const chunk = new Uint8Array(4096);
+ chunk[0] = i;
+ yield chunk;
+ }
+ }
+ const shared = share(source(), {
+ budget: 65536,
+ backpressure: 'drop-oldest',
+ });
+ const fast = shared.pull()[Symbol.asyncIterator]();
+ const slow = shared.pull()[Symbol.asyncIterator]();
+
+ const fastSeen = [];
+ for (let r = await fast.next(); !r.done; r = await fast.next()) {
+ for (const chunk of r.value) fastSeen.push(chunk[0]);
+ }
+ assert.deepStrictEqual(fastSeen, Array.from({ length: 50 }, (_, i) => i));
+
+ const slowSeen = [];
+ for (let r = await slow.next(); !r.done; r = await slow.next()) {
+ for (const chunk of r.value) slowSeen.push(chunk[0]);
+ }
+ // The slow consumer lost the oldest chunks but keeps an in-order suffix
+ // that fits the budget.
+ assert.ok(slowSeen.length > 0);
+ assert.ok(slowSeen.length * 4096 < 65536);
+ assert.deepStrictEqual(
+ slowSeen,
+ Array.from({ length: slowSeen.length }, (_, i) => 50 - slowSeen.length + i));
+}
+
async function testShareConsumerBreak() {
// Verify that a consumer breaking mid-iteration detaches properly
const enc = new TextEncoder();
@@ -495,6 +532,7 @@ Promise.all([
testShareSourceErrorFollowsBufferedData(),
testShareLateJoiningConsumer(),
testShareRetainsBufferWhenAllConsumersDetach(),
+ testShareDropOldestSplitsOversizedBatches(),
testShareConsumerBreak(),
testShareMultipleConsumersConcurrentPull(),
testShareConsumerConcurrentNextCalls(),
diff --git a/test/parallel/test-stream-iter-share-sync.js b/test/parallel/test-stream-iter-share-sync.js
index 2721905bc93..dd392d9c992 100644
--- a/test/parallel/test-stream-iter-share-sync.js
+++ b/test/parallel/test-stream-iter-share-sync.js
@@ -227,6 +227,43 @@ function testShareSyncStrictForOfDoesNotWedgeOthers() {
assert.strictEqual(count, 20);
}
+function testShareSyncDropOldestSplitsOversizedBatches() {
+ // fromSync() combines the values of this generator into a single batch
+ // that is much larger than the budget. Evicting that batch as a whole would
+ // leave the slower consumer with nothing at all.
+ function* source() {
+ for (let i = 0; i < 50; i++) {
+ const chunk = new Uint8Array(4096);
+ chunk[0] = i;
+ yield chunk;
+ }
+ }
+ const shared = shareSync(source(), {
+ budget: 65536,
+ backpressure: 'drop-oldest',
+ });
+ const fast = shared.pull()[Symbol.iterator]();
+ const slow = shared.pull()[Symbol.iterator]();
+
+ const fastSeen = [];
+ for (let r = fast.next(); !r.done; r = fast.next()) {
+ for (const chunk of r.value) fastSeen.push(chunk[0]);
+ }
+ assert.deepStrictEqual(fastSeen, Array.from({ length: 50 }, (_, i) => i));
+
+ const slowSeen = [];
+ for (let r = slow.next(); !r.done; r = slow.next()) {
+ for (const chunk of r.value) slowSeen.push(chunk[0]);
+ }
+ // The slow consumer lost the oldest chunks but keeps an in-order suffix
+ // that fits the budget.
+ assert.ok(slowSeen.length > 0);
+ assert.ok(slowSeen.length * 4096 < 65536);
+ assert.deepStrictEqual(
+ slowSeen,
+ Array.from({ length: slowSeen.length }, (_, i) => 50 - slowSeen.length + i));
+}
+
function testShareSyncStringSource() {
const shared = shareSync('hello-sync-share');
const result = textSync(shared.pull());
@@ -243,6 +280,7 @@ Promise.all([
testShareSyncSourceError(),
testShareSyncRejectsUnbounded(),
testShareSyncRejectsDropNewest(),
+ testShareSyncDropOldestSplitsOversizedBatches(),
testShareSyncStringSource(),
testShareSyncRetainsBufferWhenAllConsumersDetach(),
testShareSyncStrictBackpressureDetaches(),