From d189ea13ea02146fadaa9aed0e7d6f4bffdbd7e3 Mon Sep 17 00:00:00 2001 From: Dinh Le Date: Sat, 26 Sep 2026 14:17:32 +0700 Subject: [PATCH] fix(fetch): release the event iterator when an event fails to serialize When a yielded value could not be serialized (BigInt, circular object, throwing toJSON), toEventStream errored the stream without calling iterator.return(). An errored ReadableStream never invokes cancel(), so the source generator stayed suspended at its yield and its finally block never ran, leaking whatever it held. The node, fastify and aws-lambda adapters inherited this through the fetch implementation. Serialization failures are now handled separately from iterator errors: the suspended iterator is released first, then the stream is errored with the original serialization error. --- packages/fetch/src/event-stream.test.ts | 49 ++++++++++++++++++ packages/fetch/src/event-stream.ts | 66 ++++++++++++++++--------- 2 files changed, 93 insertions(+), 22 deletions(-) diff --git a/packages/fetch/src/event-stream.test.ts b/packages/fetch/src/event-stream.test.ts index 68b977a..a1e5785 100644 --- a/packages/fetch/src/event-stream.test.ts +++ b/packages/fetch/src/event-stream.test.ts @@ -293,6 +293,55 @@ describe('toEventStream', () => { expect((await reader.read()).done).toEqual(true) }) + it.each([ + ['a BigInt', () => ({ big: 1n }), 'BigInt'], + ['a circular reference', () => { + const value: Record = {} + value.self = value + return value + }, 'circular'], + ['a throwing toJSON', () => ({ toJSON() { throw new Error('toJSON failed') } }), 'toJSON failed'], + ])('releases the iterator when an event has %s', async (_, value, message) => { + let hasFinally = false + + async function* gen() { + try { + yield value() + yield 2 + } + finally { + hasFinally = true + } + } + + const reader = toEventStream(gen()) + .pipeThrough(new TextDecoderStream()) + .getReader() + + expect((await reader.read())).toEqual({ done: false, value: ': \n\n' }) + await expect(reader.read()).rejects.toThrow(message) + expect(hasFinally).toBe(true) + }) + + it('reports the serialization error when releasing the iterator throws', async () => { + async function* gen() { + try { + yield { big: 1n } + } + finally { + // eslint-disable-next-line no-unsafe-finally + throw new Error('cleanup') + } + } + + const reader = toEventStream(gen()) + .pipeThrough(new TextDecoderStream()) + .getReader() + + await reader.read() + await expect(reader.read()).rejects.toThrow('BigInt') + }) + it('when canceled from client - return', async () => { let hasFinally = false diff --git a/packages/fetch/src/event-stream.ts b/packages/fetch/src/event-stream.ts index 723ac72..0b34cc2 100644 --- a/packages/fetch/src/event-stream.ts +++ b/packages/fetch/src/event-stream.ts @@ -146,6 +146,8 @@ export function toEventStream( } }, async pull(controller) { + let result: IteratorResult + try { if (keepAliveEnabled) { timeout = setInterval(() => { @@ -155,28 +157,7 @@ export function toEventStream( }, keepAliveInterval) } - const result = await iterator.next() - - clearInterval(timeout) - - if (cancelled) { - return - } - - const [data, meta] = unwrapEvent(result.value) - - if (!result.done || data !== undefined || meta !== undefined || emptyCloseEventEnabled) { - const event = result.done ? 'close' : 'message' - controller.enqueue(encodeEventStreamMessage({ - ...meta, - event, - data: stringifyJSON(data), - })) - } - - if (result.done) { - controller.close() - } + result = await iterator.next() } catch (err) { clearInterval(timeout) @@ -199,6 +180,47 @@ export function toEventStream( */ controller.error(err) } + + return + } + + clearInterval(timeout) + + if (cancelled) { + return + } + + try { + const [data, meta] = unwrapEvent(result.value) + + if (!result.done || data !== undefined || meta !== undefined || emptyCloseEventEnabled) { + const event = result.done ? 'close' : 'message' + controller.enqueue(encodeEventStreamMessage({ + ...meta, + event, + data: stringifyJSON(data), + })) + } + } + catch (err) { + /** + * The event could not be serialized (e.g. BigInt, circular data, a throwing toJSON). + * An errored stream never calls `cancel()`, so release the suspended iterator here. + */ + try { + if (!result.done) { + await iterator.return?.() + } + } + finally { + controller.error(err) + } + + return + } + + if (result.done) { + controller.close() } }, async cancel() {