Skip to content
Merged
Changes from 1 commit
Commits
Show all changes
26 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Prev Previous commit
Next Next commit
fixup! separate out pump function
  • Loading branch information
debadree25 committed Jan 26, 2023
commit 2ec9a448a565900799368aeb3706960f26e01a91
74 changes: 69 additions & 5 deletions lib/internal/streams/pipeline.js
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ const {
isWritableStream,
isTransformStream,
isWebStream,
isReadableStream,
} = require('internal/streams/utils');
const { AbortController } = require('internal/abort_controller');

Expand Down Expand Up @@ -98,6 +99,65 @@ async function pump(iterable, writable, finish, { end }) {
let error;
let onresolve = null;

const resume = (err) => {
if (err) {
error = err;
}

if (onresolve) {
const callback = onresolve;
onresolve = null;
callback();
}
};

const wait = () => new Promise((resolve, reject) => {
if (error) {
reject(error);
} else {
onresolve = () => {
if (error) {
reject(error);
} else {
resolve();
}
};
}
});

writable.on('drain', resume);
const cleanup = eos(writable, { readable: false }, resume);

try {
if (writable.writableNeedDrain) {
await wait();
}

for await (const chunk of iterable) {
if (!writable.write(chunk)) {
await wait();
}
}

if (end) {
writable.end();
}

await wait();

finish();
} catch (err) {
finish(error !== err ? aggregateTwoErrors(error, err) : err);
} finally {
cleanup();
writable.off('drain', resume);
}
}

async function pumpWeb(iterable, writable, finish, { end }) {
let error;
let onresolve = null;

if (isTransformStream(writable)) {
writable = writable.writable;
}
Expand Down Expand Up @@ -353,12 +413,13 @@ function pipelineImpl(streams, callback, opts) {
if (isReadable(stream) && isLastStream) {
lastStreamCleanup.push(cleanup);
}
} else if (isTransformStream(ret) || isReadableStream(ret)) {
const toRead = ret.readable || ret;
finishCount++;
pumpWeb(toRead, stream, finish, { end });
} else if (isIterable(ret)) {
finishCount++;
pump(ret, stream, finish, { end });
} else if (isTransformStream(ret)) {
finishCount++;
pump(ret.readable, stream, finish, { end });
} else {
throw new ERR_INVALID_ARG_TYPE(
'val', ['Readable', 'Iterable', 'AsyncIterable'], ret);
Expand All @@ -368,11 +429,14 @@ function pipelineImpl(streams, callback, opts) {
if (isReadableNodeStream(ret)) {
finishCount += 2;
pipeNodeToWeb(ret, stream, finish, { end });
} else if (isReadableStream(ret)) {
finishCount++;
pumpWeb(ret, stream, finish, { end });
} else if (isTransformStream(ret)) {
pumpWeb(ret.readable, stream, finish, { end });
} else if (isIterable(ret)) {
finishCount++;
pump(ret, stream, finish, { end });
} else if (isTransformStream(ret)) {
pump(ret.readable, stream, finish, { end });
} else {
throw new ERR_INVALID_ARG_TYPE(
'val', ['Readable', 'Iterable', 'AsyncIterable'], ret);
Expand Down