Skip to content
Merged
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
49 changes: 49 additions & 0 deletions packages/fetch/src/event-stream.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, unknown> = {}
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

Expand Down
66 changes: 44 additions & 22 deletions packages/fetch/src/event-stream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -146,6 +146,8 @@ export function toEventStream(
}
},
async pull(controller) {
let result: IteratorResult<unknown>

try {
if (keepAliveEnabled) {
timeout = setInterval(() => {
Expand All @@ -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)
Expand All @@ -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() {
Expand Down
Loading