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
9 changes: 9 additions & 0 deletions .changeset/clean-stream-readers.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
---
'@tanstack/ai-client': patch
---

Cancel fetch response bodies when SSE or NDJSON parsing fails, the consumer exits early, or SSE reaches a `[DONE]` marker. This closes unfinished connections and preserves the original error if cancellation fails.

Also cancel streaming generation responses when the client stops reading, including after a `RUN_ERROR` event.

Start cancellation without waiting for custom hooks, so pending hooks cannot delay errors, iterator completion, or reader lock release.
4 changes: 4 additions & 0 deletions docs/chat/connection-adapters.md
Original file line number Diff line number Diff line change
Expand Up @@ -691,6 +691,10 @@ stop(); // aborts the active stream

For `SubscribeConnectionAdapter`, the signal in `subscribe()` ends the entire subscription (component unmount); the signal in `send()` ends just the in-flight send.

If parsing fails or you exit the chunk iterator early, the fetch adapters cancel the response body. The SSE adapter also cancels after a `[DONE]` marker. Responses that reach their normal end keep all their chunks.

Custom cancellation hooks do not delay errors or early iterator returns. The adapters release the reader lock even if cancellation fails or stays pending.

## Error Handling

Adapters should throw on transport errors (HTTP non-2xx, parse failures, dropped sockets). The `ChatClient` catches the throw, emits a `RUN_ERROR` chunk if none has been emitted yet, and surfaces it via `onError` / the `error` state:
Expand Down
4 changes: 4 additions & 0 deletions docs/media/generations.md
Original file line number Diff line number Diff line change
Expand Up @@ -165,6 +165,10 @@ const { generate, result, isLoading } = useGenerateImage({

Combines the best of both: **type-safe input** from the fetcher pattern with **streaming** from a server function that returns an SSE `Response`. When the fetcher returns a `Response` object (instead of a plain result), the client automatically parses it as an SSE stream.

If the client stops reading before the response ends, it cancels the response body. A `RUN_ERROR` event also closes the unfinished response.

Custom cancellation hooks do not delay the client error or loading state. The client releases the reader lock even if cancellation fails or stays pending.

**Server:**

```typescript ignore
Expand Down
6 changes: 6 additions & 0 deletions packages/ai-client/src/connection-adapters.ts
Original file line number Diff line number Diff line change
Expand Up @@ -359,6 +359,12 @@ async function* readStreamLines(
throw new StreamTruncatedError()
}
} finally {
try {
// Custom cancellation hooks must not delay errors or iterator return.
void reader.cancel().catch(() => {})
} catch {
// Preserve the original error if cancellation throws synchronously.
}
reader.releaseLock()
}
}
Expand Down
6 changes: 6 additions & 0 deletions packages/ai-client/src/sse-parser.ts
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,12 @@ async function* readStreamLines(
yield buffer
}
} finally {
try {
// Custom cancellation hooks must not delay errors or iterator return.
void reader.cancel().catch(() => {})
} catch {
// Preserve the original error if cancellation throws synchronously.
}
reader.releaseLock()
}
}
Expand Down
105 changes: 105 additions & 0 deletions packages/ai-client/tests/connection-adapters-cleanup.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,105 @@
import { describe, expect, it, vi } from 'vitest'
import {
fetchHttpStream,
fetchServerSentEvents,
} from '../src/connection-adapters'

const chunk = { type: 'RUN_STARTED', threadId: 'thread', runId: 'run' }

describe.each([
['SSE', fetchServerSentEvents, (data: string) => `data: ${data}\n\n`],
['NDJSON', fetchHttpStream, (data: string) => `${data}\n`],
] as const)('%s response cleanup', (_name, createAdapter, frame) => {
function connection(
data: string,
close = false,
onCancel?: () => void | Promise<void>,
) {
const cancel = vi.fn(onCancel)
const body = new ReadableStream<Uint8Array>({
start(controller) {
controller.enqueue(new TextEncoder().encode(frame(data)))
if (close) controller.close()
},
cancel,
})
const adapter = createAdapter('/chat', {
fetchClient: async () => new Response(body),
})
return { body, cancel, stream: adapter.connect([]) }
}

it('cancels the response when JSON parsing fails', async () => {
const { body, cancel, stream } = connection('{invalid json}')
await expect(stream[Symbol.asyncIterator]().next()).rejects.toBeInstanceOf(
SyntaxError,
)
expect(cancel).toHaveBeenCalledOnce()
expect(body.locked).toBe(false)
})

it.each([false, true])(
'cancels the response on early return (pending cancellation: %s)',
async (pendingCancellation) => {
const { body, cancel, stream } = connection(
JSON.stringify(chunk),
false,
pendingCancellation ? () => new Promise<void>(() => {}) : undefined,
)
for await (const received of stream) {
expect(received).toMatchObject(chunk)
break
}
expect(cancel).toHaveBeenCalledOnce()
expect(body.locked).toBe(false)
},
)

it.each(['throws', 'rejects', 'never settles'] as const)(
'preserves a parse error when cancellation %s',
async (behavior) => {
const { body, cancel, stream } = connection(
'{invalid json}',
false,
() => {
const error = new Error('cleanup failed')
if (behavior === 'throws') throw error
if (behavior === 'rejects') return Promise.reject(error)
return new Promise<void>(() => {})
},
)
await expect(
stream[Symbol.asyncIterator]().next(),
).rejects.toBeInstanceOf(SyntaxError)
expect(cancel).toHaveBeenCalledOnce()
expect(body.locked).toBe(false)
},
)

it('drains a normally closed response without canceling its source', async () => {
const { body, cancel, stream } = connection(JSON.stringify(chunk), true)
const chunks = []
for await (const received of stream) chunks.push(received)
expect(chunks).toEqual([chunk])
expect(cancel).not.toHaveBeenCalled()
expect(body.locked).toBe(false)
})
})

it('cancels an SSE response after the DONE sentinel', async () => {
const cancel = vi.fn()
const body = new ReadableStream<Uint8Array>({
start(controller) {
controller.enqueue(new TextEncoder().encode('data: [DONE]\n\n'))
},
cancel,
})
const adapter = fetchServerSentEvents('/chat', {
fetchClient: async () => new Response(body),
})
const chunks = []
for await (const received of adapter.connect([])) chunks.push(received)
expect(chunks).toEqual([expect.objectContaining({ type: 'RUN_FINISHED' })])
expect(cancel).toHaveBeenCalledOnce()
expect(body.locked).toBe(false)
})
95 changes: 95 additions & 0 deletions packages/ai-client/tests/sse-parser-cleanup.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
import { expect, it, vi } from 'vitest'
import { parseSSEResponse } from '../src/sse-parser'
import { GenerationClient } from '../src/generation-client'

const chunk = { type: 'RUN_STARTED', threadId: 'thread', runId: 'run' }
const frame = `data: ${JSON.stringify(chunk)}\n\n`

it.each(['resolves', 'throws', 'rejects', 'never settles'] as const)(
'cancels a generation response on early return when cancellation %s',
async (behavior) => {
const cancel = vi.fn(() => {
const error = new Error('cleanup failed')
if (behavior === 'throws') throw error
if (behavior === 'rejects') return Promise.reject(error)
if (behavior === 'never settles') return new Promise<void>(() => {})
return undefined
})
const body = new ReadableStream<Uint8Array>({
start(controller) {
controller.enqueue(new TextEncoder().encode(frame))
},
cancel,
})
for await (const received of parseSSEResponse(new Response(body))) {
expect(received).toEqual(chunk)
break
}
expect(cancel).toHaveBeenCalledOnce()
expect(body.locked).toBe(false)
},
)

it('reports a generation error without waiting for custom cancellation', async () => {
const cancel = vi.fn(() => new Promise<void>(() => {}))
const body = new ReadableStream<Uint8Array>({
start(controller) {
controller.enqueue(
new TextEncoder().encode(
'data: {"type":"RUN_ERROR","message":"generation failed"}\n\n',
),
)
},
cancel,
})
const client = new GenerationClient({
fetcher: async () => new Response(body),
})
try {
await client.generate({})
expect(client.getError()?.message).toBe('generation failed')
expect(client.getIsLoading()).toBe(false)
expect(cancel).toHaveBeenCalledOnce()
expect(body.locked).toBe(false)
} finally {
client.dispose()
}
})

it('preserves generation reader errors and releases the lock', async () => {
const error = new Error('response reader failed')
const body = new ReadableStream<Uint8Array>({
start(controller) {
controller.error(error)
},
})
await expect(parseSSEResponse(new Response(body)).next()).rejects.toBe(error)
expect(body.locked).toBe(false)
})

it('drains generation EOF and keeps ignoring malformed JSON and DONE markers', async () => {
const warn = vi.spyOn(console, 'warn').mockImplementation(() => {})
const cancel = vi.fn()
const body = new ReadableStream<Uint8Array>({
start(controller) {
controller.enqueue(
new TextEncoder().encode(
`data: {invalid json}\n\ndata: [DONE]\n\n${frame}`,
),
)
controller.close()
},
cancel,
})
try {
const chunks = []
for await (const received of parseSSEResponse(new Response(body)))
chunks.push(received)
expect(chunks).toEqual([chunk])
expect(warn).toHaveBeenCalledTimes(2)
expect(cancel).not.toHaveBeenCalled()
expect(body.locked).toBe(false)
} finally {
warn.mockRestore()
}
})
77 changes: 77 additions & 0 deletions testing/e2e/tests/error-handling.spec.ts
Original file line number Diff line number Diff line change
@@ -1,8 +1,85 @@
import { createServer } from 'node:http'
import { GenerationClient } from '@tanstack/ai-client'
import { test, expect } from './fixtures'
import { selectScenario, runTest, getMetadata } from './tools-test/helpers'
import { sendMessage, featureUrl } from './helpers'

test.describe('Error Handling', () => {
test('closes a generation fetcher response after RUN_ERROR', async () => {
let responseClosed = false
const server = createServer((_request, response) => {
response.writeHead(200, { 'Content-Type': 'text/event-stream' })
response.once('close', () => {
responseClosed = true
})
response.write(
'data: {"type":"RUN_ERROR","message":"generation failed"}\n\n',
)
})
const client = new GenerationClient({
fetcher: async (_input, { signal }) => {
const address = server.address()
if (!address || typeof address === 'string')
throw new Error('Missing test server port')
return fetch(`http://127.0.0.1:${address.port}`, { signal })
},
})
await new Promise<void>((resolve) => server.listen(0, '127.0.0.1', resolve))
try {
await client.generate({})
expect(client.getError()?.message).toBe('generation failed')
expect(client.getIsLoading()).toBe(false)
await expect.poll(() => responseClosed).toBe(true)
} finally {
client.dispose()
server.closeAllConnections()
await new Promise<void>((resolve) => server.close(() => resolve()))
}
})

test('closes the HTTP response after a stream parse error', async ({
page,
testId,
aimockPort,
}) => {
let responseClosed = false
// Keep the response open so only client cleanup can close the connection.
const server = createServer((request, response) => {
request.resume()
response.setHeader('Access-Control-Allow-Origin', '*')
response.setHeader('Access-Control-Allow-Headers', 'Content-Type')
if (request.method === 'OPTIONS') {
response.writeHead(204).end()
return
}
response.writeHead(200, { 'Content-Type': 'text/event-stream' })
response.once('close', () => {
responseClosed = true
})
response.write('data: {invalid json}\n\n')
})
await new Promise<void>((resolve) => server.listen(0, '127.0.0.1', resolve))
try {
const address = server.address()
if (!address || typeof address === 'string')
throw new Error('Missing test server port')
await page.route('**/api/tools-test', (route) =>
route.continue({ url: `http://127.0.0.1:${address.port}` }),
)
await selectScenario(page, 'error', testId, aimockPort)
await runTest(page)
await expect(page.locator('#error-display')).toBeVisible()
await expect(page.locator('#test-metadata')).toHaveAttribute(
'data-is-loading',
'false',
)
await expect.poll(() => responseClosed).toBe(true)
} finally {
server.closeAllConnections()
await new Promise<void>((resolve) => server.close(() => resolve()))
}
})

test('displays error when server returns RUN_ERROR', async ({
page,
testId,
Expand Down
Loading