From cb29dcf567becb34580298b12db084c7dd062b0d Mon Sep 17 00:00:00 2001 From: Alem Tuzlak Date: Mon, 21 Sep 2026 19:10:11 +0200 Subject: [PATCH 1/2] feat(ai): flush CUSTOM events immediately through durability CUSTOM progress events now flush as soon as they are emitted, so compaction:started and tool progress reach the client at emit time. High-volume names (process.stdout, process.stderr, sandbox.file, sandbox.file.diff) still batch. Pass { batch: true } on emitCustomEvent to opt an event into the batch. --- .changeset/custom-event-immediate-flush.md | 6 + docs/advanced/compaction.md | 7 +- docs/advanced/middleware.md | 4 +- docs/api/ai.md | 6 +- docs/protocol/custom-events.md | 10 + .../interfaces/ChatMiddlewareContext.md | 9 +- .../type-aliases/ToolExecutionContext.md | 10 +- docs/tools/tools.md | 4 + packages/ai-sandbox/src/bridge-events.ts | 11 +- packages/ai-sandbox/src/tool-bridge.ts | 8 +- packages/ai/src/activities/chat/index.ts | 16 +- .../src/activities/chat/middleware/types.ts | 10 +- .../src/activities/chat/tools/tool-calls.ts | 20 +- packages/ai/src/stream-to-response.ts | 22 +- packages/ai/src/strip-to-spec-middleware.ts | 3 +- packages/ai/src/types.ts | 24 +- packages/ai/src/utilities/durability-batch.ts | 48 ++++ .../stream-to-response-durability.test.ts | 236 +++++++++++++++++- 18 files changed, 419 insertions(+), 35 deletions(-) create mode 100644 .changeset/custom-event-immediate-flush.md create mode 100644 packages/ai/src/utilities/durability-batch.ts diff --git a/.changeset/custom-event-immediate-flush.md b/.changeset/custom-event-immediate-flush.md new file mode 100644 index 0000000000..02d5d24cae --- /dev/null +++ b/.changeset/custom-event-immediate-flush.md @@ -0,0 +1,6 @@ +--- +'@tanstack/ai': minor +'@tanstack/ai-sandbox': patch +--- + +Flush CUSTOM events through the durability layer as soon as they are emitted, so progress events such as `compaction:started` reach the client at emit time. High-volume events (`process.stdout`, `process.stderr`, `sandbox.file`, `sandbox.file.diff`) still batch. Pass `{ batch: true }` on `emitCustomEvent` to opt an event into the batch. diff --git a/docs/advanced/compaction.md b/docs/advanced/compaction.md index 732d599b22..5b9ab3593b 100644 --- a/docs/advanced/compaction.md +++ b/docs/advanced/compaction.md @@ -197,9 +197,10 @@ The token count is a rough `characters / 4` estimate. It is good enough to trigg After a compaction, the chat stream includes three CUSTOM events in order: `compaction:started`, `compaction:state`, then `compaction:ended`. -`compaction:started` is sent before the strategy runs, so a slow -`summarizeOldest` call still shows up as started on the client. The state and -ended events follow when the strategy returns. +`compaction:started` is sent before the strategy runs. Durability flushes each +event as soon as it is emitted, so a slow `summarizeOldest` call still shows +as started on the client while it runs. The state and ended events follow when +the strategy returns. TanStack AI DevTools has a Compaction tab on the hook. Each compact shows: - started, state, and ended rows diff --git a/docs/advanced/middleware.md b/docs/advanced/middleware.md index 7c68f5cff4..2cfde67d77 100644 --- a/docs/advanced/middleware.md +++ b/docs/advanced/middleware.md @@ -692,7 +692,7 @@ Every hook receives a `ChatMiddlewareContext` as its first argument. It provides | `chunkIndex` | `number` | Running count of chunks yielded | | `signal` | `AbortSignal \| undefined` | External abort signal | | `abort(reason?)` | `function` | Abort the run from within middleware | -| `emitCustomEvent(name, value)` | `function` | Push a `CUSTOM` chunk onto the chat stream now. The engine yields it while the current hook is still running, including during `onConfig`. | +| `emitCustomEvent(name, value, options?)` | `function` | Push a `CUSTOM` chunk onto the chat stream now. The engine yields it while the current hook is still running, including during `onConfig`. Durability flushes it immediately. Pass `{ batch: true }` to keep it in the durability batch. | | `context` | `TContext` | User-provided runtime context value | | `defer(promise)` | `function` | Register a non-blocking side-effect | @@ -999,7 +999,7 @@ const progress: ChatMiddleware = { }; ``` -The engine yields each `CUSTOM` chunk as soon as you call `emitCustomEvent`. If `RUN_STARTED` is not on the wire yet, the engine sends it first. Read these events on the client the same way as tool `emitCustomEvent` calls. See [Custom Events](../protocol/custom-events). +The engine yields each `CUSTOM` chunk as soon as you call `emitCustomEvent`. If `RUN_STARTED` is not on the wire yet, the engine sends it first. Durability then flushes the event so the client can render a live indicator. Read these events on the client the same way as tool `emitCustomEvent` calls. See [Custom Events](../protocol/custom-events). ### Rate Limiting diff --git a/docs/api/ai.md b/docs/api/ai.md index 58d04546aa..faec748a17 100644 --- a/docs/api/ai.md +++ b/docs/api/ai.md @@ -632,7 +632,11 @@ interface Tool { ```typescript ignore type ToolExecutionContext = { toolCallId?: string; - emitCustomEvent: (eventName: string, value: Record) => void; + emitCustomEvent: ( + eventName: string, + value: Record, + options?: { batch?: boolean }, + ) => void; } & (unknown extends TContext ? { context?: TContext } : { context: TContext }); ``` diff --git a/docs/protocol/custom-events.md b/docs/protocol/custom-events.md index 437d80ad59..7fd19aac63 100644 --- a/docs/protocol/custom-events.md +++ b/docs/protocol/custom-events.md @@ -162,6 +162,16 @@ const progress: ChatMiddleware = { }; ``` +With durability on the response, each of these events flushes as soon as it +is emitted, so `prepare` reaches the client while `prepare()` is still +running. High-volume names stay in the durability batch: +`process.stdout`, `process.stderr`, `sandbox.file`, and `sandbox.file.diff`. +Pass `{ batch: true }` to keep one of your own events in that batch: + +```ts +ctx.emitCustomEvent("my-app:progress", { step: "prepare" }, { batch: true }); +``` + These flow over the wire exactly like the built-in events: same `CUSTOM` chunk shape, same runtime behavior. But `'my-app:progress'` isn't one of the literal names in `KnownCustomEvent`, so it's intentionally absent from diff --git a/docs/reference/interfaces/ChatMiddlewareContext.md b/docs/reference/interfaces/ChatMiddlewareContext.md index 1525fd41bd..5baa6c88b3 100644 --- a/docs/reference/interfaces/ChatMiddlewareContext.md +++ b/docs/reference/interfaces/ChatMiddlewareContext.md @@ -181,14 +181,15 @@ after the terminal hook (onFinish/onAbort/onError). ### emitCustomEvent ```ts -emitCustomEvent: (name, value) => void; +emitCustomEvent: (name, value, options?) => void; ``` Defined in: [packages/ai/src/activities/chat/middleware/types.ts:221](https://github.com/TanStack/ai/blob/main/packages/ai/src/activities/chat/middleware/types.ts#L221) Push a `CUSTOM` chunk onto the chat stream immediately. The engine yields it as soon as it can (including while `onConfig` -is still awaiting work such as a summarize call). +is still awaiting work such as a summarize call). Durability then +flushes the event on its own, unless you pass `{ batch: true }`. #### Parameters @@ -200,6 +201,10 @@ is still awaiting work such as a summarize call). `Record`\<`string`, `any`\> +##### options? + +`EmitCustomEventOptions` + #### Returns `void` diff --git a/docs/reference/type-aliases/ToolExecutionContext.md b/docs/reference/type-aliases/ToolExecutionContext.md index 512949faa3..cf6fdf8e91 100644 --- a/docs/reference/type-aliases/ToolExecutionContext.md +++ b/docs/reference/type-aliases/ToolExecutionContext.md @@ -27,11 +27,13 @@ e.g. MCP `callTool` — should forward this to cancel in-flight work. ### emitCustomEvent ```ts -emitCustomEvent: (eventName, value) => void; +emitCustomEvent: (eventName, value, options?) => void; ``` Emit a custom event during tool execution. Events are streamed to the client in real-time as AG-UI CUSTOM events. +Durability flushes each event immediately. Pass `{ batch: true }` to keep +the event in the durability batch. #### Parameters @@ -47,6 +49,12 @@ Name of the custom event Event payload value +##### options? + +`EmitCustomEventOptions` + +Pass `{ batch: true }` to keep this event in the durability batch + #### Returns `void` diff --git a/docs/tools/tools.md b/docs/tools/tools.md index 7dacae4a2f..68f0a59d02 100644 --- a/docs/tools/tools.md +++ b/docs/tools/tools.md @@ -431,6 +431,10 @@ const importData = importDataDef.server(async (input, { context, }); ``` +Each `emitCustomEvent` call flushes through durability immediately, so the +client can show progress while the tool still runs. Pass `{ batch: true }` +only for a high-volume stream. See [Custom Events](../protocol/custom-events). + See [Server Tools](./server-tools) for the full runtime-context pattern. ## Tool States diff --git a/packages/ai-sandbox/src/bridge-events.ts b/packages/ai-sandbox/src/bridge-events.ts index 1dbdaf5910..903f5dff3c 100644 --- a/packages/ai-sandbox/src/bridge-events.ts +++ b/packages/ai-sandbox/src/bridge-events.ts @@ -11,11 +11,15 @@ * runs (e.g. code mode's `code_mode:console` logs during a long execution). */ import { EventType, withTanstackMetadata } from '@tanstack/ai' -import type { StreamChunk } from '@tanstack/ai' +import type { EmitCustomEventOptions, StreamChunk } from '@tanstack/ai' export interface BridgeEventChannel { /** Pass as the bridge's `emitCustomEvent`; buffers a CUSTOM chunk for the stream. */ - emitCustomEvent: (eventName: string, value: Record) => void + emitCustomEvent: ( + eventName: string, + value: Record, + options?: EmitCustomEventOptions, + ) => void /** Live CUSTOM-chunk stream; ends after {@link close} once drained. */ stream: AsyncIterable /** Stop the stream (call when the run's main output is done). */ @@ -48,7 +52,7 @@ export function createBridgeEventChannel(meta: { } return { - emitCustomEvent(eventName, value) { + emitCustomEvent(eventName, value, options) { if (closed) return buffer.push( withTanstackMetadata( @@ -62,6 +66,7 @@ export function createBridgeEventChannel(meta: { model: meta.model, ...(meta.threadId !== undefined ? { threadId: meta.threadId } : {}), ...(meta.runId !== undefined ? { runId: meta.runId } : {}), + ...(options?.batch === true ? { batch: true } : {}), }, ) as StreamChunk, ) diff --git a/packages/ai-sandbox/src/tool-bridge.ts b/packages/ai-sandbox/src/tool-bridge.ts index 04d28e3eb1..5d7a117a31 100644 --- a/packages/ai-sandbox/src/tool-bridge.ts +++ b/packages/ai-sandbox/src/tool-bridge.ts @@ -28,7 +28,7 @@ import { ListToolsRequestSchema, } from '@modelcontextprotocol/sdk/types.js' import type { AddressInfo } from 'node:net' -import type { AnyTool } from '@tanstack/ai' +import type { AnyTool, EmitCustomEventOptions } from '@tanstack/ai' /** * Name of the bridged MCP server. The agent sees tools as @@ -72,7 +72,11 @@ export interface ToolBridgeCoreOptions { * `emitCustomEvent` never reaches a bridged tool. The harness adapter supplies * one that injects a CUSTOM chunk into its live output stream. */ - emitCustomEvent?: (eventName: string, value: Record) => void + emitCustomEvent?: ( + eventName: string, + value: Record, + options?: EmitCustomEventOptions, + ) => void /** * Optional permission-prompt tool (e.g. for Claude Code's * `--permission-prompt-tool`). When set, the bridge exposes an extra MCP tool diff --git a/packages/ai/src/activities/chat/index.ts b/packages/ai/src/activities/chat/index.ts index 4d2f8c4424..07302084f6 100644 --- a/packages/ai/src/activities/chat/index.ts +++ b/packages/ai/src/activities/chat/index.ts @@ -38,6 +38,7 @@ import { tanstackMetadata, withTanstackMetadata, } from '../../utilities/merge-metadata' +import { withDurabilityBatchHint } from '../../utilities/durability-batch' import { normalizeStreamChunk } from '../../utilities/normalize-stream-chunk' import { restorePublicUsage } from '../../utilities/restore-inbound-chunk' import type { AdapterYieldChunk } from '../../utilities/adapter-yield-chunk' @@ -101,6 +102,7 @@ import type { ChatStream, ConstrainedModelMessage, CustomEvent, + EmitCustomEventOptions, InferSchemaType, Interrupt, JSONSchema, @@ -1000,9 +1002,9 @@ class TextEngine< this.abortReason = reason this.middlewareAbortController?.abort(reason) }, - emitCustomEvent: (name, value) => { + emitCustomEvent: (name, value, options) => { this.middlewareCustomQueue.push( - this.createCustomEventChunk(name, value), + this.createCustomEventChunk(name, value, options), ) const waiters = this.middlewareCustomWaiters this.middlewareCustomWaiters = [] @@ -2019,7 +2021,8 @@ class TextEngine< this.resolveExecutableTools(executablePendingCalls), approvals, clientToolResults, - (eventName, data) => this.createCustomEventChunk(eventName, data), + (eventName, data, options) => + this.createCustomEventChunk(eventName, data, options), { onBeforeToolCall: async (toolCall, tool, args) => { this.logger.tools(`phase=before name=${toolCall.function.name}`, { @@ -2199,7 +2202,8 @@ class TextEngine< this.resolveExecutableTools(executableToolCalls), approvals, clientToolResults, - (eventName, data) => this.createCustomEventChunk(eventName, data), + (eventName, data, options) => + this.createCustomEventChunk(eventName, data, options), { onBeforeToolCall: async (toolCall, tool, args) => { this.logger.tools(`phase=before name=${toolCall.function.name}`, { @@ -4583,13 +4587,15 @@ class TextEngine< private createCustomEventChunk( eventName: string, value: Record, + options?: EmitCustomEventOptions, ): CustomEvent { - return { + const chunk: CustomEvent = { type: EventType.CUSTOM, timestamp: Date.now(), name: eventName, value, } + return options?.batch ? withDurabilityBatchHint(chunk) : chunk } private createId(prefix: string): string { diff --git a/packages/ai/src/activities/chat/middleware/types.ts b/packages/ai/src/activities/chat/middleware/types.ts index 933c4ea271..7dd48e9ff8 100644 --- a/packages/ai/src/activities/chat/middleware/types.ts +++ b/packages/ai/src/activities/chat/middleware/types.ts @@ -4,6 +4,7 @@ import type { } from '@standard-schema/spec' import type { AgentLoopState, + EmitCustomEventOptions, JSONSchema, ModelMessage, RunAgentResumeItem, @@ -216,9 +217,14 @@ export interface ChatMiddlewareContext { /** * Push a `CUSTOM` chunk onto the chat stream immediately. * The engine yields it as soon as it can (including while `onConfig` - * is still awaiting work such as a summarize call). + * is still awaiting work such as a summarize call). Durability then + * flushes the event on its own, unless you pass `{ batch: true }`. */ - emitCustomEvent: (name: string, value: Record) => void + emitCustomEvent: ( + name: string, + value: Record, + options?: EmitCustomEventOptions, + ) => void /** Runtime context provided by chat() options */ context: TContext /** diff --git a/packages/ai/src/activities/chat/tools/tool-calls.ts b/packages/ai/src/activities/chat/tools/tool-calls.ts index 4d5cf811f6..fc8aca0e67 100644 --- a/packages/ai/src/activities/chat/tools/tool-calls.ts +++ b/packages/ai/src/activities/chat/tools/tool-calls.ts @@ -7,6 +7,7 @@ import type { AnyTool, ContentPart, CustomEvent, + EmitCustomEventOptions, ModelMessage, RunFinishedEvent, Tool, @@ -760,6 +761,7 @@ export async function* executeToolCalls( createCustomEventChunk?: ( eventName: string, value: Record, + options?: EmitCustomEventOptions, ) => CustomEvent, middlewareHooks?: ToolExecutionMiddlewareHooks, userContext?: TContext, @@ -874,13 +876,21 @@ export async function* executeToolCalls( toolCallId: toolCall.id, context: userContext, abortSignal, - emitCustomEvent: (eventName: string, value: Record) => { + emitCustomEvent: ( + eventName: string, + value: Record, + options?: EmitCustomEventOptions, + ) => { if (createCustomEventChunk) { pendingEvents.push( - createCustomEventChunk(eventName, { - ...value, - toolCallId: toolCall.id, - }), + createCustomEventChunk( + eventName, + { + ...value, + toolCallId: toolCall.id, + }, + options, + ), ) } }, diff --git a/packages/ai/src/stream-to-response.ts b/packages/ai/src/stream-to-response.ts index 3227bc07d0..40858c3063 100644 --- a/packages/ai/src/stream-to-response.ts +++ b/packages/ai/src/stream-to-response.ts @@ -9,6 +9,10 @@ import { notifyRunDisconnected } from './delivery-disconnect' import { resolveResumeRunId } from './stream-durability' import { EventType } from './types' import { toWireChunk } from './strip-to-spec-middleware' +import { + isDurabilityBatchedCustom, + stripDurabilityBatchHint, +} from './utilities/durability-batch' import { resolveDebugOption } from './logger/resolve' import { runErrorEventToError } from './utilities/errors' import type { LockStore } from './activities/chat/middleware/locks' @@ -319,23 +323,29 @@ function resolveBatchSize(batch: number | undefined): number { /** * Boundaries at which the batching producer flushes early, regardless of the - * batch size — the run-start marker, terminal events, and tool-call ends. - * Flushing here keeps the durability log promptly consistent at semantically - * meaningful points. + * batch size: run-start, terminals, tool-call ends, and CUSTOM events that + * are not high-volume adapter output. * * `RUN_STARTED` matters especially for one-shot activities (image, speech, * transcription, summarize): they emit `RUN_STARTED`, then await the provider * for seconds, then a terminal. Without flushing `RUN_STARTED` the log stays * empty for the whole run, so a mount-time `joinRun` finds nothing and its - * empty-log deadline fast-fails as "run gone" — even though the run is alive. + * empty-log deadline fast-fails as "run gone" even though the run is alive. * Flushing it immediately makes the run resumable from the instant it starts. + * + * CUSTOM progress events (compaction, tool progress, middleware) flush at + * emit time so a live indicator can render. `process.stdout`, + * `process.stderr`, `sandbox.file`, and `sandbox.file.diff` stay batched. + * `emitCustomEvent(name, value, { batch: true })` opts a single event into + * that same batch. */ function isDurabilityFlushBoundary(chunk: StreamChunk): boolean { return ( chunk.type === 'RUN_STARTED' || chunk.type === 'RUN_FINISHED' || chunk.type === 'RUN_ERROR' || - chunk.type === 'TOOL_CALL_END' + chunk.type === 'TOOL_CALL_END' || + (chunk.type === 'CUSTOM' && !isDurabilityBatchedCustom(chunk)) ) } @@ -446,7 +456,7 @@ export function durableStreamSource( async function* flush(): AsyncIterable { if (batch.length === 0) return - const toForward = batch + const toForward = batch.map(stripDurabilityBatchHint) batch = [] // Tag each chunk with the exact backend offset. Requiring one opaque // token per chunk preserves exact-once resume at any batch size. diff --git a/packages/ai/src/strip-to-spec-middleware.ts b/packages/ai/src/strip-to-spec-middleware.ts index 833ce5bbf7..f62e5be09f 100644 --- a/packages/ai/src/strip-to-spec-middleware.ts +++ b/packages/ai/src/strip-to-spec-middleware.ts @@ -6,6 +6,7 @@ import { tanstackMetadata, withTanstackMetadata, } from './utilities/merge-metadata' +import { stripDurabilityBatchHint } from './utilities/durability-batch' import { normalizeStreamChunk } from './utilities/normalize-stream-chunk' import { isSpecTopLevelKey } from './utilities/spec-event-keys' @@ -54,5 +55,5 @@ export function toWireChunk( chunk: StreamChunk | AdapterYieldChunk, ): StreamChunk { const [normalized] = normalizeStreamChunk(chunk) - return stripToSpec(normalized ?? chunk) + return stripDurabilityBatchHint(stripToSpec(normalized ?? chunk)) } diff --git a/packages/ai/src/types.ts b/packages/ai/src/types.ts index 81db7f37d4..07d3693e96 100644 --- a/packages/ai/src/types.ts +++ b/packages/ai/src/types.ts @@ -629,6 +629,22 @@ type RuntimeContextField = context: TContext } +/** + * Options for a single `emitCustomEvent` call, on both the tool-execution and + * middleware contexts. + */ +export interface EmitCustomEventOptions { + /** + * Keep this event in the durability batch with later chunks. + * CUSTOM events flush as soon as they are emitted, so a progress + * indicator can render at emit time. Pass `{ batch: true }` for a + * high-volume stream that should share appends with later output. + * `process.stdout`, `process.stderr`, `sandbox.file`, and + * `sandbox.file.diff` already batch. + */ + batch?: boolean +} + /** * Context passed to tool execute functions, providing capabilities like * emitting custom events during execution. @@ -649,6 +665,8 @@ export type ToolExecutionContext = * * @param eventName - Name of the custom event * @param value - Event payload value + * @param options - Pass `{ batch: true }` to keep this event in the + * durability batch instead of flushing it immediately * * @example * ```ts @@ -661,7 +679,11 @@ export type ToolExecutionContext = * }) * ``` */ - emitCustomEvent: (eventName: string, value: Record) => void + emitCustomEvent: ( + eventName: string, + value: Record, + options?: EmitCustomEventOptions, + ) => void } export type ToolExecuteFunction< diff --git a/packages/ai/src/utilities/durability-batch.ts b/packages/ai/src/utilities/durability-batch.ts new file mode 100644 index 0000000000..49f1e88e8d --- /dev/null +++ b/packages/ai/src/utilities/durability-batch.ts @@ -0,0 +1,48 @@ +import { CUSTOM_EVENT } from '../custom-events' +import type { CustomEvent, StreamChunk } from '../types' +import { tanstackMetadata, withTanstackMetadata } from './merge-metadata' + +/** + * High-volume CUSTOM names that stay in the durability batch. + * Everything else flushes as soon as it is emitted. + */ +const BATCHED_CUSTOM_EVENT_NAMES = new Set([ + CUSTOM_EVENT.PROCESS_STDOUT, + CUSTOM_EVENT.PROCESS_STDERR, + 'sandbox.file', + 'sandbox.file.diff', +]) + +function hasBatchHint(chunk: StreamChunk): boolean { + const tanstack = tanstackMetadata(chunk) + if (tanstack == null) return false + return Reflect.get(tanstack, 'batch') === true +} + +/** Mark a CUSTOM chunk so the durability producer keeps it in the batch. */ +export function withDurabilityBatchHint(chunk: CustomEvent): CustomEvent { + return withTanstackMetadata(chunk, { batch: true }) +} + +export function isDurabilityBatchedCustom(chunk: StreamChunk): boolean { + if (chunk.type !== 'CUSTOM') return false + if (BATCHED_CUSTOM_EVENT_NAMES.has(chunk.name)) return true + return hasBatchHint(chunk) +} + +/** + * Drop the in-process batch hint so it does not sit in the log or on the wire. + */ +export function stripDurabilityBatchHint(chunk: StreamChunk): StreamChunk { + const tanstack = tanstackMetadata(chunk) + if (tanstack == null || Reflect.get(tanstack, 'batch') !== true) return chunk + Reflect.deleteProperty(tanstack, 'batch') + const metadata = chunk.metadata + if (metadata != null && Object.keys(tanstack).length === 0) { + Reflect.deleteProperty(metadata, 'tanstack') + if (Object.keys(metadata).length === 0) { + Reflect.deleteProperty(chunk, 'metadata') + } + } + return chunk +} diff --git a/packages/ai/tests/stream-to-response-durability.test.ts b/packages/ai/tests/stream-to-response-durability.test.ts index af00697935..93f166468e 100644 --- a/packages/ai/tests/stream-to-response-durability.test.ts +++ b/packages/ai/tests/stream-to-response-durability.test.ts @@ -1,4 +1,5 @@ import { describe, expect, it, vi } from 'vitest' +import { z } from 'zod' import { memoryStream } from '../src/stream-durability' import { RUN_ACCEPTED_EVENT, @@ -9,9 +10,14 @@ import { } from '../src/stream-to-response' import { EventType } from '../src/types' import { chat } from '../src/activities/chat/index' +import { toolDefinition } from '../src/activities/chat/tools/tool-definition' +import { CUSTOM_EVENT } from '../src/custom-events' +import { withDurabilityBatchHint } from '../src/utilities/durability-batch' +import { tanstackMetadata } from '../src/utilities/merge-metadata' import { createMockAdapter, ev } from './test-utils' +import type { ChatMiddleware } from '../src/activities/chat/middleware/types' import type { StreamDurability } from '../src/stream-durability' -import type { StreamChunk } from '../src/types' +import type { CustomEvent, StreamChunk } from '../src/types' import type { AdapterYieldChunk } from '../src/utilities/adapter-yield-chunk' function fiveChunkStream(): { @@ -834,3 +840,231 @@ describe('a failing middleware onFinish surfaces past the transport', () => { ) }) }) + +function customThenText(custom: StreamChunk): AsyncIterable { + return { + async *[Symbol.asyncIterator]() { + yield custom + yield ev.textStart() + yield ev.textContent('a') + yield ev.textContent('b') + yield ev.textEnd() + yield ev.runFinished('stop') + }, + } +} + +function customEvent(name: string): CustomEvent { + return { + type: EventType.CUSTOM, + name, + value: {}, + timestamp: Date.now(), + } +} + +const hasModelText = (chunks: Array): boolean => + chunks.some((chunk) => chunk.type === EventType.TEXT_MESSAGE_CONTENT) + +function namedCustom(name: string): (chunk: StreamChunk) => boolean { + return (chunk) => chunk.type === EventType.CUSTOM && chunk.name === name +} + +describe('CUSTOM events flush through durability at emit time', () => { + it('flushes a progress CUSTOM event on its own, before model text', async () => { + const durability = memoryStream( + new Request('https://example.test/api/chat?runId=custom-flush', { + method: 'POST', + }), + ) + const appendSpy = vi.spyOn(durability, 'append') + + await readBody( + toServerSentEventsResponse(customThenText(customEvent('work:started')), { + durability: { adapter: durability }, + }), + ) + + const batches = appendSpy.mock.calls.map(([chunks]) => chunks) + const startedBatch = batches.find((chunks) => + chunks.some(namedCustom('work:started')), + ) + if (startedBatch === undefined) { + throw new Error('work:started was never appended') + } + expect(startedBatch).toHaveLength(1) + expect(hasModelText(startedBatch)).toBe(false) + }) + + it('keeps process.stdout in the batch with later text', async () => { + const durability = memoryStream( + new Request('https://example.test/api/chat?runId=stdout-batch', { + method: 'POST', + }), + ) + const appendSpy = vi.spyOn(durability, 'append') + + await readBody( + toServerSentEventsResponse( + customThenText(customEvent(CUSTOM_EVENT.PROCESS_STDOUT)), + { durability: { adapter: durability } }, + ), + ) + + const batches = appendSpy.mock.calls.map(([chunks]) => chunks) + const stdoutBatch = batches.find((chunks) => + chunks.some(namedCustom(CUSTOM_EVENT.PROCESS_STDOUT)), + ) + if (stdoutBatch === undefined) { + throw new Error('process.stdout was never appended') + } + expect(hasModelText(stdoutBatch)).toBe(true) + }) + + it('keeps emitCustomEvent({ batch: true }) in the batch with later text', async () => { + const durability = memoryStream( + new Request('https://example.test/api/chat?runId=opt-in-batch', { + method: 'POST', + }), + ) + const appendSpy = vi.spyOn(durability, 'append') + const marked = withDurabilityBatchHint(customEvent('work:started')) + + await readBody( + toServerSentEventsResponse(customThenText(marked), { + durability: { adapter: durability }, + }), + ) + + const batches = appendSpy.mock.calls.map(([chunks]) => chunks) + const startedBatch = batches.find((chunks) => + chunks.some(namedCustom('work:started')), + ) + if (startedBatch === undefined) { + throw new Error('work:started was never appended') + } + expect(hasModelText(startedBatch)).toBe(true) + expect( + startedBatch.some((chunk) => { + const tanstack = tanstackMetadata(chunk) + return tanstack != null && 'batch' in tanstack + }), + ).toBe(false) + }) + + it('appends a middleware CUSTOM event while onConfig is still waiting', async () => { + let release = () => {} + const gate = new Promise((resolve) => { + release = resolve + }) + const durability = memoryStream( + new Request('https://example.test/api/chat?runId=live-middleware', { + method: 'POST', + }), + ) + const appendSpy = vi.spyOn(durability, 'append') + const { adapter } = createMockAdapter({ + iterations: [ + [ + ev.runStarted(), + ev.textStart(), + ev.textContent('reply'), + ev.textEnd(), + ev.runFinished('stop'), + ], + ], + }) + const middleware: ChatMiddleware = { + name: 'slow-prepare', + async onConfig(ctx) { + if (ctx.phase !== 'beforeModel') return + ctx.emitCustomEvent('work:started', { step: 'prepare' }) + await gate + ctx.emitCustomEvent('work:ended', { step: 'prepare' }) + }, + } + + const body = readBody( + toServerSentEventsResponse( + chat({ + adapter, + messages: [{ role: 'user', content: 'go' }], + middleware: [middleware], + }), + { durability: { adapter: durability } }, + ), + ) + + await vi.waitFor(() => { + const appended = appendSpy.mock.calls.some(([chunks]) => + chunks.some(namedCustom('work:started')), + ) + expect(appended).toBe(true) + }) + expect( + appendSpy.mock.calls.some(([chunks]) => + chunks.some(namedCustom('work:ended')), + ), + ).toBe(false) + + release() + await body + }) + + it('flushes each tool progress event in its own append', async () => { + const durability = memoryStream( + new Request('https://example.test/api/chat?runId=tool-progress', { + method: 'POST', + }), + ) + const appendSpy = vi.spyOn(durability, 'append') + const tool = toolDefinition({ + name: 'longJob', + description: 'Reports progress while it runs', + inputSchema: z.object({}), + }).server((_args, ctx) => { + for (const step of [1, 2, 3]) { + ctx?.emitCustomEvent('tool:progress', { step }) + } + return { done: true } + }) + const { adapter } = createMockAdapter({ + iterations: [ + [ + ev.runStarted(), + ev.toolStart('call_1', 'longJob'), + ev.toolArgs('call_1', '{}'), + ev.runFinished('tool_calls'), + ], + [ + ev.runStarted(), + ev.textStart(), + ev.textContent('done'), + ev.textEnd(), + ev.runFinished('stop'), + ], + ], + }) + + await readBody( + toServerSentEventsResponse( + chat({ + adapter, + messages: [{ role: 'user', content: 'go' }], + tools: [tool], + }), + { durability: { adapter: durability } }, + ), + ) + + const progressBatches = appendSpy.mock.calls + .map(([chunks]) => chunks) + .filter((chunks) => chunks.some(namedCustom('tool:progress'))) + expect(progressBatches).toHaveLength(3) + expect( + progressBatches.every( + (chunks) => chunks.filter(namedCustom('tool:progress')).length === 1, + ), + ).toBe(true) + }) +}) From 3ecd5ff3ebc186b12d197bc4a29402ef587a5394 Mon Sep 17 00:00:00 2001 From: Alem Tuzlak Date: Mon, 21 Sep 2026 19:45:02 +0200 Subject: [PATCH 2/2] fix(docs): give kiira a complete batch option snippet The { batch: true } example used ctx with no declaration, so root:test:kiira failed in CI. --- docs/protocol/custom-events.md | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/docs/protocol/custom-events.md b/docs/protocol/custom-events.md index 7fd19aac63..7145fbf0fe 100644 --- a/docs/protocol/custom-events.md +++ b/docs/protocol/custom-events.md @@ -169,7 +169,14 @@ running. High-volume names stay in the durability batch: Pass `{ batch: true }` to keep one of your own events in that batch: ```ts -ctx.emitCustomEvent("my-app:progress", { step: "prepare" }, { batch: true }); +import { type ChatMiddleware } from "@tanstack/ai"; + +const noisy: ChatMiddleware = { + name: "noisy", + async onConfig(ctx) { + ctx.emitCustomEvent("my-app:ticks", { n: 1 }, { batch: true }); + }, +}; ``` These flow over the wire exactly like the built-in events: same `CUSTOM`