Skip to content
Open
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
9 changes: 9 additions & 0 deletions lib/internal/streams/duplexify.js
Original file line number Diff line number Diff line change
Expand Up @@ -103,6 +103,8 @@ module.exports = function duplexify(body, name) {
const then = value?.then;
if (typeof then === 'function') {
let d;
// Tracks whether writable final() has started.
let finalized = false;

const promise = FunctionPrototypeCall(
then,
Expand All @@ -111,6 +113,12 @@ module.exports = function duplexify(body, name) {
if (val != null) {
throw new ERR_INVALID_RETURN_VALUE('nully', 'body', val);
}
// The async function returned without (fully) consuming the input.
// Destroy the duplex so that pipeline propagates destruction
// upstream. See https://github.com/nodejs/node/issues/55077.
if (!finalized) {
destroyer(d);
}
},
(err) => {
destroyer(d, err);
Expand All @@ -123,6 +131,7 @@ module.exports = function duplexify(body, name) {
readable: false,
write,
final(cb) {
finalized = true;
final(async () => {
try {
await promise;
Expand Down
16 changes: 16 additions & 0 deletions test/parallel/test-stream-duplex-from.js
Original file line number Diff line number Diff line change
Expand Up @@ -418,3 +418,19 @@ function makeATestWritableStream(writeFunc) {
}));
r.destroy(expectedErr);
}

// Regression for https://github.com/nodejs/node/issues/55077:
// An AsyncFunction passed to Duplex.from() that returns without consuming its
// input must still allow pipeline() to complete and destroy the upstream.
{
const r = Readable.from(['foo', 'bar', 'baz']);
pipeline(
r,
Duplex.from(async function() {
// Intentionally do not consume the async iterable input.
}),
common.mustCall(() => {
assert.strictEqual(r.destroyed, true);
}),
);
}
Loading