From 2ea65e898c67f5aeeb73925eb5c5a055352d458d Mon Sep 17 00:00:00 2001 From: Alem Tuzlak Date: Mon, 5 Oct 2026 12:23:32 +0200 Subject: [PATCH 1/4] fix(ai): stream durable responses live instead of in batches of 32 With `durability` set, a chunk reached the client only after its batch was appended, and a batch was appended only at 32 chunks or the run end. A short reply showed up all at once at the end. The batch now also flushes when the producer sends no new chunk for `batchWaitMs` (default 50ms). `batch` stays the most chunks in one append. Set `batchWaitMs` on `durability` or on `toWebSocketStream`. Claude-Session: https://claude.ai/code/session_01APYv1qshKyjPPpkFyRZhfZ --- .changeset/durable-batch-live-flush.md | 7 + docs/config.json | 2 +- docs/resumable-streams/advanced.md | 21 ++- packages/ai/src/stream-to-response.ts | 145 ++++++++++++++++-- packages/ai/src/stream-to-websocket.ts | 8 + .../stream-to-response-durability.test.ts | 88 +++++++++++ .../e2e/src/routes/api.durable-delivery.ts | 31 +++- testing/e2e/tests/delivery-durability.spec.ts | 39 +++++ 8 files changed, 318 insertions(+), 23 deletions(-) create mode 100644 .changeset/durable-batch-live-flush.md diff --git a/.changeset/durable-batch-live-flush.md b/.changeset/durable-batch-live-flush.md new file mode 100644 index 0000000000..44682e5c1b --- /dev/null +++ b/.changeset/durable-batch-live-flush.md @@ -0,0 +1,7 @@ +--- +'@tanstack/ai': patch +--- + +Stream durable responses live again. With `durability` set, a chunk waited in the append batch until 32 chunks arrived or the run finished, so a short reply showed up all at once at the end. Now the batch also flushes when the model sends no new chunk for 50ms. You do not need `batch: 1` for live text any more. `batch` stays the largest number of chunks in one append. + +Set the wait with the new `batchWaitMs` option on `durability` (`toServerSentEventsResponse`, `toHttpResponse`) and on `toWebSocketStream` / `toWebSocketResponse`. A higher value means fewer writes to the log but slower live text. `0` appends every chunk on its own. diff --git a/docs/config.json b/docs/config.json index 2d5fa8669f..d2567b4e12 100644 --- a/docs/config.json +++ b/docs/config.json @@ -608,7 +608,7 @@ "label": "Advanced", "to": "resumable-streams/advanced", "addedAt": "2026-08-04", - "updatedAt": "2026-08-20" + "updatedAt": "2026-10-05" }, { "label": "WebSockets", diff --git a/docs/resumable-streams/advanced.md b/docs/resumable-streams/advanced.md index 8b427cd4bd..2c3f25271e 100644 --- a/docs/resumable-streams/advanced.md +++ b/docs/resumable-streams/advanced.md @@ -48,14 +48,17 @@ export async function POST(request: Request) { runId, }) return toServerSentEventsResponse(stream, { - durability: { adapter: durableStream(request, durableOptions), batch: 32 }, + durability: { + adapter: durableStream(request, durableOptions), + batch: 32, + batchWaitMs: 50, + }, }) } ``` - `headers` takes a static object for fixed credentials or an async resolver for rotating tokens. The resolver runs for every create, append, read, and close. -- `batch` controls how many chunks are buffered per log append (default 32). - The backend must return a non-empty `Stream-Next-Offset` header on create, append, and close. A missing header fails loudly. The adapter never guesses an offset. @@ -69,6 +72,20 @@ export async function POST(request: Request) { not need it: see [Resuming a run without duplicating what you already streamed](#resuming-a-run-without-duplicating-what-you-already-streamed). +### Batch the log writes + +Every append to a remote backend is one request. A chunk reaches the client only +after its append. So two options on `durability` trade writes against how fast +the text appears: + +- `batch`: the most chunks in one append. The default is 32. +- `batchWaitMs`: the most milliseconds that a chunk waits for more chunks. Then + its batch is appended and sent. The default is 50. + +If your backend charges per write, raise `batchWaitMs` and `batch`. The text then +arrives in bigger steps. If the reply arrives in large blocks, lower +`batchWaitMs`. Set it to `0` to append every chunk on its own. + ## Attaching to a run by id Reconnect after a drop is automatic. To attach to a run from the start on diff --git a/packages/ai/src/stream-to-response.ts b/packages/ai/src/stream-to-response.ts index 40858c3063..cf95394fbb 100644 --- a/packages/ai/src/stream-to-response.ts +++ b/packages/ai/src/stream-to-response.ts @@ -305,6 +305,63 @@ function sseEncoders( /** Default number of chunks buffered before a durability `append`. */ const DEFAULT_DURABILITY_BATCH = 32 +/** + * Default `batchWaitMs`: the longest a buffered chunk waits for the producer's + * next chunk before its batch flushes anyway. Chunks are forwarded only AFTER + * they are appended, so a size-only batch held live text until 32 chunks + * arrived: a short reply showed up all at once at `RUN_FINISHED`. This bounds + * that wait and keeps bursts in one `append`, so a remote log still does not + * pay one write per token. + */ +const DEFAULT_DURABILITY_BATCH_WAIT_MS = 50 + +/** Largest delay `setTimeout` honors; a larger one fires at once. */ +const MAX_TIMER_DELAY_MS = 2_147_483_647 + +/** + * Resolve and validate `batchWaitMs`. `NaN`, a negative value, `Infinity`, and + * anything past the timer limit are rejected: `setTimeout` would turn each of + * them into an immediate flush, the opposite of what a large value asks for. + */ +function resolveBatchWaitMs(batchWaitMs: number | undefined): number { + if (batchWaitMs === undefined) return DEFAULT_DURABILITY_BATCH_WAIT_MS + if ( + !Number.isFinite(batchWaitMs) || + batchWaitMs < 0 || + batchWaitMs > MAX_TIMER_DELAY_MS + ) { + throw new Error( + `Invalid durability batchWaitMs: ${batchWaitMs}. Must be a number from 0 to ${MAX_TIMER_DELAY_MS}.`, + ) + } + return batchWaitMs +} + +/** + * True when `promise` settles within `ms`. The promise is observed either way, + * so a rejection that loses the race is not reported as unhandled. + */ +async function settlesWithin( + promise: Promise, + ms: number, +): Promise { + if (ms <= 0) return false + let timer: ReturnType | undefined + try { + return await Promise.race([ + promise.then( + () => true, + () => true, + ), + new Promise((resolve) => { + timer = setTimeout(() => resolve(false), ms) + }), + ]) + } finally { + clearTimeout(timer) + } +} + /** * Resolve and validate the durability batch size. A non-positive-integer (0, * negative, fractional, or `NaN`) is rejected rather than clamped: silently @@ -378,8 +435,9 @@ export const RUN_ACCEPTED_EVENT = 'run.accepted' * so `chat()`'s lazy iterator never fires the provider — the untouched * generator is simply GC'd. This is what makes resume free of re-invocation. * - **Fresh** (`resumeFrom()` null): iterate `stream`, buffering up to `batch` - * chunks (flushing early at terminal / tool-call boundaries), `append` each - * batch to the log, then forward. Appending BEFORE forwarding guarantees a + * chunks (flushing early at terminal / tool-call boundaries, or once the + * oldest buffered chunk has waited `batchWaitMs`), `append` each batch to + * the log, then forward. Appending BEFORE forwarding guarantees a * reconnecting client can always replay exactly what it already saw. * * The returned `getId` maps each forwarded chunk to the exact opaque offset @@ -391,6 +449,7 @@ export function durableStreamSource( options: { abortController: AbortController batch?: number + batchWaitMs?: number logger?: InternalLogger }, ): { @@ -399,6 +458,7 @@ export function durableStreamSource( } { const resumeOffset = durability.resumeFrom() const batchSize = resolveBatchSize(options.batch) + const batchWaitMs = resolveBatchWaitMs(options.batchWaitMs) const abortController = options.abortController const logger = options.logger const idByChunk = new WeakMap() @@ -502,12 +562,45 @@ export function durableStreamSource( timestamp: Date.now(), }) yield* flush() - for await (const chunk of stream) { - if (isAborted(abortController.signal)) break - batch.push(chunk) - if (batch.length >= batchSize || isDurabilityFlushBoundary(chunk)) { - yield* flush() + // Iterated by hand, not with `for await`, so the batch can flush while + // the next chunk is still pending. The `finally` does what `for await` + // would: close the producer on an early exit, but not after it finished + // or threw. + const iterator = stream[Symbol.asyncIterator]() + let iteratorFinished = false + let flushBy = 0 + try { + for (;;) { + const next = iterator.next() + if ( + batch.length > 0 && + !(await settlesWithin(next, flushBy - Date.now())) + ) { + yield* flush() + } + let result: IteratorResult + try { + result = await next + } catch (error) { + iteratorFinished = true + throw error + } + if (result.done) { + iteratorFinished = true + break + } + if (isAborted(abortController.signal)) break + const chunk = result.value + if (batch.length === 0) { + flushBy = Date.now() + batchWaitMs + } + batch.push(chunk) + if (batch.length >= batchSize || isDurabilityFlushBoundary(chunk)) { + yield* flush() + } } + } finally { + if (!iteratorFinished) await iterator.return?.() } if (!isAborted(abortController.signal)) yield* flush() } catch (error) { @@ -695,10 +788,11 @@ export function durableStreamSource( * to make the stream resumable: fresh runs are appended to the log and each SSE * event is tagged with an `id:` offset; a reconnect (native `Last-Event-ID`) or * a `?offset` join replays from the log without re-running the producer. `batch` - * controls how many chunks are buffered per `append` (default 32). + * controls how many chunks are buffered per `append` (default 32), and + * `batchWaitMs` how long a buffered chunk waits for more (default 50). * * @param stream - AsyncIterable of StreamChunks from chat() - * @param init - Optional Response initialization options (including `abortController`, `durability` with its optional `batch`, and `debug`) + * @param init - Optional Response initialization options (including `abortController`, `durability` with its optional `batch` and `batchWaitMs`, and `debug`) * @returns Response in Server-Sent Events format * * @example @@ -713,7 +807,18 @@ export function toServerSentEventsResponse( stream: AsyncIterable, init?: ResponseInit & { abortController?: AbortController - durability?: { adapter: StreamDurability; batch?: number } + durability?: { + adapter: StreamDurability + /** Most chunks in one `append` (default 32). */ + batch?: number + /** + * Most ms a chunk waits in the batch for the next one before the batch + * is appended and sent anyway (default 50). Chunks reach the client only + * after their append, so a higher value means fewer writes but slower + * live text, and `0` appends every chunk on its own. + */ + batchWaitMs?: number + } /** * Customize logging for durability failure paths (terminal-append and * close). These failures are always logged server-side by default (the @@ -761,6 +866,7 @@ export function toServerSentEventsResponse( const { source, getId } = durableStreamSource(stream, durability.adapter, { abortController: producerAbortController, batch: durability.batch, + batchWaitMs: durability.batchWaitMs, // `errors` category is on by default even when `debug` is undefined, so // durability terminal-append / close failures always surface server-side — // including on the client-disconnect path where there is no live consumer. @@ -1109,12 +1215,13 @@ function ndjsonEncoders( * NDJSON line is emitted as an `{ id, chunk }` envelope carrying an opaque * offset; a reconnect (native `Last-Event-ID` header) or a `?offset` join * replays from the log without re-running the producer. `batch` controls how - * many chunks are buffered per `append` (default 32). This shares the exact + * many chunks are buffered per `append` (default 32), and `batchWaitMs` how + * long a buffered chunk waits for more (default 50). This shares the exact * `durableStreamSource` used by `toServerSentEventsResponse` — only the wire * encoding differs. * * @param stream - AsyncIterable of StreamChunks from chat() - * @param init - Optional Response initialization options (including `abortController`, `durability` with its optional `batch`, and `debug`) + * @param init - Optional Response initialization options (including `abortController`, `durability` with its optional `batch` and `batchWaitMs`, and `debug`) * @returns Response in HTTP stream format (newline-delimited JSON) * * @example @@ -1129,7 +1236,18 @@ export function toHttpResponse( stream: AsyncIterable, init?: ResponseInit & { abortController?: AbortController - durability?: { adapter: StreamDurability; batch?: number } + durability?: { + adapter: StreamDurability + /** Most chunks in one `append` (default 32). */ + batch?: number + /** + * Most ms a chunk waits in the batch for the next one before the batch + * is appended and sent anyway (default 50). Chunks reach the client only + * after their append, so a higher value means fewer writes but slower + * live text, and `0` appends every chunk on its own. + */ + batchWaitMs?: number + } /** * Customize logging for durability failure paths (terminal-append and * close). These failures are always logged server-side by default (the @@ -1172,6 +1290,7 @@ export function toHttpResponse( const { source, getId } = durableStreamSource(stream, durability.adapter, { abortController: producerAbortController, batch: durability.batch, + batchWaitMs: durability.batchWaitMs, // Errors-on-by-default logger (see toServerSentEventsResponse). logger: resolveDebugOption(debug), }) diff --git a/packages/ai/src/stream-to-websocket.ts b/packages/ai/src/stream-to-websocket.ts index ec67f31fb1..0bdc0cfa27 100644 --- a/packages/ai/src/stream-to-websocket.ts +++ b/packages/ai/src/stream-to-websocket.ts @@ -94,6 +94,11 @@ export interface WebSocketStreamInit { durability?: (ctx: WsRunContext) => StreamDurability /** Chunks buffered per durability append (default 32). */ batch?: number + /** + * Most ms a chunk waits in the batch for the next one before the batch is + * appended and sent anyway (default 50). `0` appends every chunk on its own. + */ + batchWaitMs?: number /** Heartbeat ping interval in ms (default 30_000). */ heartbeatMs?: number /** @@ -250,6 +255,9 @@ export function toWebSocketStream( { abortController: turnAbort, ...(init.batch === undefined ? {} : { batch: init.batch }), + ...(init.batchWaitMs === undefined + ? {} + : { batchWaitMs: init.batchWaitMs }), logger, }, ) diff --git a/packages/ai/tests/stream-to-response-durability.test.ts b/packages/ai/tests/stream-to-response-durability.test.ts index 93f166468e..bfb4b4dedf 100644 --- a/packages/ai/tests/stream-to-response-durability.test.ts +++ b/packages/ai/tests/stream-to-response-durability.test.ts @@ -258,6 +258,94 @@ describe('toServerSentEventsResponse with durability', () => { expect(batchSizes.reduce((sum, size) => sum + size, 0)).toBe(6) }) + /** + * Ms until 'a' reaches the reader while the model is still busy after it, or + * `Infinity` when 'a' is still held after 1s. 'a' sits in a batch of 32 that + * only 1 chunk filled, so only the batch wait can send it before the run ends. + */ + async function msUntilFirstText( + runId: string, + options: { batchWaitMs?: number } = {}, + ): Promise { + let releaseModel = (): void => undefined + const modelThinking = new Promise((resolve) => { + releaseModel = resolve + }) + const stream: AsyncIterable = { + async *[Symbol.asyncIterator]() { + yield ev.textContent('a') + await modelThinking + yield ev.textContent('b') + yield ev.runFinished('stop') + }, + } + const startedAt = Date.now() + const response = toServerSentEventsResponse(stream, { + durability: { + adapter: memoryStream( + new Request(`https://example.test/api/chat?runId=${runId}`, { + method: 'POST', + }), + ), + ...options, + }, + }) + if (!response.body) throw new Error('Expected a response body') + const reader = response.body.getReader() + const decoder = new TextDecoder() + let body = '' + const readUntilA = async (): Promise => { + while (!body.includes('"delta":"a"')) { + const result = await reader.read() + if (result.done) throw new Error('Stream ended before "a"') + body += decoder.decode(result.value) + } + return Date.now() - startedAt + } + + const elapsed = await Promise.race([ + readUntilA(), + new Promise((resolve) => + setTimeout(() => resolve(Infinity), 1000), + ), + ]) + releaseModel() + for (;;) { + const result = await reader.read() + if (result.done) break + } + return elapsed + } + + it('delivers a buffered chunk while the producer waits for the next one', async () => { + expect(await msUntilFirstText('live-while-idle')).toBeLessThan(1000) + }) + + it('holds a buffered chunk for batchWaitMs', async () => { + const elapsed = await msUntilFirstText('live-wait-300', { + batchWaitMs: 300, + }) + + expect(elapsed).toBeGreaterThanOrEqual(250) + expect(elapsed).toBeLessThan(1000) + }) + + it('rejects an invalid batchWaitMs', () => { + const durability = memoryStream( + new Request('https://example.test/api/chat?runId=response-bad-wait', { + method: 'POST', + }), + ) + const { stream } = fiveChunkStream() + for (const batchWaitMs of [-1, NaN, Infinity, 2_147_483_648]) { + expect(() => + toServerSentEventsResponse(stream, { + durability: { adapter: durability, batchWaitMs }, + }), + ).toThrow(/Invalid durability batchWaitMs/) + } + }) + it('rejects a non-positive-integer batch size', () => { const durability = memoryStream( new Request('https://example.test/api/chat?runId=response-bad-batch', { diff --git a/testing/e2e/src/routes/api.durable-delivery.ts b/testing/e2e/src/routes/api.durable-delivery.ts index 8703edda11..3721e7bc8a 100644 --- a/testing/e2e/src/routes/api.durable-delivery.ts +++ b/testing/e2e/src/routes/api.durable-delivery.ts @@ -23,14 +23,23 @@ import type { StreamChunk } from '@tanstack/ai' * `?transport=ndjson` switches the wire encoding from SSE to newline-delimited * JSON (each durable line is an `{ id, chunk }` envelope). The durability layer * — logging, offsets, resume, terminalization — is identical for both. + * + * `?scenario=slow` spaces the content chunks 300ms apart and keeps the default + * batch size, so a test can check that text reaches the client while the run + * is still going. */ // Emits bare TEXT_MESSAGE_CONTENT chunks without TEXT_MESSAGE_START/END // bracketing: this harness deliberately exercises raw chunk delivery + resume, // not UIMessage reassembly. The durability layer terminalizes on RUN_FINISHED // (emitted below), which is all resume/join needs. +// +// `delayMs` waits before each content chunk, like a model that is still +// producing tokens. The `slow` scenario uses it to prove that the durability +// batch does not hold live text until the run ends. export function fixedRun( threadId: string, runId: string, + delayMs = 0, ): AsyncIterable { return (async function* () { yield { @@ -40,6 +49,9 @@ export function fixedRun( timestamp: Date.now(), } as StreamChunk for (let i = 1; i <= 5; i++) { + if (delayMs > 0) { + await new Promise((resolve) => setTimeout(resolve, delayMs)) + } yield { type: 'TEXT_MESSAGE_CONTENT', messageId: 'm', @@ -139,11 +151,11 @@ function agentLoopRun( })() } -function isAgentLoop(request: Request): boolean { +function scenarioOf(request: Request): string | null { try { - return new URL(request.url).searchParams.get('scenario') === 'agent-loop' + return new URL(request.url).searchParams.get('scenario') } catch { - return false + return null } } @@ -183,9 +195,11 @@ function durableResponse( durability: ReturnType, batch?: number, ): Response { - const stream = isAgentLoop(request) - ? agentLoopRun('thread-durable', runId) - : fixedRun('thread-durable', runId) + const scenario = scenarioOf(request) + const stream = + scenario === 'agent-loop' + ? agentLoopRun('thread-durable', runId) + : fixedRun('thread-durable', runId, scenario === 'slow' ? 300 : 0) const durabilityOption = { adapter: durability, ...(batch ? { batch } : {}) } return isNdjson(request) ? toHttpResponse(stream, { durability: durabilityOption }) @@ -197,8 +211,11 @@ export const Route = createFileRoute('/api/durable-delivery')({ handlers: { POST: async ({ request }) => { const { durability, runId, advertiseRunId } = durableRun(request) + // `slow` keeps the default batch (32): its point is that live text + // still flows while the batch is far from full. + const batch = scenarioOf(request) === 'slow' ? undefined : 2 return withRunId( - durableResponse(request, runId, durability, 2), + durableResponse(request, runId, durability, batch), advertiseRunId, ) }, diff --git a/testing/e2e/tests/delivery-durability.spec.ts b/testing/e2e/tests/delivery-durability.spec.ts index 4003a0b437..1e80682d46 100644 --- a/testing/e2e/tests/delivery-durability.spec.ts +++ b/testing/e2e/tests/delivery-durability.spec.ts @@ -111,6 +111,45 @@ test.describe('delivery durability', () => { }) }) +test.describe('delivery durability (live text)', () => { + test('streams text while the run is still going, with the default batch', async ({ + page, + }) => { + // The `slow` run sends 5 text chunks 300ms apart. Chunks reach the client + // only after they are appended to the log, and a batch of 32 never fills + // here. So the batch must flush on its own while the run waits, not hold + // every chunk until RUN_FINISHED. + await page.goto('/') + const timing = await page.evaluate(async () => { + const startedAt = performance.now() + const response = await fetch('/api/durable-delivery?scenario=slow', { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: '{}', + }) + if (!response.body) throw new Error('Expected a response body') + const reader = response.body.getReader() + const decoder = new TextDecoder() + let body = '' + let firstTextAt = -1 + for (;;) { + const { done, value } = await reader.read() + if (done) break + body += decoder.decode(value, { stream: true }) + if (firstTextAt < 0 && body.includes('"TEXT_MESSAGE_CONTENT"')) { + firstTextAt = performance.now() - startedAt + } + } + return { firstTextAt, endAt: performance.now() - startedAt } + }) + + expect(timing.firstTextAt).toBeGreaterThan(0) + // The first chunk comes about 1.2s before the end. Held until the end, the + // gap would be about 0. + expect(timing.endAt - timing.firstTextAt).toBeGreaterThan(800) + }) +}) + test.describe('delivery durability (agent loop)', () => { test('delivers a full tool-calling run, not truncated at the first RUN_FINISHED', async ({ request, From 1ed42fe51b7f41483f0348119088d97276ee01b3 Mon Sep 17 00:00:00 2001 From: Alem Tuzlak Date: Mon, 5 Oct 2026 12:23:33 +0200 Subject: [PATCH 2/4] fix(ai-client): send one hydrate GET when a view re-attaches React Strict Mode runs mount effects twice in dev: attach, detach, attach. Each attach of a `persistence: true` chat sent its own hydrate GET. A re-attach now reuses the GET that is still in flight, and that GET still paints the view when it returns. Claude-Session: https://claude.ai/code/session_01APYv1qshKyjPPpkFyRZhfZ --- .changeset/chat-hydrate-strict-mode.md | 5 +++ packages/ai-client/src/chat-client.ts | 11 +++++ .../ai-client/tests/dispose-tail-leak.test.ts | 42 +++++++++++++++++++ testing/e2e/src/routeTree.gen.ts | 21 ++++++++++ .../e2e/src/routes/client-mount-hydrate.tsx | 41 ++++++++++++++++++ .../e2e/tests/persistence-durability.spec.ts | 29 +++++++++++++ 6 files changed, 149 insertions(+) create mode 100644 .changeset/chat-hydrate-strict-mode.md create mode 100644 testing/e2e/src/routes/client-mount-hydrate.tsx diff --git a/.changeset/chat-hydrate-strict-mode.md b/.changeset/chat-hydrate-strict-mode.md new file mode 100644 index 0000000000..06a593c6d0 --- /dev/null +++ b/.changeset/chat-hydrate-strict-mode.md @@ -0,0 +1,5 @@ +--- +'@tanstack/ai-client': patch +--- + +Send one hydrate `GET` when a chat with `persistence: true` mounts in React Strict Mode. Strict Mode attaches, detaches and attaches the client again in dev, and each attach sent its own `GET`. A re-attach now reuses the `GET` that is still in flight. diff --git a/packages/ai-client/src/chat-client.ts b/packages/ai-client/src/chat-client.ts index 4a89f6952d..92b696c9d0 100644 --- a/packages/ai-client/src/chat-client.ts +++ b/packages/ai-client/src/chat-client.ts @@ -540,6 +540,8 @@ export class ChatClient< private hydrationError: Error | undefined /** Whether a view is currently watching. See `attach` / `detach`. */ private tailing = false + /** `historyGeneration` of the mount hydration GET still in flight, if any. */ + private hydrationInFlight: number | undefined /** Constructor inputs `attach()` needs on every re-attach, not just the first. */ private readonly rejoinRunId: string | null | undefined private readonly cachesMessages: boolean @@ -1148,12 +1150,17 @@ export class ChatClient< if (!hydrate) return if (this.isLoading || this.abortController) return if (this.disposed) return + // A re-attach while the last GET is still out reuses that GET. React Strict + // Mode runs mount effects twice (attach, detach, attach), and the GET + // applies its result because a view is attached again when it returns. + if (this.hydrationInFlight === this.historyGeneration) return const pageSize = this.historyPageSize const hydrateOptions: ChatHydrateOptions | undefined = pageSize === undefined ? undefined : { limit: pageSize } void (async () => { let result: ChatHydrationResult const generation = this.historyGeneration + this.hydrationInFlight = generation try { result = await hydrate(this.threadId, hydrateOptions) } catch (cause) { @@ -1161,6 +1168,10 @@ export class ChatClient< // older attempt must not touch the state of a newer one. if (generation === this.historyGeneration) this.failHydration(cause) return + } finally { + if (this.hydrationInFlight === generation) { + this.hydrationInFlight = undefined + } } if (generation !== this.historyGeneration) return // NO VIEW IS WATCHING ANY MORE (it unmounted while this fetch was in diff --git a/packages/ai-client/tests/dispose-tail-leak.test.ts b/packages/ai-client/tests/dispose-tail-leak.test.ts index 0a444e6bda..5e77adb126 100644 --- a/packages/ai-client/tests/dispose-tail-leak.test.ts +++ b/packages/ai-client/tests/dispose-tail-leak.test.ts @@ -338,6 +338,48 @@ describe('dispose stops a late hydration from opening a tail', () => { client.dispose() }) + it('a re-attach while hydration is in flight reuses that GET', async () => { + // React Strict Mode runs mount effects twice in dev: attach, detach, attach, + // all before the first hydrate GET returns. That must stay one GET, and its + // result must still paint the re-attached view. + let releaseHydration: (() => void) | undefined + const gate = new Promise((resolve) => { + releaseHydration = resolve + }) + const hydrate = vi.fn(async () => { + await gate + return { + messages: [ + { + id: 'u1', + role: 'user' as const, + parts: [{ type: 'text' as const, content: 'hi' }], + }, + ], + activeRun: null, + interrupts: null, + } + }) + const connection: ResumableConnectConnectionAdapter = { + connect: async function* () {}, + joinRun: async function* () {}, + hydrate, + } + + const client = mountedChatClient({ + threadId: 't1', + connection, + persistence: true, + }) + client.detach() + client.attach() + releaseHydration?.() + + await vi.waitFor(() => expect(client.getMessages()).toHaveLength(1)) + expect(hydrate).toHaveBeenCalledTimes(1) + client.dispose() + }) + it('still joins the run when hydration resolves before dispose', async () => { // The guard must not break the ordinary path it protects. const joinRun = vi.fn(async function* () { diff --git a/testing/e2e/src/routeTree.gen.ts b/testing/e2e/src/routeTree.gen.ts index ffefe68803..a6dc8d8f7c 100644 --- a/testing/e2e/src/routeTree.gen.ts +++ b/testing/e2e/src/routeTree.gen.ts @@ -13,6 +13,7 @@ import { Route as IndexRouteImport } from './routes/index' import { Route as ByokRouteImport } from './routes/byok' import { Route as ChatClientDefaultBridgeRouteImport } from './routes/chat-client-default-bridge' import { Route as ChatClientStreamProcessingRouteImport } from './routes/chat-client-stream-processing' +import { Route as ClientMountHydrateRouteImport } from './routes/client-mount-hydrate' import { Route as DevtoolsChatRouteImport } from './routes/devtools-chat' import { Route as DevtoolsGenerationHooksRouteImport } from './routes/devtools-generation-hooks' import { Route as DevtoolsMemoryRouteImport } from './routes/devtools-memory' @@ -157,6 +158,11 @@ const ChatClientStreamProcessingRoute = path: '/chat-client-stream-processing', getParentRoute: () => rootRouteImport, } as any) +const ClientMountHydrateRoute = ClientMountHydrateRouteImport.update({ + id: '/client-mount-hydrate', + path: '/client-mount-hydrate', + getParentRoute: () => rootRouteImport, +} as any) const DevtoolsChatRoute = DevtoolsChatRouteImport.update({ id: '/devtools-chat', path: '/devtools-chat', @@ -801,6 +807,7 @@ export interface FileRoutesByFullPath { '/byok': typeof ByokRoute '/chat-client-default-bridge': typeof ChatClientDefaultBridgeRoute '/chat-client-stream-processing': typeof ChatClientStreamProcessingRoute + '/client-mount-hydrate': typeof ClientMountHydrateRoute '/devtools-chat': typeof DevtoolsChatRoute '/devtools-generation-hooks': typeof DevtoolsGenerationHooksRoute '/devtools-memory': typeof DevtoolsMemoryRoute @@ -929,6 +936,7 @@ export interface FileRoutesByTo { '/byok': typeof ByokRoute '/chat-client-default-bridge': typeof ChatClientDefaultBridgeRoute '/chat-client-stream-processing': typeof ChatClientStreamProcessingRoute + '/client-mount-hydrate': typeof ClientMountHydrateRoute '/devtools-chat': typeof DevtoolsChatRoute '/devtools-generation-hooks': typeof DevtoolsGenerationHooksRoute '/devtools-memory': typeof DevtoolsMemoryRoute @@ -1058,6 +1066,7 @@ export interface FileRoutesById { '/byok': typeof ByokRoute '/chat-client-default-bridge': typeof ChatClientDefaultBridgeRoute '/chat-client-stream-processing': typeof ChatClientStreamProcessingRoute + '/client-mount-hydrate': typeof ClientMountHydrateRoute '/devtools-chat': typeof DevtoolsChatRoute '/devtools-generation-hooks': typeof DevtoolsGenerationHooksRoute '/devtools-memory': typeof DevtoolsMemoryRoute @@ -1188,6 +1197,7 @@ export interface FileRouteTypes { | '/byok' | '/chat-client-default-bridge' | '/chat-client-stream-processing' + | '/client-mount-hydrate' | '/devtools-chat' | '/devtools-generation-hooks' | '/devtools-memory' @@ -1316,6 +1326,7 @@ export interface FileRouteTypes { | '/byok' | '/chat-client-default-bridge' | '/chat-client-stream-processing' + | '/client-mount-hydrate' | '/devtools-chat' | '/devtools-generation-hooks' | '/devtools-memory' @@ -1444,6 +1455,7 @@ export interface FileRouteTypes { | '/byok' | '/chat-client-default-bridge' | '/chat-client-stream-processing' + | '/client-mount-hydrate' | '/devtools-chat' | '/devtools-generation-hooks' | '/devtools-memory' @@ -1573,6 +1585,7 @@ export interface RootRouteChildren { ByokRoute: typeof ByokRoute ChatClientDefaultBridgeRoute: typeof ChatClientDefaultBridgeRoute ChatClientStreamProcessingRoute: typeof ChatClientStreamProcessingRoute + ClientMountHydrateRoute: typeof ClientMountHydrateRoute DevtoolsChatRoute: typeof DevtoolsChatRoute DevtoolsGenerationHooksRoute: typeof DevtoolsGenerationHooksRoute DevtoolsMemoryRoute: typeof DevtoolsMemoryRoute @@ -1722,6 +1735,13 @@ declare module '@tanstack/react-router' { preLoaderRoute: typeof ChatClientStreamProcessingRouteImport parentRoute: typeof rootRouteImport } + '/client-mount-hydrate': { + id: '/client-mount-hydrate' + path: '/client-mount-hydrate' + fullPath: '/client-mount-hydrate' + preLoaderRoute: typeof ClientMountHydrateRouteImport + parentRoute: typeof rootRouteImport + } '/devtools-chat': { id: '/devtools-chat' path: '/devtools-chat' @@ -2642,6 +2662,7 @@ const rootRouteChildren: RootRouteChildren = { ByokRoute: ByokRoute, ChatClientDefaultBridgeRoute: ChatClientDefaultBridgeRoute, ChatClientStreamProcessingRoute: ChatClientStreamProcessingRoute, + ClientMountHydrateRoute: ClientMountHydrateRoute, DevtoolsChatRoute: DevtoolsChatRoute, DevtoolsGenerationHooksRoute: DevtoolsGenerationHooksRoute, DevtoolsMemoryRoute: DevtoolsMemoryRoute, diff --git a/testing/e2e/src/routes/client-mount-hydrate.tsx b/testing/e2e/src/routes/client-mount-hydrate.tsx new file mode 100644 index 0000000000..63c19a43f8 --- /dev/null +++ b/testing/e2e/src/routes/client-mount-hydrate.tsx @@ -0,0 +1,41 @@ +import { useEffect, useState } from 'react' +import { createFileRoute } from '@tanstack/react-router' +import { fetchServerSentEvents, useChat } from '@tanstack/ai-react' + +/** + * A server-authoritative chat (`persistence: true`) that mounts on the client, + * after hydration, like a chat that waits for a thread id from localStorage. + * + * React Strict Mode replays the effects of a client mount (attach, detach, + * attach). A component hydrated from the server render does not get that + * replay, so `/persistence-durability` cannot show it. The GET answers with + * the `server-interrupt` scenario's pending interrupt. + */ +const connection = fetchServerSentEvents( + '/api/persistence-durability?scenario=server-interrupt', +) + +export const Route = createFileRoute('/client-mount-hydrate')({ + component: ClientMountHydratePage, +}) + +function ClientMountHydratePage() { + const [mounted, setMounted] = useState(false) + useEffect(() => setMounted(true), []) + return ( +
+ {mounted ? : null} +
+ ) +} + +function Chat() { + const { interrupts } = useChat({ + threadId: 'client-mount-hydrate', + connection, + persistence: true, + }) + return ( +
+ ) +} diff --git a/testing/e2e/tests/persistence-durability.spec.ts b/testing/e2e/tests/persistence-durability.spec.ts index 7ee0c84691..438c0858ea 100644 --- a/testing/e2e/tests/persistence-durability.spec.ts +++ b/testing/e2e/tests/persistence-durability.spec.ts @@ -143,6 +143,35 @@ test.describe('persistence durability (browser refresh)', () => { ) expect(stored).toBeNull() }) + + test('sends one hydrate GET when a chat mounts under React Strict Mode (persistence: true)', async ({ + page, + }) => { + // The dev server replays a client mount's effects (attach, detach, attach) + // before the first GET returns. Each attach used to send its own hydrate + // GET. The page mounts the chat after hydration, because a hydrated mount + // gets no replay. + const hydrateGets: Array = [] + page.on('request', (req) => { + const url = new URL(req.url()) + if ( + req.method() === 'GET' && + url.pathname === '/api/persistence-durability' && + url.searchParams.has('threadId') + ) { + hydrateGets.push(req.url()) + } + }) + + await page.goto('/client-mount-hydrate') + + // The interrupt comes from the hydrate response, so every hydrate GET of + // this mount was already sent when it shows. + await expect + .poll(() => interruptCount(page), { timeout: 15_000 }) + .toBeGreaterThanOrEqual(1) + expect(hydrateGets).toHaveLength(1) + }) }) test.describe('structured output persistence', () => { From f2dad7dbeaae6c7eb9d6206cbb61f1b9594accf2 Mon Sep 17 00:00:00 2001 From: Alem Tuzlak Date: Mon, 5 Oct 2026 12:23:33 +0200 Subject: [PATCH 3/4] docs(examples): add a durable persistence demo page `/durable-persistence` in ts-react-chat runs `persistence: true` with a durable POST. Switches turn on the old batching and a slow model, and a counter shows how text chunks arrive. Reload mid-answer to see the run continue. Claude-Session: https://claude.ai/code/session_01APYv1qshKyjPPpkFyRZhfZ --- .../ts-react-chat/src/components/Header.tsx | 13 + examples/ts-react-chat/src/routeTree.gen.ts | 42 +++ .../src/routes/api.durable-persistence.ts | 82 ++++++ .../src/routes/durable-persistence.tsx | 239 ++++++++++++++++++ 4 files changed, 376 insertions(+) create mode 100644 examples/ts-react-chat/src/routes/api.durable-persistence.ts create mode 100644 examples/ts-react-chat/src/routes/durable-persistence.tsx diff --git a/examples/ts-react-chat/src/components/Header.tsx b/examples/ts-react-chat/src/components/Header.tsx index 20d344a520..c703e7f96e 100644 --- a/examples/ts-react-chat/src/components/Header.tsx +++ b/examples/ts-react-chat/src/components/Header.tsx @@ -408,6 +408,19 @@ export default function Header() { Persistent Chat + setIsOpen(false)} + className="flex items-center gap-3 p-3 rounded-lg hover:bg-gray-800 transition-colors mb-2" + activeProps={{ + className: + 'flex items-center gap-3 p-3 rounded-lg bg-cyan-600 hover:bg-cyan-700 transition-colors mb-2', + }} + > + + Durable Persistence + + setIsOpen(false)} diff --git a/examples/ts-react-chat/src/routeTree.gen.ts b/examples/ts-react-chat/src/routeTree.gen.ts index 5fb8b067d2..895cc91f9e 100644 --- a/examples/ts-react-chat/src/routeTree.gen.ts +++ b/examples/ts-react-chat/src/routeTree.gen.ts @@ -13,6 +13,7 @@ import { Route as IndexRouteImport } from './routes/index' import { Route as AppStudioRouteImport } from './routes/app-studio' import { Route as CapabilityDemoRouteImport } from './routes/capability-demo' import { Route as CompactionRouteImport } from './routes/compaction' +import { Route as DurablePersistenceRouteImport } from './routes/durable-persistence' import { Route as GenerationHooksRouteImport } from './routes/generation-hooks' import { Route as GenericInterruptsRouteImport } from './routes/generic-interrupts' import { Route as ImageGenRouteImport } from './routes/image-gen' @@ -38,6 +39,7 @@ import { Route as ApiAppStudioForkRouteImport } from './routes/api.app-studio-fo import { Route as ApiArtifactsRouteImport } from './routes/api.artifacts' import { Route as ApiCapabilityDemoRouteImport } from './routes/api.capability-demo' import { Route as ApiCompactionRouteImport } from './routes/api.compaction' +import { Route as ApiDurablePersistenceRouteImport } from './routes/api.durable-persistence' import { Route as ApiGenericInterruptsRouteImport } from './routes/api.generic-interrupts' import { Route as ApiImageGenRouteImport } from './routes/api.image-gen' import { Route as ApiImageToolReproRouteImport } from './routes/api.image-tool-repro' @@ -104,6 +106,11 @@ const CompactionRoute = CompactionRouteImport.update({ path: '/compaction', getParentRoute: () => rootRouteImport, } as any) +const DurablePersistenceRoute = DurablePersistenceRouteImport.update({ + id: '/durable-persistence', + path: '/durable-persistence', + getParentRoute: () => rootRouteImport, +} as any) const GenerationHooksRoute = GenerationHooksRouteImport.update({ id: '/generation-hooks', path: '/generation-hooks', @@ -229,6 +236,11 @@ const ApiCompactionRoute = ApiCompactionRouteImport.update({ path: '/api/compaction', getParentRoute: () => rootRouteImport, } as any) +const ApiDurablePersistenceRoute = ApiDurablePersistenceRouteImport.update({ + id: '/api/durable-persistence', + path: '/api/durable-persistence', + getParentRoute: () => rootRouteImport, +} as any) const ApiGenericInterruptsRoute = ApiGenericInterruptsRouteImport.update({ id: '/api/generic-interrupts', path: '/api/generic-interrupts', @@ -466,6 +478,7 @@ export interface FileRoutesByFullPath { '/app-studio': typeof AppStudioRoute '/capability-demo': typeof CapabilityDemoRoute '/compaction': typeof CompactionRoute + '/durable-persistence': typeof DurablePersistenceRoute '/generation-hooks': typeof GenerationHooksRoute '/generic-interrupts': typeof GenericInterruptsRoute '/image-gen': typeof ImageGenRoute @@ -491,6 +504,7 @@ export interface FileRoutesByFullPath { '/api/artifacts': typeof ApiArtifactsRoute '/api/capability-demo': typeof ApiCapabilityDemoRoute '/api/compaction': typeof ApiCompactionRoute + '/api/durable-persistence': typeof ApiDurablePersistenceRoute '/api/generic-interrupts': typeof ApiGenericInterruptsRoute '/api/image-gen': typeof ApiImageGenRoute '/api/image-tool-repro': typeof ApiImageToolReproRoute @@ -542,6 +556,7 @@ export interface FileRoutesByTo { '/app-studio': typeof AppStudioRoute '/capability-demo': typeof CapabilityDemoRoute '/compaction': typeof CompactionRoute + '/durable-persistence': typeof DurablePersistenceRoute '/generation-hooks': typeof GenerationHooksRoute '/generic-interrupts': typeof GenericInterruptsRoute '/image-gen': typeof ImageGenRoute @@ -567,6 +582,7 @@ export interface FileRoutesByTo { '/api/artifacts': typeof ApiArtifactsRoute '/api/capability-demo': typeof ApiCapabilityDemoRoute '/api/compaction': typeof ApiCompactionRoute + '/api/durable-persistence': typeof ApiDurablePersistenceRoute '/api/generic-interrupts': typeof ApiGenericInterruptsRoute '/api/image-gen': typeof ApiImageGenRoute '/api/image-tool-repro': typeof ApiImageToolReproRoute @@ -619,6 +635,7 @@ export interface FileRoutesById { '/app-studio': typeof AppStudioRoute '/capability-demo': typeof CapabilityDemoRoute '/compaction': typeof CompactionRoute + '/durable-persistence': typeof DurablePersistenceRoute '/generation-hooks': typeof GenerationHooksRoute '/generic-interrupts': typeof GenericInterruptsRoute '/image-gen': typeof ImageGenRoute @@ -644,6 +661,7 @@ export interface FileRoutesById { '/api/artifacts': typeof ApiArtifactsRoute '/api/capability-demo': typeof ApiCapabilityDemoRoute '/api/compaction': typeof ApiCompactionRoute + '/api/durable-persistence': typeof ApiDurablePersistenceRoute '/api/generic-interrupts': typeof ApiGenericInterruptsRoute '/api/image-gen': typeof ApiImageGenRoute '/api/image-tool-repro': typeof ApiImageToolReproRoute @@ -697,6 +715,7 @@ export interface FileRouteTypes { | '/app-studio' | '/capability-demo' | '/compaction' + | '/durable-persistence' | '/generation-hooks' | '/generic-interrupts' | '/image-gen' @@ -722,6 +741,7 @@ export interface FileRouteTypes { | '/api/artifacts' | '/api/capability-demo' | '/api/compaction' + | '/api/durable-persistence' | '/api/generic-interrupts' | '/api/image-gen' | '/api/image-tool-repro' @@ -773,6 +793,7 @@ export interface FileRouteTypes { | '/app-studio' | '/capability-demo' | '/compaction' + | '/durable-persistence' | '/generation-hooks' | '/generic-interrupts' | '/image-gen' @@ -798,6 +819,7 @@ export interface FileRouteTypes { | '/api/artifacts' | '/api/capability-demo' | '/api/compaction' + | '/api/durable-persistence' | '/api/generic-interrupts' | '/api/image-gen' | '/api/image-tool-repro' @@ -849,6 +871,7 @@ export interface FileRouteTypes { | '/app-studio' | '/capability-demo' | '/compaction' + | '/durable-persistence' | '/generation-hooks' | '/generic-interrupts' | '/image-gen' @@ -874,6 +897,7 @@ export interface FileRouteTypes { | '/api/artifacts' | '/api/capability-demo' | '/api/compaction' + | '/api/durable-persistence' | '/api/generic-interrupts' | '/api/image-gen' | '/api/image-tool-repro' @@ -926,6 +950,7 @@ export interface RootRouteChildren { AppStudioRoute: typeof AppStudioRoute CapabilityDemoRoute: typeof CapabilityDemoRoute CompactionRoute: typeof CompactionRoute + DurablePersistenceRoute: typeof DurablePersistenceRoute GenerationHooksRoute: typeof GenerationHooksRoute GenericInterruptsRoute: typeof GenericInterruptsRoute ImageGenRoute: typeof ImageGenRoute @@ -951,6 +976,7 @@ export interface RootRouteChildren { ApiArtifactsRoute: typeof ApiArtifactsRoute ApiCapabilityDemoRoute: typeof ApiCapabilityDemoRoute ApiCompactionRoute: typeof ApiCompactionRoute + ApiDurablePersistenceRoute: typeof ApiDurablePersistenceRoute ApiGenericInterruptsRoute: typeof ApiGenericInterruptsRoute ApiImageGenRoute: typeof ApiImageGenRoute ApiImageToolReproRoute: typeof ApiImageToolReproRoute @@ -1027,6 +1053,13 @@ declare module '@tanstack/react-router' { preLoaderRoute: typeof CompactionRouteImport parentRoute: typeof rootRouteImport } + '/durable-persistence': { + id: '/durable-persistence' + path: '/durable-persistence' + fullPath: '/durable-persistence' + preLoaderRoute: typeof DurablePersistenceRouteImport + parentRoute: typeof rootRouteImport + } '/generation-hooks': { id: '/generation-hooks' path: '/generation-hooks' @@ -1202,6 +1235,13 @@ declare module '@tanstack/react-router' { preLoaderRoute: typeof ApiCompactionRouteImport parentRoute: typeof rootRouteImport } + '/api/durable-persistence': { + id: '/api/durable-persistence' + path: '/api/durable-persistence' + fullPath: '/api/durable-persistence' + preLoaderRoute: typeof ApiDurablePersistenceRouteImport + parentRoute: typeof rootRouteImport + } '/api/generic-interrupts': { id: '/api/generic-interrupts' path: '/api/generic-interrupts' @@ -1536,6 +1576,7 @@ const rootRouteChildren: RootRouteChildren = { AppStudioRoute: AppStudioRoute, CapabilityDemoRoute: CapabilityDemoRoute, CompactionRoute: CompactionRoute, + DurablePersistenceRoute: DurablePersistenceRoute, GenerationHooksRoute: GenerationHooksRoute, GenericInterruptsRoute: GenericInterruptsRoute, ImageGenRoute: ImageGenRoute, @@ -1561,6 +1602,7 @@ const rootRouteChildren: RootRouteChildren = { ApiArtifactsRoute: ApiArtifactsRoute, ApiCapabilityDemoRoute: ApiCapabilityDemoRoute, ApiCompactionRoute: ApiCompactionRoute, + ApiDurablePersistenceRoute: ApiDurablePersistenceRoute, ApiGenericInterruptsRoute: ApiGenericInterruptsRoute, ApiImageGenRoute: ApiImageGenRoute, ApiImageToolReproRoute: ApiImageToolReproRoute, diff --git a/examples/ts-react-chat/src/routes/api.durable-persistence.ts b/examples/ts-react-chat/src/routes/api.durable-persistence.ts new file mode 100644 index 0000000000..5f941d4396 --- /dev/null +++ b/examples/ts-react-chat/src/routes/api.durable-persistence.ts @@ -0,0 +1,82 @@ +import { createFileRoute } from '@tanstack/react-router' +import { + chat, + chatParamsFromRequestBody, + memoryStream, + resumeServerSentEventsResponse, + toServerSentEventsResponse, +} from '@tanstack/ai' +import { openaiText } from '@tanstack/ai-openai' +import { + memoryPersistence, + reconstructChat, + withPersistence, +} from '@tanstack/ai-persistence' +import type { StreamChunk } from '@tanstack/ai' + +// In memory, so this page never touches the SQLite threads of /persistent-chat. +// A server restart clears it. +const persistence = memoryPersistence() + +// "Slow model" on the page puts 30ms between text chunks, like a slower model, +// so the batching is easy to see by eye. +async function* slowText( + stream: AsyncIterable, +): AsyncIterable { + for await (const chunk of stream) { + if (chunk.type === 'TEXT_MESSAGE_CONTENT') { + await new Promise((resolve) => setTimeout(resolve, 30)) + } + yield chunk + } +} + +/** + * Server-authoritative persistence with a resumable POST: the persistence + * guide's POST plus `durability`. A reload mid-answer rejoins the run from the + * memoryStream log, because the GET reports the thread's active run. + */ +export const Route = createFileRoute('/api/durable-persistence')({ + server: { + handlers: { + POST: async ({ request }) => { + const params = await chatParamsFromRequestBody(await request.json()) + const stream = chat({ + adapter: openaiText('gpt-5.5'), + messages: params.messages, + threadId: params.threadId, + runId: params.runId, + ...(params.resume ? { resume: params.resume } : {}), + middleware: [withPersistence(persistence)], + }) + // "Old batching" on the page turns the batch wait off. A chunk then + // waits until 32 chunks or the run end, which is how durable batching + // worked before `batchWaitMs`. + const oldBatching = params.forwardedProps.oldBatching === true + const slowModel = params.forwardedProps.slowModel === true + return toServerSentEventsResponse( + slowModel ? slowText(stream) : stream, + { + durability: { + adapter: memoryStream(request), + ...(oldBatching ? { batchWaitMs: 2_147_483_647 } : {}), + }, + }, + ) + }, + + GET: ({ request }) => { + // A rejoin (`?offset` / `Last-Event-ID`) replays the run's log. + const durability = memoryStream(request) + if (durability.resumeFrom() !== null) { + return resumeServerSentEventsResponse({ adapter: durability }) + } + // Otherwise return the stored thread plus a cursor to any active run. + // Demo only: a real app checks that the session owns the thread. + return reconstructChat(persistence, request, { + authorize: async (threadId) => threadId.length > 0, + }) + }, + }, + }, +}) diff --git a/examples/ts-react-chat/src/routes/durable-persistence.tsx b/examples/ts-react-chat/src/routes/durable-persistence.tsx new file mode 100644 index 0000000000..989df3a834 --- /dev/null +++ b/examples/ts-react-chat/src/routes/durable-persistence.tsx @@ -0,0 +1,239 @@ +import { createFileRoute } from '@tanstack/react-router' +import { useEffect, useMemo, useRef, useState } from 'react' +import { fetchServerSentEvents } from '@tanstack/ai-client' +import { useChat } from '@tanstack/ai-react' + +export const Route = createFileRoute('/durable-persistence')({ + component: DurablePersistencePage, +}) + +/** + * A fetch that reports how many text chunks each network read of an SSE + * response carries. It counts before the client parses anything, so the + * numbers show what the server sent and when, not when React rendered it. + */ +function countingFetch(onRead: (textChunks: number) => void): typeof fetch { + return async (input, init) => { + const response = await fetch(input, init) + const isSse = response.headers + .get('content-type') + ?.includes('text/event-stream') + if (!response.body || !isSse) return response + const decoder = new TextDecoder() + let pending = '' + const counted = response.body.pipeThrough( + new TransformStream({ + transform(bytes, controller) { + pending += decoder.decode(bytes, { stream: true }) + const events = pending.split('\n\n') + pending = events.pop() ?? '' + const textChunks = events.filter((event) => + event.includes('"type":"TEXT_MESSAGE_CONTENT"'), + ).length + if (textChunks > 0) onRead(textChunks) + controller.enqueue(bytes) + }, + }), + ) + return new Response(counted, { + status: response.status, + statusText: response.statusText, + headers: response.headers, + }) + } +} + +// The thread id lives in localStorage, so a reload opens the same thread. The +// transcript itself lives on the server. +const THREAD_KEY = 'durable-persistence:thread' + +function newThreadId(): string { + const id = `thread-${crypto.randomUUID()}` + window.localStorage.setItem(THREAD_KEY, id) + return id +} + +function DurablePersistencePage() { + // Read on the client only: the server render has no localStorage. + const [threadId, setThreadId] = useState(null) + useEffect(() => { + setThreadId(window.localStorage.getItem(THREAD_KEY) ?? newThreadId()) + }, []) + + return ( +
+

Durable persistence

+

+ persistence: true on the client. The server POST is the + persistence guide's POST plus{' '} + durability: {'{ adapter: memoryStream(request) }'}. Send a + long prompt and watch the counter. Reload while it streams: the reply + continues where it was. After a reload, the part you missed arrives as + one step, then the rest streams. +

+ + {threadId ? : null} +
+ ) +} + +interface ChunkStats { + chunks: number + steps: number + biggestStep: number +} + +const NO_STATS: ChunkStats = { chunks: 0, steps: 0, biggestStep: 0 } + +function ChatPane({ threadId }: { threadId: string }) { + const [oldBatching, setOldBatching] = useState(false) + const [slowModel, setSlowModel] = useState(true) + const [stats, setStats] = useState(NO_STATS) + // Network reads within 5ms of each other count as one step on screen. + const tracker = useRef({ lastAt: 0, currentStep: 0 }) + + const connection = useMemo( + () => + fetchServerSentEvents('/api/durable-persistence', { + fetchClient: countingFetch((textChunks) => { + const now = performance.now() + const t = tracker.current + const newStep = now - t.lastAt > 5 + t.currentStep = (newStep ? 0 : t.currentStep) + textChunks + t.lastAt = now + const step = t.currentStep + setStats((s) => ({ + chunks: s.chunks + textChunks, + steps: s.steps + (newStep ? 1 : 0), + biggestStep: Math.max(s.biggestStep, step), + })) + }), + }), + [], + ) + + const { messages, sendMessage, isLoading, connectionStatus } = useChat({ + threadId, + connection, + persistence: true, + body: { oldBatching, slowModel }, + }) + const [input, setInput] = useState( + 'Write a 700-word story about a lighthouse keeper.', + ) + + const handleSubmit = (e: React.FormEvent) => { + e.preventDefault() + const text = input.trim() + if (!text || isLoading) return + setInput('') + setStats(NO_STATS) + tracker.current = { lastAt: 0, currentStep: 0 } + void sendMessage(text) + } + + return ( + <> + + +
+ thread: {threadId} | connection:{' '} + {connectionStatus} +
+ text chunks: {stats.chunks} | arrived in{' '} + {stats.steps} steps | biggest step:{' '} + {stats.biggestStep} chunks +
+ +
+ {messages.map((message) => ( +
+
{message.role}
+ {message.parts.map((part, index) => + part.type === 'text' && part.content ? ( +

+ {part.content} +

+ ) : null, + )} +
+ ))} +
+ +
+ setInput(e.target.value)} + placeholder="Ask something..." + style={{ flex: 1, padding: 8 }} + /> + +
+ + ) +} + +const page: React.CSSProperties = { + maxWidth: 720, + margin: '0 auto', + padding: 24, + fontFamily: 'system-ui, sans-serif', +} + +const bubble: React.CSSProperties = { + borderRadius: 8, + padding: '10px 14px', + maxWidth: '85%', +} + +const userBubble: React.CSSProperties = { + ...bubble, + alignSelf: 'flex-end', + background: '#eef2ff', +} + +const assistantBubble: React.CSSProperties = { + ...bubble, + alignSelf: 'flex-start', + background: '#f6f6f6', +} + +const roleLabel: React.CSSProperties = { + fontSize: 11, + textTransform: 'uppercase', + letterSpacing: '0.05em', + color: '#999', + marginBottom: 4, +} From d3ba6525688551f3791bf0285517417caf6aa9a5 Mon Sep 17 00:00:00 2001 From: Alem Tuzlak Date: Mon, 5 Oct 2026 13:08:04 +0200 Subject: [PATCH 4/4] fix(ai): end the run when a timed durability flush fails A timed flush runs while the next chunk is still pending. If its append failed, the producer awaited `return()` on the source, and an async generator runs that only after its pending pull. So RUN_ERROR and the log close waited for the model's next chunk, maybe forever. After a failed timed flush, close the source in the background. The failure path then persists RUN_ERROR and ends the response at once. Claude-Session: https://claude.ai/code/session_01APYv1qshKyjPPpkFyRZhfZ --- packages/ai/src/stream-to-response.ts | 26 +++++++++-- .../stream-to-response-durability.test.ts | 43 +++++++++++++++++++ 2 files changed, 66 insertions(+), 3 deletions(-) diff --git a/packages/ai/src/stream-to-response.ts b/packages/ai/src/stream-to-response.ts index cf95394fbb..1911383021 100644 --- a/packages/ai/src/stream-to-response.ts +++ b/packages/ai/src/stream-to-response.ts @@ -565,9 +565,11 @@ export function durableStreamSource( // Iterated by hand, not with `for await`, so the batch can flush while // the next chunk is still pending. The `finally` does what `for await` // would: close the producer on an early exit, but not after it finished - // or threw. + // or threw. After a failed timed flush it closes it in the background. const iterator = stream[Symbol.asyncIterator]() let iteratorFinished = false + // A timed flush failed while `next` was still pending. + let pullPending = false let flushBy = 0 try { for (;;) { @@ -576,7 +578,12 @@ export function durableStreamSource( batch.length > 0 && !(await settlesWithin(next, flushBy - Date.now())) ) { - yield* flush() + try { + yield* flush() + } catch (error) { + pullPending = true + throw error + } } let result: IteratorResult try { @@ -600,7 +607,20 @@ export function durableStreamSource( } } } finally { - if (!iteratorFinished) await iterator.return?.() + if (pullPending) { + // An async generator runs `return()` only after its pending pull, + // so awaiting it here waits for the model's next chunk, maybe + // forever. Close it in the background: the failure path below must + // persist RUN_ERROR now, and a failure is never a detach, so nothing + // below needs the producer's `onAbort` chain to finish first. + Promise.resolve(iterator.return?.()).catch((error: unknown) => { + logger?.errors('closing the producer after a failed flush failed', { + error, + }) + }) + } else if (!iteratorFinished) { + await iterator.return?.() + } } if (!isAborted(abortController.signal)) yield* flush() } catch (error) { diff --git a/packages/ai/tests/stream-to-response-durability.test.ts b/packages/ai/tests/stream-to-response-durability.test.ts index bfb4b4dedf..93ea8771f7 100644 --- a/packages/ai/tests/stream-to-response-durability.test.ts +++ b/packages/ai/tests/stream-to-response-durability.test.ts @@ -602,6 +602,49 @@ function runFinished(): AdapterYieldChunk { } describe('durability producer robustness', () => { + it('ends with RUN_ERROR when a timed flush fails while the model is still busy', async () => { + const log = memoryStream( + new Request('https://example.test/api/chat?runId=timed-flush-fails', { + method: 'POST', + }), + ) + let appends = 0 + const durability: StreamDurability = { + ...log, + // Append 1 is the run-accepted marker. Append 2 is the timed flush of + // 'a', made while the model has not sent its next chunk. + append: (chunks) => { + appends += 1 + return appends === 2 + ? Promise.reject(new Error('log write failed')) + : log.append(chunks) + }, + } + const stream: AsyncIterable = { + async *[Symbol.asyncIterator]() { + yield ev.textContent('a') + // The model never sends another chunk. + await new Promise(() => undefined) + }, + } + + const outcome = await Promise.race([ + readBody( + toServerSentEventsResponse(stream, { + durability: { adapter: durability }, + }), + ), + new Promise<'hung'>((resolve) => setTimeout(() => resolve('hung'), 1000)), + ]) + if (outcome === 'hung') throw new Error('The response never ended') + + const events = parseSseEvents(outcome) + expect(field(events.at(-1)!, 'type')).toBe(EventType.RUN_ERROR) + const logged: Array = [] + for await (const { chunk } of log.read('-1')) logged.push(chunk) + expect(logged.at(-1)?.type).toBe(EventType.RUN_ERROR) + }) + it('flushes buffered chunks to the durable log before the terminal on abort', async () => { const durability = memoryStream( new Request('https://example.test/api/chat?runId=h1-abort', {