Skip to content

Commit ca45ad4

Browse files
committed
stream: read sync sources synchronously in stream/iter pipeTo()
pipeTo() iterated the from() normalization of a sync iterable source with for await...of, which costs three promises and several ticks per batch, even though the batches are read synchronously. Let the iterator returned by from() for a sync iterable read the next batch synchronously when no operation is running or queued and the value read needs no asynchronous normalization, and have pipeTo() use it when there are no transforms and no signal. Other values are normalized through next() as before, a write error still closes the source, and an error reading the source still does not. The source and the writer see the same calls in the same order; the batches are no longer written on separate ticks unless a write is asynchronous. Piping a sync generator yielding 16-byte chunks, one per batch, is about 2.5 times faster, the same as pipeToSync(). Assisted-by: OpenCode
1 parent 3c7362f commit ca45ad4

4 files changed

Lines changed: 137 additions & 1 deletion

File tree

‎lib/internal/streams/iter/from.js‎

Lines changed: 55 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -139,9 +139,36 @@ function waitForNormalization(value, context) {
139139
return promise;
140140
}
141141

142+
// The method of an iterator returned by from() for a sync iterable that reads
143+
// the next batch synchronously when it can, for pipeTo(): see nextSyncBatch()
144+
// in createSyncSourceNormalizer().
145+
const kNextSyncBatch = Symbol('kNextSyncBatch');
146+
142147
function createNormalizationIterator(createIterator) {
143148
const context = createNormalizationContext();
144149
const iterator = createIterator(context);
150+
if (iterator[kNextSyncBatch] !== undefined) {
151+
return ObjectSetPrototypeOf({
152+
next(value) {
153+
return FunctionPrototypeCall(iterator.next, iterator, value);
154+
},
155+
return(value) {
156+
cancelNormalization(
157+
context, lazyDOMException('Aborted', 'AbortError'));
158+
return FunctionPrototypeCall(iterator.return, iterator, value);
159+
},
160+
throw(error) {
161+
cancelNormalization(context, error, true);
162+
return FunctionPrototypeCall(iterator.throw, iterator, error);
163+
},
164+
[kNextSyncBatch]() {
165+
return iterator[kNextSyncBatch]();
166+
},
167+
[SymbolAsyncIterator]() {
168+
return this;
169+
},
170+
}, null);
171+
}
145172
return ObjectSetPrototypeOf({
146173
next(value) {
147174
return FunctionPrototypeCall(iterator.next, iterator, value);
@@ -938,7 +965,7 @@ function createSyncSourceNormalizer(source, context) {
938965
let done = false;
939966
// A normalizeAsyncSourceValue() generator for the value being normalized.
940967
let valueBatches = null;
941-
const { run, settled, release } = createOperationQueue();
968+
const { run, settled, release, idle } = createOperationQueue();
942969

943970
// Results are often produced synchronously. An async generator stays busy
944971
// until the tick after a yield or return (both await their operand), so
@@ -1063,10 +1090,36 @@ function createSyncSourceNormalizer(source, context) {
10631090
});
10641091
}
10651092

1093+
// Read the next batch like next() but synchronously, when it is read
1094+
// synchronously from the source and no operation is running or queued:
1095+
// returns the batch, or null when done. Returns undefined, having started
1096+
// nothing but the normalization of a value, when next() must be used.
1097+
// Throws on error, like next() rejects.
1098+
function nextSyncBatch() {
1099+
if (!idle() || valueBatches !== null) return undefined;
1100+
if (done) return null;
1101+
let result;
1102+
try {
1103+
result = reader.next();
1104+
} catch (error) {
1105+
done = true;
1106+
throw error;
1107+
}
1108+
if (result.done) {
1109+
done = true;
1110+
return null;
1111+
}
1112+
const value = result.value;
1113+
if (ArrayIsArray(value)) return value;
1114+
valueBatches = normalizeAsyncSourceValue(value.value, context, false);
1115+
return undefined;
1116+
}
1117+
10661118
return ObjectSetPrototypeOf({
10671119
next() { return run(doNext); },
10681120
return(value) { return run(doReturn, value); },
10691121
throw(error) { return run(doThrow, error); },
1122+
[kNextSyncBatch]: nextSyncBatch,
10701123
}, null);
10711124
}
10721125

@@ -1290,6 +1343,7 @@ module.exports = {
12901343
isPrimitiveChunk,
12911344
isSyncIterable,
12921345
isUint8ArrayBatch,
1346+
kNextSyncBatch,
12931347
normalizeAsyncSource,
12941348
normalizeAsyncValue,
12951349
normalizeSyncSource,

‎lib/internal/streams/iter/pull.js‎

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,7 @@ const {
5252
isSyncIterable,
5353
isAsyncIterable,
5454
isUint8ArrayBatch,
55+
kNextSyncBatch,
5556
} = require('internal/streams/iter/from');
5657

5758
const {
@@ -1708,6 +1709,34 @@ async function pipeTo(source, ...args) {
17081709
}
17091710
}
17101711

1712+
// Write the batches of a sync source as the for await...of loop below
1713+
// does, reading them synchronously when possible: an error reading a
1714+
// batch ends the loop, and an error writing one closes the source first,
1715+
// ignoring errors closing it.
1716+
async function pipeSyncSource(iterator) {
1717+
for (;;) {
1718+
let batch = iterator[kNextSyncBatch]();
1719+
if (batch === undefined) {
1720+
const result = await iterator.next();
1721+
if (result.done) return;
1722+
batch = result.value;
1723+
} else if (batch === null) {
1724+
return;
1725+
}
1726+
try {
1727+
const p = writeBatch(batch);
1728+
if (p) await p;
1729+
} catch (error) {
1730+
try {
1731+
await iterator.return();
1732+
} catch {
1733+
// The error writing the batch is thrown.
1734+
}
1735+
throw error;
1736+
}
1737+
}
1738+
}
1739+
17111740
function writeBatch(batch) {
17121741
// Single chunk, the common case: check the view around writeSync()
17131742
// without allocating a batch entry, and create one only to fall back to
@@ -1767,6 +1796,8 @@ async function pipeTo(source, ...args) {
17671796
const p = writeBatch(batch);
17681797
if (p) await p;
17691798
}
1799+
} else if (normalized[kNextSyncBatch] !== undefined) {
1800+
await pipeSyncSource(normalized);
17701801
} else {
17711802
for await (const batch of normalized) {
17721803
const p = writeBatch(batch);

‎lib/internal/streams/iter/utils.js‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -564,6 +564,10 @@ function createOperationQueue() {
564564
busy = false;
565565
drain();
566566
},
567+
// Whether no operation is running or queued.
568+
idle() {
569+
return !busy && (queue === null || queue.length === 0);
570+
},
567571
};
568572
}
569573

‎test/parallel/test-stream-iter-pipeto.js‎

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -345,6 +345,52 @@ async function testPipeToSyncIterableAsyncValue() {
345345
assert.strictEqual(result, 'ab');
346346
}
347347

348+
// pipeTo() reads sync iterables synchronously when it can. An error writing
349+
// a batch must still close the source, ignoring an error closing it, and an
350+
// error reading the source must not close it.
351+
async function testPipeToSyncIterableWriteError() {
352+
for (const returnThrows of [false, true]) {
353+
const error = new Error('write');
354+
let closed = false;
355+
const source = {
356+
[Symbol.iterator]() {
357+
let i = 0;
358+
return {
359+
next() {
360+
return { done: false, value: [new Uint8Array([i++])] };
361+
},
362+
return() {
363+
closed = true;
364+
if (returnThrows) throw new Error('return');
365+
return { done: true };
366+
},
367+
};
368+
},
369+
};
370+
await assert.rejects(pipeTo(source, {
371+
write: common.mustNotCall(),
372+
writeSync: common.mustCall(() => { throw error; }),
373+
fail: common.mustCall((reason) => assert.strictEqual(reason, error)),
374+
}), error);
375+
assert.strictEqual(closed, true);
376+
}
377+
378+
const error = new Error('source');
379+
const source = {
380+
[Symbol.iterator]() {
381+
return {
382+
next() { throw error; },
383+
return: common.mustNotCall(),
384+
};
385+
},
386+
};
387+
await assert.rejects(pipeTo(source, {
388+
write: common.mustNotCall(),
389+
writeSync: common.mustNotCall(),
390+
fail: common.mustCall((reason) => assert.strictEqual(reason, error)),
391+
}), error);
392+
}
393+
348394
Promise.all([
349395
testPipeToSync(),
350396
testPipeTo(),
@@ -365,5 +411,6 @@ Promise.all([
365411
testPipeToSyncIterableUsesFromBatching(),
366412
testPipeToSyncIterableWriteFallback(),
367413
testPipeToSyncIterableAsyncValue(),
414+
testPipeToSyncIterableWriteError(),
368415
testPipeToSourceNormalizationIndependentOfWriter(),
369416
]).then(common.mustCall());

0 commit comments

Comments
 (0)