From 005e3b00fd38e641f1add05ef88d77b90a0cd7b9 Mon Sep 17 00:00:00 2001 From: Jacob Quinn Date: Fri, 2 Oct 2026 00:07:01 -0600 Subject: [PATCH 1/2] fix(ai-client): cancel abandoned fetch response streams --- .changeset/clean-stream-readers.md | 7 ++ docs/chat/connection-adapters.md | 2 + docs/media/generations.md | 2 + packages/ai-client/src/connection-adapters.ts | 5 ++ packages/ai-client/src/sse-parser.ts | 5 ++ .../tests/connection-adapters-cleanup.test.ts | 88 +++++++++++++++++++ .../tests/sse-parser-cleanup.test.ts | 64 ++++++++++++++ testing/e2e/tests/error-handling.spec.ts | 77 ++++++++++++++++ 8 files changed, 250 insertions(+) create mode 100644 .changeset/clean-stream-readers.md create mode 100644 packages/ai-client/tests/connection-adapters-cleanup.test.ts create mode 100644 packages/ai-client/tests/sse-parser-cleanup.test.ts diff --git a/.changeset/clean-stream-readers.md b/.changeset/clean-stream-readers.md new file mode 100644 index 0000000000..277fdbcc90 --- /dev/null +++ b/.changeset/clean-stream-readers.md @@ -0,0 +1,7 @@ +--- +'@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. diff --git a/docs/chat/connection-adapters.md b/docs/chat/connection-adapters.md index cee42186ce..60ffabbb38 100644 --- a/docs/chat/connection-adapters.md +++ b/docs/chat/connection-adapters.md @@ -691,6 +691,8 @@ 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. + ## 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: diff --git a/docs/media/generations.md b/docs/media/generations.md index 6699ccba95..18ada3cdec 100644 --- a/docs/media/generations.md +++ b/docs/media/generations.md @@ -165,6 +165,8 @@ 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. + **Server:** ```typescript ignore diff --git a/packages/ai-client/src/connection-adapters.ts b/packages/ai-client/src/connection-adapters.ts index d5ebafab20..c9c091bab7 100644 --- a/packages/ai-client/src/connection-adapters.ts +++ b/packages/ai-client/src/connection-adapters.ts @@ -359,6 +359,11 @@ async function* readStreamLines( throw new StreamTruncatedError() } } finally { + try { + await reader.cancel() + } catch { + // A failed stream can reject cancellation. Preserve the original error. + } reader.releaseLock() } } diff --git a/packages/ai-client/src/sse-parser.ts b/packages/ai-client/src/sse-parser.ts index 5c85a842ae..6b8cd7501e 100644 --- a/packages/ai-client/src/sse-parser.ts +++ b/packages/ai-client/src/sse-parser.ts @@ -37,6 +37,11 @@ async function* readStreamLines( yield buffer } } finally { + try { + await reader.cancel() + } catch { + // A failed stream can reject cancellation. Preserve the original error. + } reader.releaseLock() } } diff --git a/packages/ai-client/tests/connection-adapters-cleanup.test.ts b/packages/ai-client/tests/connection-adapters-cleanup.test.ts new file mode 100644 index 0000000000..198fa208d1 --- /dev/null +++ b/packages/ai-client/tests/connection-adapters-cleanup.test.ts @@ -0,0 +1,88 @@ +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, cancelError?: Error) { + const cancel = vi.fn(() => { + if (cancelError) throw cancelError + }) + const body = new ReadableStream({ + 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('cancels the response when the consumer stops early', async () => { + const { body, cancel, stream } = connection(JSON.stringify(chunk)) + for await (const received of stream) { + expect(received).toMatchObject(chunk) + break + } + expect(cancel).toHaveBeenCalledOnce() + expect(body.locked).toBe(false) + }) + + it('preserves a parse error when cancellation rejects', async () => { + const { body, cancel, stream } = connection( + '{invalid json}', + false, + new Error('cleanup failed'), + ) + 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({ + 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) +}) diff --git a/packages/ai-client/tests/sse-parser-cleanup.test.ts b/packages/ai-client/tests/sse-parser-cleanup.test.ts new file mode 100644 index 0000000000..a4d71cc4c0 --- /dev/null +++ b/packages/ai-client/tests/sse-parser-cleanup.test.ts @@ -0,0 +1,64 @@ +import { expect, it, vi } from 'vitest' +import { parseSSEResponse } from '../src/sse-parser' + +const chunk = { type: 'RUN_STARTED', threadId: 'thread', runId: 'run' } +const frame = `data: ${JSON.stringify(chunk)}\n\n` + +it.each([false, true])( + 'cancels a generation response on early return (cancellation rejects: %s)', + async (rejectCancellation) => { + const cancel = vi.fn(() => { + if (rejectCancellation) throw new Error('cleanup failed') + }) + const body = new ReadableStream({ + 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('preserves generation reader errors and releases the lock', async () => { + const error = new Error('response reader failed') + const body = new ReadableStream({ + 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({ + 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() + } +}) diff --git a/testing/e2e/tests/error-handling.spec.ts b/testing/e2e/tests/error-handling.spec.ts index 2779d0e72d..f5b5dc89ee 100644 --- a/testing/e2e/tests/error-handling.spec.ts +++ b/testing/e2e/tests/error-handling.spec.ts @@ -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((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((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((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((resolve) => server.close(() => resolve())) + } + }) + test('displays error when server returns RUN_ERROR', async ({ page, testId, From 981c49ab61cc9c28323ba189f2e0d42b49e96dde Mon Sep 17 00:00:00 2001 From: Jacob Quinn Date: Fri, 2 Oct 2026 00:33:12 -0600 Subject: [PATCH 2/2] fix(ai-client): release readers without waiting for cancellation --- .changeset/clean-stream-readers.md | 2 + docs/chat/connection-adapters.md | 2 + docs/media/generations.md | 2 + packages/ai-client/src/connection-adapters.ts | 5 +- packages/ai-client/src/sse-parser.ts | 5 +- .../tests/connection-adapters-cleanup.test.ts | 67 ++++++++++++------- .../tests/sse-parser-cleanup.test.ts | 39 +++++++++-- 7 files changed, 89 insertions(+), 33 deletions(-) diff --git a/.changeset/clean-stream-readers.md b/.changeset/clean-stream-readers.md index 277fdbcc90..64ce96f74f 100644 --- a/.changeset/clean-stream-readers.md +++ b/.changeset/clean-stream-readers.md @@ -5,3 +5,5 @@ 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. diff --git a/docs/chat/connection-adapters.md b/docs/chat/connection-adapters.md index 60ffabbb38..2e6e123755 100644 --- a/docs/chat/connection-adapters.md +++ b/docs/chat/connection-adapters.md @@ -693,6 +693,8 @@ For `SubscribeConnectionAdapter`, the signal in `subscribe()` ends the entire su 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: diff --git a/docs/media/generations.md b/docs/media/generations.md index 18ada3cdec..1af7f551ce 100644 --- a/docs/media/generations.md +++ b/docs/media/generations.md @@ -167,6 +167,8 @@ Combines the best of both: **type-safe input** from the fetcher pattern with **s 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 diff --git a/packages/ai-client/src/connection-adapters.ts b/packages/ai-client/src/connection-adapters.ts index c9c091bab7..424ab275cc 100644 --- a/packages/ai-client/src/connection-adapters.ts +++ b/packages/ai-client/src/connection-adapters.ts @@ -360,9 +360,10 @@ async function* readStreamLines( } } finally { try { - await reader.cancel() + // Custom cancellation hooks must not delay errors or iterator return. + void reader.cancel().catch(() => {}) } catch { - // A failed stream can reject cancellation. Preserve the original error. + // Preserve the original error if cancellation throws synchronously. } reader.releaseLock() } diff --git a/packages/ai-client/src/sse-parser.ts b/packages/ai-client/src/sse-parser.ts index 6b8cd7501e..0c8a70a159 100644 --- a/packages/ai-client/src/sse-parser.ts +++ b/packages/ai-client/src/sse-parser.ts @@ -38,9 +38,10 @@ async function* readStreamLines( } } finally { try { - await reader.cancel() + // Custom cancellation hooks must not delay errors or iterator return. + void reader.cancel().catch(() => {}) } catch { - // A failed stream can reject cancellation. Preserve the original error. + // Preserve the original error if cancellation throws synchronously. } reader.releaseLock() } diff --git a/packages/ai-client/tests/connection-adapters-cleanup.test.ts b/packages/ai-client/tests/connection-adapters-cleanup.test.ts index 198fa208d1..ed3e3d8b68 100644 --- a/packages/ai-client/tests/connection-adapters-cleanup.test.ts +++ b/packages/ai-client/tests/connection-adapters-cleanup.test.ts @@ -10,10 +10,12 @@ 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, cancelError?: Error) { - const cancel = vi.fn(() => { - if (cancelError) throw cancelError - }) + function connection( + data: string, + close = false, + onCancel?: () => void | Promise, + ) { + const cancel = vi.fn(onCancel) const body = new ReadableStream({ start(controller) { controller.enqueue(new TextEncoder().encode(frame(data))) @@ -36,28 +38,43 @@ describe.each([ expect(body.locked).toBe(false) }) - it('cancels the response when the consumer stops early', async () => { - const { body, cancel, stream } = connection(JSON.stringify(chunk)) - for await (const received of stream) { - expect(received).toMatchObject(chunk) - break - } - 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(() => {}) : undefined, + ) + for await (const received of stream) { + expect(received).toMatchObject(chunk) + break + } + expect(cancel).toHaveBeenCalledOnce() + expect(body.locked).toBe(false) + }, + ) - it('preserves a parse error when cancellation rejects', async () => { - const { body, cancel, stream } = connection( - '{invalid json}', - false, - new Error('cleanup failed'), - ) - await expect(stream[Symbol.asyncIterator]().next()).rejects.toBeInstanceOf( - SyntaxError, - ) - 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(() => {}) + }, + ) + 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) diff --git a/packages/ai-client/tests/sse-parser-cleanup.test.ts b/packages/ai-client/tests/sse-parser-cleanup.test.ts index a4d71cc4c0..76d0a52842 100644 --- a/packages/ai-client/tests/sse-parser-cleanup.test.ts +++ b/packages/ai-client/tests/sse-parser-cleanup.test.ts @@ -1,14 +1,19 @@ 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([false, true])( - 'cancels a generation response on early return (cancellation rejects: %s)', - async (rejectCancellation) => { +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(() => { - if (rejectCancellation) throw new Error('cleanup failed') + 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(() => {}) + return undefined }) const body = new ReadableStream({ start(controller) { @@ -25,6 +30,32 @@ it.each([false, true])( }, ) +it('reports a generation error without waiting for custom cancellation', async () => { + const cancel = vi.fn(() => new Promise(() => {})) + const body = new ReadableStream({ + 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({