Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
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
72 changes: 32 additions & 40 deletions lib/internal/streams/pipeline.js
Original file line number Diff line number Diff line change
Expand Up @@ -20,20 +20,16 @@ const {
ERR_INVALID_RETURN_VALUE,
ERR_MISSING_ARGS,
ERR_STREAM_DESTROYED,
ERR_STREAM_PREMATURE_CLOSE,
},
} = require('internal/errors');

const { validateCallback } = require('internal/validators');

function noop() {}

const {
isIterable,
isReadable,
isStream,
} = require('internal/streams/utils');
const assert = require('internal/assert');

let PassThrough;
let Readable;
Expand Down Expand Up @@ -109,62 +105,58 @@ async function* fromReadable(val) {

async function pump(iterable, writable, finish) {
let error;
let callback = noop;
let onresolve = null;

const resume = (err) => {
error = aggregateTwoErrors(error, err);
const _callback = callback;
callback = noop;
_callback();
};
const onClose = () => {
resume(new ERR_STREAM_PREMATURE_CLOSE());
if (err) {
error = err;
}

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

const waitForDrain = () => new Promise((resolve) => {
assert(callback === noop);
if (error || writable.destroyed) {
resolve();
const wait = () => new Promise((resolve, reject) => {
if (error) {
reject(error);
} else {
callback = resolve;
onresolve = () => {
if (error) {
reject(error);
} else {
resolve();
}
};
}
});

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

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

if (error) {
return;
await wait();
}

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

if (error) {
return;
}

writable.end();

await wait();

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

Expand Down
33 changes: 0 additions & 33 deletions test/parallel/test-stream-pipeline.js
Original file line number Diff line number Diff line change
Expand Up @@ -1387,36 +1387,3 @@ const net = require('net');
assert.strictEqual(res, content);
}));
}

{
const writableLike = new Stream();
writableLike.writableNeedDrain = true;

pipeline(
async function *() {},
writableLike,
common.mustCall((err) => {
assert.strictEqual(err.code, 'ERR_STREAM_PREMATURE_CLOSE');
})
);

writableLike.emit('close');
}

{
const writableLike = new Stream();
writableLike.write = () => false;

pipeline(
async function *() {
yield null;
yield null;
},
writableLike,
common.mustCall((err) => {
assert.strictEqual(err.code, 'ERR_STREAM_PREMATURE_CLOSE');
})
);

writableLike.emit('close');
}
Copy link
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

These test are just for very edge case mostly made up compat scenarios and actually do not align with compat as implemented in finished. I just removed them,