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() {