diff --git a/lib/core/util.js b/lib/core/util.js index f00ddfce0d1..14709525863 100644 --- a/lib/core/util.js +++ b/lib/core/util.js @@ -661,8 +661,7 @@ function ReadableStreamFrom (iterable) { return iterator.next().then(({ done, value }) => { if (done) { return queueMicrotask(() => { - controller.close() - controller.byobRequest?.respond(0) + readableStreamClose(controller) }) } else { const buf = Buffer.isBuffer(value) ? value : Buffer.from(value) diff --git a/test/fetch/readable-stream-from.js b/test/fetch/readable-stream-from.js index 7b047d31eec..efebcb2f8ad 100644 --- a/test/fetch/readable-stream-from.js +++ b/test/fetch/readable-stream-from.js @@ -27,3 +27,38 @@ test('ReadableStream empty enqueue', async (t) => { const response = new Response(iterable) t.assert.deepStrictEqual(await response.text(), '') }) + +// https://github.com/nodejs/undici/issues/5715 +test('ReadableStream cancellation while iterator.next() is in flight', async (t) => { + let resolveNext + let startNext + const nextStarted = new Promise((resolve) => { + startNext = resolve + }) + const iterable = { + [Symbol.asyncIterator] () { + return { + next () { + startNext() + return new Promise((resolve) => { + resolveNext = resolve + }) + }, + return () { + return Promise.resolve({ done: true, value: undefined }) + } + } + } + } + + const reader = new Response(iterable).body.getReader() + const read = reader.read() + await nextStarted + await reader.cancel() + resolveNext({ done: true, value: undefined }) + t.assert.deepStrictEqual(await read, { done: true, value: undefined }) + + // Let the microtask that closes the controller run. It must not throw if + // cancellation already closed the stream. + await new Promise((resolve) => setImmediate(resolve)) +})