From d7da6b9a0730c49b046a6ec96e2ebe975f587493 Mon Sep 17 00:00:00 2001 From: Jack Herrington Date: Sat, 26 Sep 2026 08:05:41 -0700 Subject: [PATCH 01/39] feat(ai-harness): add the AG-UI bridge subpath New `@tanstack/ai-harness/ag-ui` export: a thin, documented seam for AG-UI clients over the already AG-UI-native session stream. - `sessionEventsToAgUi()` normalizes a session's SessionEvent stream to pure AG-UI: run usage under `metadata.tanstack.usage`, interrupt/approval waits as `RUN_FINISHED` with `outcome.type === 'interrupt'`, subagent attribution preserved. Optional interim `tanstack.spend` CUSTOM ticks; optional dropping of harness-native CUSTOM control events. - `createAgUiHandler()` serves a session as SSE via `@ag-ui/encoder`. POST a RunAgentInput to run a prompt or resume interrupts. Strict AG-UI by default (stream begins with RUN_STARTED) so a bare `@ag-ui/client` HttpAgent consumes it directly; control detail stays on the harness protocol. Pinned to `@ag-ui/core`/`@ag-ui/encoder`/`@ag-ui/client` @ 1.0.0. Tests cover every mapped event type, the interrupt/resume round-trip, and bare-client consumption. Adds `docs/harness/ag-ui.md` and a changeset. The harness-native approval path and the relay dashboard are unchanged. Generated with Claude Code. Co-Authored-By: Claude Opus 4.8 (1M context) --- .changeset/harness-ag-ui-bridge.md | 12 + docs/config.json | 5 + docs/harness/ag-ui.md | 146 +++++++++ packages/ai-harness/package.json | 7 + packages/ai-harness/src/ag-ui.ts | 377 ++++++++++++++++++++++ packages/ai-harness/tests/ag-ui.test.ts | 407 ++++++++++++++++++++++++ packages/ai-harness/vite.config.ts | 1 + 7 files changed, 955 insertions(+) create mode 100644 .changeset/harness-ag-ui-bridge.md create mode 100644 docs/harness/ag-ui.md create mode 100644 packages/ai-harness/src/ag-ui.ts create mode 100644 packages/ai-harness/tests/ag-ui.test.ts diff --git a/.changeset/harness-ag-ui-bridge.md b/.changeset/harness-ag-ui-bridge.md new file mode 100644 index 0000000000..6e77940180 --- /dev/null +++ b/.changeset/harness-ag-ui-bridge.md @@ -0,0 +1,12 @@ +--- +'@tanstack/ai-harness': minor +--- + +New subpath `@tanstack/ai-harness/ag-ui`: the AG-UI bridge. A harness session already streams AG-UI events, so this is a thin, documented seam for AG-UI clients. + +- `sessionEventsToAgUi(events, options?)` normalizes a session's `SessionEvent` stream into a pure AG-UI event stream: run usage is surfaced under `metadata.tanstack.usage` (`normalizeUsage` handles both the AG-UI spec array and the TanStack prompt/completion shapes), interrupt/approval waits stay as `RUN_FINISHED` with `outcome.type === 'interrupt'`, and subagent attribution is preserved. It can optionally drop the harness-native `CUSTOM` control events or emit interim `tanstack.spend` `CUSTOM` ticks for a live spend meter. +- `createAgUiHandler(options)` is a `fetch` handler a bare `@ag-ui/client` `HttpAgent` can point at. `POST` a `RunAgentInput` to run a prompt, or one with `resume` entries to answer the last turn's interrupts. It streams AG-UI events as SSE via `@ag-ui/encoder`'s `EventEncoder` (protobuf framing when the client's `Accept` prefers it). The stream is strict AG-UI by default (harness control events dropped so it begins with `RUN_STARTED`); the harness-native approval path and the relay dashboard are unchanged. + +The AG-UI wire is pinned: built and tested against `@ag-ui/core@1.0.0`, `@ag-ui/encoder@1.0.0`, and `@ag-ui/client@1.0.0`. + +Generated with Claude Code. diff --git a/docs/config.json b/docs/config.json index d74f1aeba1..fbfb0a95a6 100644 --- a/docs/config.json +++ b/docs/config.json @@ -951,6 +951,11 @@ "label": "Self-host the dashboard", "to": "harness/dashboard", "addedAt": "2026-09-26" + }, + { + "label": "Stream a harness over AG-UI", + "to": "harness/ag-ui", + "addedAt": "2026-09-26" } ] }, diff --git a/docs/harness/ag-ui.md b/docs/harness/ag-ui.md new file mode 100644 index 0000000000..0239f59800 --- /dev/null +++ b/docs/harness/ag-ui.md @@ -0,0 +1,146 @@ +--- +title: Stream a harness over AG-UI +id: harness-ag-ui +order: 13 +description: "Serve a harness session as a pure AG-UI event stream a bare @ag-ui/client can consume, with approvals as interrupts and token usage as metadata." +keywords: + - tanstack ai + - harness + - AG-UI + - SSE + - interrupts + - usage +--- + +A harness session already speaks AG-UI: `session.events()` yields events whose payloads are AG-UI protocol events. The `@tanstack/ai-harness/ag-ui` subpath is the thin, documented seam an AG-UI client consumes — a normalizer and an SSE handler a bare [`@ag-ui/client`](https://docs.ag-ui.com) `HttpAgent` can point at. + +This does not replace the [harness protocol](./connect.md) or the [relay dashboard](./dashboard.md). AG-UI is the *run stream*; the harness protocol carries the *control plane* (pairing, snapshots, history, config). Use both: AG-UI for the conversation, the harness protocol for everything around it. + +## Serve a session over AG-UI + +`createAgUiHandler` is one `fetch` handler. `authorize` is required, so no endpoint is open by accident. + +```ts group=harness-ag-ui +import { defineHarness, createHarnessHost } from '@tanstack/ai-harness' +import { createAgUiHandler } from '@tanstack/ai-harness/ag-ui' +import { memoryPersistence } from '@tanstack/ai-persistence' +import { openaiText } from '@tanstack/ai-openai' + +const assistant = defineHarness({ + name: 'acme/assistant', + adapter: openaiText('gpt-5.6'), +}) + +const host = createHarnessHost({ persistence: memoryPersistence() }) + +export const handler = createAgUiHandler({ + host, + harness: assistant, + authorize: (request) => + request.headers.get('authorization') === + `Bearer ${process.env.HARNESS_TOKEN}` + ? { id: 'user-1' } + : null, + canAccess: (principal, threadId) => threadId.startsWith(principal.id), +}) +``` + +`POST` a `RunAgentInput` to it. A body with a trailing user message runs a prompt; a body with `resume` entries answers the last turn's interrupts and continues. The response is `text/event-stream` of AG-UI events, encoded with `@ag-ui/encoder` so a protobuf-accepting client gets binary framing for free. The stream closes when the run reaches a terminal state. + +## Consume it from a client + +Point a bare `@ag-ui/client` `HttpAgent` at the handler's URL: + +```ts group=harness-ag-ui-client +import { HttpAgent } from '@ag-ui/client' + +const agent = new HttpAgent({ + url: 'https://example.com/agent', + threadId: 'user-1/main', + headers: { authorization: `Bearer ${process.env.HARNESS_TOKEN}` }, +}) + +agent.addMessage({ id: 'u1', role: 'user', content: 'Summarize incidents' }) +await agent.runAgent() + +console.log(agent.messages.at(-1)) // the assistant's reply +``` + +The stream is strict AG-UI by default: it begins with `RUN_STARTED`, as a bare client requires. (The harness-native `CUSTOM` control events — `harness.operation.*`, `harness.question`, and friends — are dropped from the AG-UI stream; read them over the harness protocol's `/events` tier instead. Pass `stream: { includeHarnessEvents: true }` to keep them for a lenient consumer.) + +## Event mapping + +| Harness concept | AG-UI event | +| --- | --- | +| run start / end | `RUN_STARTED` / `RUN_FINISHED` | +| assistant text | `TEXT_MESSAGE_START` / `TEXT_MESSAGE_CONTENT` / `TEXT_MESSAGE_END` | +| tool call | `TOOL_CALL_START` / `TOOL_CALL_ARGS` / `TOOL_CALL_END` / `TOOL_CALL_RESULT` | +| approval / wait | `RUN_FINISHED` with `outcome.type === 'interrupt'` | +| resume | a new run whose `RunAgentInput.resume` answers the interrupts | +| subagent | `SUBAGENT_STARTED` / `SUBAGENT_FINISHED` (with `subagentRunId`) | +| token usage | `metadata.tanstack.usage` on `RUN_FINISHED` | +| live spend (interim) | a `tanstack.spend` `CUSTOM` event | + +## Approvals as interrupts + +A turn that stops for approval finishes with a `RUN_FINISHED` whose `outcome.type === 'interrupt'`. `@ag-ui/client` collects these on `agent.pendingInterrupts`: + +```ts group=harness-ag-ui-client +await agent.runAgent() + +for (const interrupt of agent.pendingInterrupts) { + console.log(interrupt.id, interrupt.message) // "Approval required to run remove" +} +``` + +Answer by starting a new run whose `resume` entries reference the interrupt ids: + +```ts group=harness-ag-ui-resume +await fetch('https://example.com/agent', { + method: 'POST', + headers: { + 'content-type': 'application/json', + authorization: `Bearer ${process.env.HARNESS_TOKEN}`, + }, + body: JSON.stringify({ + threadId: 'user-1/main', + runId: 'run-2', + messages: [], + tools: [], + context: [], + resume: [{ interruptId: 'approval_call_1', status: 'resolved', payload: true }], + }), +}) +``` + +The harness-native approval path (`session.resolve`, the CLI, and the relay dashboard) keeps working; AG-UI interrupt/resume is the AG-UI-shaped view of the same wait. + +## Token usage and spend + +`RUN_FINISHED` carries usage under `metadata.tanstack.usage` as a `NormalizedUsage` (`inputTokens`, `outputTokens`, `totalTokens`). Read it directly, or normalize any provider's usage with `normalizeUsage`: + +```ts group=harness-ag-ui-usage +import { normalizeUsage } from '@tanstack/ai-harness/ag-ui' + +normalizeUsage([{ inputTokens: 10, outputTokens: 5 }]) +// { inputTokens: 10, outputTokens: 5, totalTokens: 15 } +normalizeUsage({ promptTokens: 8, completionTokens: 2, totalTokens: 10 }) +// { inputTokens: 8, outputTokens: 2, totalTokens: 10 } +``` + +For a live spend meter, opt into interim ticks. After each run that reports usage, the stream carries a `tanstack.spend` `CUSTOM` event with that run's usage and a running cumulative total: + +```ts group=harness-ag-ui +const meteredHandler = createAgUiHandler({ + host, + harness: assistant, + authorize: () => ({ id: 'user-1' }), + stream: { emitSpendEvents: true }, +}) +``` + +Spend rides `CUSTOM` because the AG-UI spec has no first-class spend event yet — track that convention deliberately. + +## Pin the AG-UI version + +The AG-UI spec is still evolving. This bridge is built and tested against `@ag-ui/core@1.0.0`, `@ag-ui/encoder@1.0.0`, and `@ag-ui/client@1.0.0`. Pin those versions and move them on purpose, not by floating range. diff --git a/packages/ai-harness/package.json b/packages/ai-harness/package.json index 00f9bd8934..dee6c3da0f 100644 --- a/packages/ai-harness/package.json +++ b/packages/ai-harness/package.json @@ -49,6 +49,10 @@ "./worker": { "types": "./dist/esm/worker.d.ts", "import": "./dist/esm/worker.js" + }, + "./ag-ui": { + "types": "./dist/esm/ag-ui.d.ts", + "import": "./dist/esm/ag-ui.js" } }, "scripts": { @@ -63,9 +67,12 @@ "test:types": "tsc" }, "dependencies": { + "@ag-ui/core": "1.0.0", + "@ag-ui/encoder": "1.0.0", "@tanstack/store": "^0.11.1" }, "devDependencies": { + "@ag-ui/client": "1.0.0", "@tanstack/ai": "workspace:*", "@tanstack/ai-persistence": "workspace:*", "@vitest/coverage-v8": "4.1.10", diff --git a/packages/ai-harness/src/ag-ui.ts b/packages/ai-harness/src/ag-ui.ts new file mode 100644 index 0000000000..c7ed7752e2 --- /dev/null +++ b/packages/ai-harness/src/ag-ui.ts @@ -0,0 +1,377 @@ +/** + * `@tanstack/ai-harness/ag-ui` — the AG-UI bridge. + * + * A harness session already streams AG-UI events: `session.events()` yields + * {@link SessionEvent}s whose `event` is a `StreamChunk`, and a `StreamChunk` + * is an AG-UI protocol event (`@tanstack/ai` builds on `@ag-ui/core`). This + * module is the thin, documented seam an AG-UI client consumes: + * + * - {@link sessionEventsToAgUi} — normalize a session's `SessionEvent` stream + * into a pure AG-UI event stream (usage surfaced under a documented metadata + * key, optional live-spend `CUSTOM` events, optional dropping of the + * harness-native `CUSTOM` control events). + * - {@link createAgUiHandler} — a `fetch` handler that a bare `@ag-ui/client` + * `HttpAgent` can point at. It accepts a `RunAgentInput` (prompt or resume) + * and streams AG-UI events as SSE via `@ag-ui/encoder`'s `EventEncoder`. + * + * The AG-UI wire is pinned: this bridge is built and tested against + * `@ag-ui/core@1.0.0`, `@ag-ui/encoder@1.0.0`, and `@ag-ui/client@1.0.0`. + * Track that pin deliberately when the spec moves (see `docs/harness/ag-ui.md`). + * + * Approvals: a harness turn that stops for outside input finishes with a + * `RUN_FINISHED` whose `outcome.type === 'interrupt'`. An AG-UI client answers + * by starting a new run whose `RunAgentInput.resume` entries reference the + * interrupt ids. This does not replace the harness-native approval path (the + * CLI and the relay dashboard keep using `session.resolve` / the control tier); + * it is the AG-UI-shaped view of the same thing. + */ +import { EventEncoder } from '@ag-ui/encoder' +import { EventType, chatParamsFromRequestBody } from '@tanstack/ai' +import { HARNESS_EVENTS } from './types' +import type { BaseEvent } from '@ag-ui/core' +import type { StreamChunk } from '@tanstack/ai' +import type { AnyHarness } from './define' +import type { HarnessHost } from './host' +import type { HarnessSession } from './session' +import type { Authorize } from './http' +import type { Principal, SessionEvent } from './types' + +/** The `metadata.tanstack` namespace carries TanStack extras on AG-UI events. */ +export const TANSTACK_METADATA_NAMESPACE = 'tanstack' + +/** + * The `CUSTOM` event name for an interim live-spend tick. Emitted by + * {@link sessionEventsToAgUi} when `emitSpendEvents` is on. Its `value` is a + * {@link SpendSnapshot}. Interim by design: the AG-UI spec has no first-class + * spend event yet, so this rides `CUSTOM` (see the spec's event-mapping table). + */ +export const TANSTACK_SPEND_EVENT = 'tanstack.spend' + +/** A normalized token-usage rollup, provider-agnostic. */ +export interface NormalizedUsage { + inputTokens: number + outputTokens: number + totalTokens: number +} + +/** The payload of a {@link TANSTACK_SPEND_EVENT} `CUSTOM` event. */ +export interface SpendSnapshot { + /** The AG-UI run this tick closes. */ + runId?: string + threadId?: string + /** Usage reported by the just-finished run. */ + usage: NormalizedUsage + /** Running total across every run seen on this stream so far. */ + cumulative: NormalizedUsage +} + +const isHarnessCustom = (event: StreamChunk): boolean => + event.type === EventType.CUSTOM && + typeof (event as { name?: unknown }).name === 'string' && + (event as { name: string }).name.startsWith('harness.') + +const zeroUsage = (): NormalizedUsage => ({ + inputTokens: 0, + outputTokens: 0, + totalTokens: 0, +}) + +const isRecord = (value: unknown): value is Record => + typeof value === 'object' && value !== null + +const num = (value: unknown): number => (typeof value === 'number' ? value : 0) + +/** + * Normalize AG-UI/TanStack token usage into a single rollup. Accepts the AG-UI + * spec array (`SpecTokenUsage[]`, `inputTokens`/`outputTokens`) and the + * TanStack shape (`promptTokens`/`completionTokens`), summing arrays. + */ +export function normalizeUsage(usage: unknown): NormalizedUsage | undefined { + if (usage == null) return undefined + const entries = Array.isArray(usage) ? usage : [usage] + const total = zeroUsage() + let sawAny = false + for (const entry of entries) { + if (!isRecord(entry)) continue + sawAny = true + const input = num(entry.inputTokens) || num(entry.promptTokens) + const output = num(entry.outputTokens) || num(entry.completionTokens) + const combined = num(entry.totalTokens) || input + output + total.inputTokens += input + total.outputTokens += output + total.totalTokens += combined + } + return sawAny ? total : undefined +} + +/** Read a RUN_FINISHED event's usage, checking `usage` then `metadata`. */ +function usageOf(event: StreamChunk): NormalizedUsage | undefined { + const record = event as Record + const direct = normalizeUsage(record.usage) + if (direct) return direct + const metadata = isRecord(record.metadata) ? record.metadata : undefined + const ns = metadata && isRecord(metadata[TANSTACK_METADATA_NAMESPACE]) + ? (metadata[TANSTACK_METADATA_NAMESPACE] as Record) + : undefined + return ns ? normalizeUsage(ns.usage) : undefined +} + +/** Non-destructively set `metadata.tanstack.usage` on a RUN_FINISHED event. */ +function withUsageMetadata( + event: StreamChunk, + usage: NormalizedUsage, +): StreamChunk { + const record = event as Record + const metadata = isRecord(record.metadata) ? { ...record.metadata } : {} + const ns = isRecord(metadata[TANSTACK_METADATA_NAMESPACE]) + ? { ...(metadata[TANSTACK_METADATA_NAMESPACE] as object) } + : {} + if ((ns as Record).usage !== undefined) return event + ;(ns as Record).usage = usage + metadata[TANSTACK_METADATA_NAMESPACE] = ns + return { ...(event as object), metadata } as StreamChunk +} + +const addUsage = (a: NormalizedUsage, b: NormalizedUsage): NormalizedUsage => ({ + inputTokens: a.inputTokens + b.inputTokens, + outputTokens: a.outputTokens + b.outputTokens, + totalTokens: a.totalTokens + b.totalTokens, +}) + +/** Options for {@link sessionEventsToAgUi}. */ +export interface SessionEventsToAgUiOptions { + /** + * Keep the harness-native `CUSTOM` control events (`harness.operation.*`, + * `harness.input.*`, `harness.question`, ...). They are valid AG-UI `CUSTOM` + * events a strict AG-UI client can ignore, but they carry the control-plane + * detail the dashboard wants. Default: `true`. + */ + includeHarnessEvents?: boolean + /** + * After each `RUN_FINISHED` that reports usage, emit an interim + * {@link TANSTACK_SPEND_EVENT} `CUSTOM` event carrying that run's usage and + * the running total, for a live spend meter. Default: `false`. + */ + emitSpendEvents?: boolean +} + +/** + * Map a session's `SessionEvent` stream to a pure AG-UI event stream. + * + * The harness stream is already AG-UI-native, so this is mostly a documented + * pass-through. What it guarantees: + * + * - `RUN_FINISHED` usage is surfaced under `metadata.tanstack.usage` + * ({@link NormalizedUsage}) so consumers read one key regardless of provider. + * - Interrupt/approval waits stay as `RUN_FINISHED` with + * `outcome.type === 'interrupt'`; subagent events keep their `subagentRunId`. + * - Optionally drops the harness-native `CUSTOM` events, or emits interim + * spend ticks. + */ +export async function* sessionEventsToAgUi( + events: AsyncIterable, + options: SessionEventsToAgUiOptions = {}, +): AsyncGenerator { + const includeHarnessEvents = options.includeHarnessEvents ?? true + const emitSpendEvents = options.emitSpendEvents ?? false + let cumulative = zeroUsage() + + for await (const entry of events) { + const event = entry.event + if (!includeHarnessEvents && isHarnessCustom(event)) continue + + if (event.type === EventType.RUN_FINISHED) { + const usage = usageOf(event) + // Emit the spend tick *before* RUN_FINISHED so it stays inside the + // RUN_STARTED..RUN_FINISHED window a strict AG-UI consumer expects. + if (usage && emitSpendEvents) { + cumulative = addUsage(cumulative, usage) + const value: SpendSnapshot = { + ...(typeof (event as { runId?: unknown }).runId === 'string' + ? { runId: (event as { runId: string }).runId } + : {}), + ...(typeof (event as { threadId?: unknown }).threadId === 'string' + ? { threadId: (event as { threadId: string }).threadId } + : {}), + usage, + cumulative, + } + yield { + type: EventType.CUSTOM, + name: TANSTACK_SPEND_EVENT, + value, + timestamp: (event as { timestamp?: number }).timestamp ?? Date.now(), + } as StreamChunk + } + yield usage ? withUsageMetadata(event, usage) : event + continue + } + + yield event + } +} + +/** Session events for one operation, up to and including its terminal event. */ +async function* followOperation( + events: AsyncIterable, + operationId: string, +): AsyncGenerator { + for await (const entry of events) { + if (entry.operationId !== operationId) continue + yield entry + if ( + entry.event.type === EventType.CUSTOM && + (entry.event as { name?: string }).name === HARNESS_EVENTS.operationFinished + ) { + return + } + } +} + +/** The text of the last user message of an AG-UI request. */ +function lastUserText(messages: ReadonlyArray): string | undefined { + const message = messages.findLast( + (entry) => isRecord(entry) && entry.role === 'user', + ) + if (!isRecord(message)) return undefined + if (typeof message.content === 'string') return message.content + const parts = Array.isArray(message.content) ? message.content : [] + return parts + .map((part: unknown) => + isRecord(part) && part.type === 'text' && typeof part.text === 'string' + ? part.text + : '', + ) + .join('') +} + +const json = (body: unknown, status = 200) => + new Response(JSON.stringify(body), { + status, + headers: { 'Content-Type': 'application/json' }, + }) + +/** Options for {@link createAgUiHandler}. */ +export interface AgUiHandlerOptions { + host: HarnessHost + harness: AnyHarness + /** Decide who sends a request. Return `null` to refuse it with 401. */ + authorize: Authorize + /** May this principal use this thread? Default: yes. */ + canAccess?: ( + principal: Principal, + threadId: string, + ) => boolean | Promise + /** + * Passed through to {@link sessionEventsToAgUi} for every run. Note the + * handler defaults `includeHarnessEvents` to `false` (strict AG-UI); set it + * to `true` here to keep the harness-native `CUSTOM` control events. + */ + stream?: SessionEventsToAgUiOptions +} + +/** + * A `fetch` handler an AG-UI client (`@ag-ui/client`'s `HttpAgent`) can point + * at. `POST` a `RunAgentInput`: + * + * - with a trailing user message → runs a prompt; + * - with `resume` entries → answers the last turn's interrupts, then continues. + * + * The response is `text/event-stream` of AG-UI events, encoded with + * `@ag-ui/encoder` so protobuf-accepting clients get binary framing for free. + * The stream closes when the operation reaches a terminal state (completed, + * interrupted, failed, or cancelled). + * + * The stream is strict AG-UI by default: harness-native `CUSTOM` control + * events are dropped so it begins with `RUN_STARTED`, as a bare `@ag-ui/client` + * requires (that control detail belongs on the harness protocol's `/events` + * tier, which the dashboard consumes separately). Pass + * `stream.includeHarnessEvents: true` for a lenient consumer that wants them. + */ +export function createAgUiHandler( + options: AgUiHandlerOptions, +): (request: Request) => Promise { + const { host, harness, authorize } = options + const canAccess = options.canAccess ?? (() => true) + + const openFor = async ( + principal: Principal, + threadId: string, + ): Promise => { + if (!(await canAccess(principal, threadId))) return null + return host.open(harness, { threadId, principal }) + } + + return async (request) => { + if (request.method !== 'POST') return json({ error: 'method' }, 405) + const principal = await authorize(request) + if (!principal) return json({ error: 'unauthorized' }, 401) + + let operationId: string + let session: HarnessSession + try { + const params = await chatParamsFromRequestBody(await request.json()) + const opened = await openFor(principal, params.threadId) + if (!opened) return json({ error: 'forbidden' }, 403) + session = opened + + if (params.resume && params.resume.length > 0) { + const receipt = await session.resolve(params.resume) + if (receipt.status === 'rejected' || !receipt.operationId) { + return json({ error: receipt.reason ?? 'rejected' }, 409) + } + operationId = receipt.operationId + } else { + const message = lastUserText(params.messages) + if (message === undefined) return json({ error: 'no user message' }, 400) + operationId = session.prompt(message).id + } + } catch (error) { + // `chatParamsFromRequestBody` throws a `Response` on a malformed body. + if (error instanceof Response) return error + return json( + { error: error instanceof Error ? error.message : String(error) }, + 400, + ) + } + + // `encodeBinary` returns SSE bytes normally and protobuf framing when the + // client's `Accept` prefers it — paired with `getContentType()` below. + const encoder = new EventEncoder({ + accept: request.headers.get('accept') ?? undefined, + }) + // Not `request.signal`: some servers abort it once the body is read. The + // response aborts this controller when the client goes away. + const reader = new AbortController() + const body = new ReadableStream({ + cancel: () => reader.abort(), + async start(controller) { + try { + const entries = followOperation( + session.events({ signal: reader.signal }), + operationId, + ) + for await (const event of sessionEventsToAgUi(entries, { + includeHarnessEvents: false, + ...options.stream, + })) { + controller.enqueue(encoder.encodeBinary(event as BaseEvent)) + } + } catch (error) { + if (!reader.signal.aborted) controller.error(error) + return + } finally { + if (!reader.signal.aborted) controller.close() + } + }, + }) + + return new Response(body, { + headers: { + 'Content-Type': encoder.getContentType(), + 'Cache-Control': 'no-cache', + Connection: 'keep-alive', + }, + }) + } +} diff --git a/packages/ai-harness/tests/ag-ui.test.ts b/packages/ai-harness/tests/ag-ui.test.ts new file mode 100644 index 0000000000..3af515dd0e --- /dev/null +++ b/packages/ai-harness/tests/ag-ui.test.ts @@ -0,0 +1,407 @@ +import { describe, expect, it, vi } from 'vitest' +import { z } from 'zod' +import { HttpAgent } from '@ag-ui/client' +import { EventType, toolDefinition } from '@tanstack/ai' +import { memoryPersistence } from '@tanstack/ai-persistence' +import { createHarnessHost, defineHarness } from '../src' +import { + TANSTACK_SPEND_EVENT, + createAgUiHandler, + normalizeUsage, + sessionEventsToAgUi, +} from '../src/ag-ui' +import { mockAdapter, text, toolCall } from './helpers' +import type { SessionEvent } from '../src' +import type { StreamChunk } from '@tanstack/ai' + +/** Wrap raw AG-UI chunks in `SessionEvent`s, as `session.events()` would. */ +async function* feed( + chunks: Array, + operationId = 'op-1', +): AsyncGenerator { + let n = 0 + for (const event of chunks) { + yield { cursor: String(++n), operationId, event } + } +} + +async function collect(it: AsyncIterable): Promise> { + const out: Array = [] + for await (const value of it) out.push(value) + return out +} + +const runFinished = (extra: Record): StreamChunk => + ({ + type: EventType.RUN_FINISHED, + runId: 'r', + threadId: 't', + timestamp: 0, + ...extra, + }) as unknown as StreamChunk + +const custom = (name: string, value: unknown = {}): StreamChunk => + ({ type: EventType.CUSTOM, name, value, timestamp: 0 }) as unknown as StreamChunk + +const typesOf = (events: Array) => events.map((e) => e.type) + +describe('normalizeUsage', () => { + it('sums the AG-UI spec array shape', () => { + expect( + normalizeUsage([ + { inputTokens: 10, outputTokens: 5 }, + { inputTokens: 1, outputTokens: 2 }, + ]), + ).toEqual({ inputTokens: 11, outputTokens: 7, totalTokens: 18 }) + }) + + it('maps the TanStack prompt/completion shape', () => { + expect( + normalizeUsage({ promptTokens: 8, completionTokens: 2, totalTokens: 10 }), + ).toEqual({ inputTokens: 8, outputTokens: 2, totalTokens: 10 }) + }) + + it('returns undefined for empty or missing usage', () => { + expect(normalizeUsage(undefined)).toBeUndefined() + expect(normalizeUsage(null)).toBeUndefined() + expect(normalizeUsage([])).toBeUndefined() + }) +}) + +describe('sessionEventsToAgUi mapper', () => { + it('passes native AG-UI events through, in order', async () => { + const out = await collect(sessionEventsToAgUi(feed(text('hi')))) + expect(typesOf(out)).toEqual([ + EventType.RUN_STARTED, + EventType.TEXT_MESSAGE_START, + EventType.TEXT_MESSAGE_CONTENT, + EventType.TEXT_MESSAGE_END, + EventType.RUN_FINISHED, + ]) + }) + + it('surfaces run usage under metadata.tanstack.usage', async () => { + const [event] = await collect( + sessionEventsToAgUi( + feed([runFinished({ usage: [{ inputTokens: 10, outputTokens: 5 }] })]), + ), + ) + expect((event as any).metadata.tanstack.usage).toEqual({ + inputTokens: 10, + outputTokens: 5, + totalTokens: 15, + }) + }) + + it('does not overwrite an existing metadata.tanstack.usage', async () => { + const preset = { inputTokens: 1, outputTokens: 1, totalTokens: 2 } + const [event] = await collect( + sessionEventsToAgUi( + feed([ + runFinished({ + usage: [{ inputTokens: 99, outputTokens: 99 }], + metadata: { tanstack: { usage: preset } }, + }), + ]), + ), + ) + expect((event as any).metadata.tanstack.usage).toEqual(preset) + }) + + it('keeps interrupt outcomes intact for the AG-UI resume flow', async () => { + const [event] = await collect( + sessionEventsToAgUi( + feed([ + runFinished({ + outcome: { + type: 'interrupt', + interrupts: [ + { id: 'i1', reason: 'tool_call', toolCallId: 'call_1' }, + ], + }, + }), + ]), + ), + ) + expect((event as any).outcome.type).toBe('interrupt') + expect((event as any).outcome.interrupts[0].id).toBe('i1') + }) + + it('preserves subagent attribution on tool calls', async () => { + const toolStart = { + type: EventType.TOOL_CALL_START, + toolCallId: 'c1', + toolCallName: 'lookup', + subagentRunId: 'sub-1', + timestamp: 0, + } as unknown as StreamChunk + const [event] = await collect(sessionEventsToAgUi(feed([toolStart]))) + expect((event as any).subagentRunId).toBe('sub-1') + }) + + it('keeps harness CUSTOM events by default and drops them on request', async () => { + const chunks = [ + custom('harness.operation.started'), + custom('tool.progress'), + ...text('done'), + ] + const kept = await collect(sessionEventsToAgUi(feed(chunks))) + expect( + kept.some((e) => (e as any).name === 'harness.operation.started'), + ).toBe(true) + + const dropped = await collect( + sessionEventsToAgUi(feed(chunks), { includeHarnessEvents: false }), + ) + expect( + dropped.some((e) => (e as any).name === 'harness.operation.started'), + ).toBe(false) + // A non-harness CUSTOM event still flows through. + expect(dropped.some((e) => (e as any).name === 'tool.progress')).toBe(true) + }) + + it('emits interim spend events with a running cumulative total', async () => { + const out = await collect( + sessionEventsToAgUi( + feed([ + runFinished({ runId: 'r1', usage: [{ inputTokens: 10, outputTokens: 5 }] }), + runFinished({ runId: 'r2', usage: [{ inputTokens: 2, outputTokens: 3 }] }), + ]), + { emitSpendEvents: true }, + ), + ) + const spends = out.filter( + (e) => e.type === EventType.CUSTOM && (e as any).name === TANSTACK_SPEND_EVENT, + ) + expect(spends).toHaveLength(2) + expect((spends[0] as any).value.cumulative.totalTokens).toBe(15) + expect((spends[1] as any).value.usage.totalTokens).toBe(5) + expect((spends[1] as any).value.cumulative.totalTokens).toBe(20) + expect((spends[1] as any).value.runId).toBe('r2') + }) + + it('does not emit spend events when a run reports no usage', async () => { + const out = await collect( + sessionEventsToAgUi(feed(text('hi')), { emitSpendEvents: true }), + ) + expect( + out.some((e) => (e as any).name === TANSTACK_SPEND_EVENT), + ).toBe(false) + }) +}) + +/** Build a POST request shaped like an AG-UI `RunAgentInput`. */ +function runRequest(body: Record): Request { + return new Request('http://localhost/agent', { + method: 'POST', + headers: { + 'content-type': 'application/json', + accept: 'text/event-stream', + }, + body: JSON.stringify({ tools: [], context: [], ...body }), + }) +} + +/** Parse an SSE response body into AG-UI event objects. */ +async function readSse(res: Response): Promise> { + const raw = await res.text() + return raw + .split('\n\n') + .filter(Boolean) + .flatMap((block) => { + const line = block.split('\n').find((l) => l.startsWith('data:')) + return line ? [JSON.parse(line.slice(5).trim())] : [] + }) +} + +describe('createAgUiHandler', () => { + it('streams a prompt run as AG-UI SSE', async () => { + const host = createHarnessHost({ persistence: memoryPersistence() }) + const { adapter } = mockAdapter([() => text('hello there')]) + const harness = defineHarness({ name: 'test/agui', adapter }) + const handler = createAgUiHandler({ + host, + harness, + authorize: () => ({ id: 'u' }), + }) + + const res = await handler( + runRequest({ + threadId: 't1', + runId: 'run-1', + messages: [{ id: 'u1', role: 'user', content: 'hi' }], + }), + ) + expect(res.status).toBe(200) + expect(res.headers.get('content-type')).toBe('text/event-stream') + + const events = await readSse(res) + expect(events.some((e) => e.type === EventType.RUN_STARTED)).toBe(true) + const textOut = events + .filter((e) => e.type === EventType.TEXT_MESSAGE_CONTENT) + .map((e) => e.delta) + .join('') + expect(textOut).toContain('hello there') + expect(events.some((e) => e.type === EventType.RUN_FINISHED)).toBe(true) + await host.close() + }) + + it('rejects an unauthorized request with 401', async () => { + const host = createHarnessHost({ persistence: memoryPersistence() }) + const { adapter } = mockAdapter([() => text('x')]) + const handler = createAgUiHandler({ + host, + harness: defineHarness({ name: 'test/agui-401', adapter }), + authorize: () => null, + }) + const res = await handler( + runRequest({ threadId: 't1', runId: 'r', messages: [] }), + ) + expect(res.status).toBe(401) + await host.close() + }) + + it('returns 400 for a malformed body', async () => { + const host = createHarnessHost({ persistence: memoryPersistence() }) + const { adapter } = mockAdapter([() => text('x')]) + const handler = createAgUiHandler({ + host, + harness: defineHarness({ name: 'test/agui-400', adapter }), + authorize: () => ({ id: 'u' }), + }) + // Missing threadId/runId → chatParamsFromRequestBody throws a Response. + const res = await handler( + new Request('http://localhost/agent', { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ messages: [] }), + }), + ) + expect(res.status).toBe(400) + await host.close() + }) + + it('runs the interrupt → resume round-trip over AG-UI', async () => { + const host = createHarnessHost({ persistence: memoryPersistence() }) + const execute = vi.fn(async () => ({ ok: true })) + const remove = toolDefinition({ + name: 'remove', + description: 'Remove a file', + needsApproval: true, + inputSchema: z.object({ path: z.string() }), + }).server(execute) + const { adapter } = mockAdapter([ + () => toolCall('remove', { path: 'a.txt' }, 'call_1'), + () => text('removed'), + ]) + const harness = defineHarness({ + name: 'test/agui-approve', + adapter, + tools: [remove], + }) + const handler = createAgUiHandler({ + host, + harness, + authorize: () => ({ id: 'u' }), + }) + + const messages = [{ id: 'u1', role: 'user', content: 'remove a.txt' }] + const res1 = await handler( + runRequest({ threadId: 't1', runId: 'run-1', messages }), + ) + const events1 = await readSse(res1) + + expect( + events1.some( + (e) => + e.type === EventType.TOOL_CALL_START && e.toolCallName === 'remove', + ), + ).toBe(true) + const interrupted = events1.find( + (e) => e.type === EventType.RUN_FINISHED && e.outcome?.type === 'interrupt', + ) + expect(interrupted).toBeDefined() + expect(execute).not.toHaveBeenCalled() + const interruptId = interrupted.outcome.interrupts[0].id + + const res2 = await handler( + runRequest({ + threadId: 't1', + runId: 'run-2', + messages, + resume: [{ interruptId, status: 'resolved', payload: true }], + }), + ) + const events2 = await readSse(res2) + const textOut = events2 + .filter((e) => e.type === EventType.TEXT_MESSAGE_CONTENT) + .map((e) => e.delta) + .join('') + expect(textOut).toContain('removed') + expect(execute).toHaveBeenCalledWith({ path: 'a.txt' }, expect.anything()) + await host.close() + }) +}) + +describe('consumed by a bare @ag-ui/client HttpAgent', () => { + it('runs a prompt and surfaces the assistant reply', async () => { + const host = createHarnessHost({ persistence: memoryPersistence() }) + const { adapter } = mockAdapter([() => text('hello there')]) + const harness = defineHarness({ name: 'test/client-run', adapter }) + const handler = createAgUiHandler({ + host, + harness, + authorize: () => ({ id: 'u' }), + }) + + const agent = new HttpAgent({ + url: 'http://localhost/agent', + threadId: 't1', + fetch: (url: string, init: RequestInit) => + handler(new Request(url, init)), + }) + agent.addMessage({ id: 'u1', role: 'user', content: 'hi' }) + await agent.runAgent() + + expect(JSON.stringify(agent.messages)).toContain('hello there') + expect(agent.pendingInterrupts).toHaveLength(0) + await host.close() + }) + + it('surfaces an approval as a pending interrupt', async () => { + const host = createHarnessHost({ persistence: memoryPersistence() }) + const remove = toolDefinition({ + name: 'remove', + description: 'Remove a file', + needsApproval: true, + inputSchema: z.object({ path: z.string() }), + }).server(async () => ({ ok: true })) + const { adapter } = mockAdapter([ + () => toolCall('remove', { path: 'a.txt' }, 'call_1'), + () => text('removed'), + ]) + const harness = defineHarness({ + name: 'test/client-approve', + adapter, + tools: [remove], + }) + const handler = createAgUiHandler({ + host, + harness, + authorize: () => ({ id: 'u' }), + }) + + const agent = new HttpAgent({ + url: 'http://localhost/agent', + threadId: 't1', + fetch: (url: string, init: RequestInit) => + handler(new Request(url, init)), + }) + agent.addMessage({ id: 'u1', role: 'user', content: 'remove a.txt' }) + await agent.runAgent() + + expect(agent.pendingInterrupts).toHaveLength(1) + expect(agent.pendingInterrupts[0]?.toolCallId).toBe('call_1') + await host.close() + }) +}) diff --git a/packages/ai-harness/vite.config.ts b/packages/ai-harness/vite.config.ts index 17373da93a..440985d714 100644 --- a/packages/ai-harness/vite.config.ts +++ b/packages/ai-harness/vite.config.ts @@ -35,6 +35,7 @@ export default mergeConfig( './src/first-party/index.ts', './src/build.ts', './src/worker.ts', + './src/ag-ui.ts', ], srcDir: './src', cjs: false, From c82a5ad45bda02e2bbe40f1863220e7bed6ef3fa Mon Sep 17 00:00:00 2001 From: Jack Herrington Date: Sat, 26 Sep 2026 08:44:55 -0700 Subject: [PATCH 02/39] fix(ai-harness): coalesce multi-turn operations for strict AG-UI clients The AG-UI bridge emitted one RUN_STARTED/RUN_FINISHED per model turn, so a strict @ag-ui/client rejected the tool result between turns ("run already finished"), and it emitted RUN_FINISHED.usage as a TokenUsage object where the spec wants SpecTokenUsage[]. - `operationToAgUiRun` coalesces one operation into a single AG-UI run: one RUN_STARTED (synthesized when a resumed run leads with a tool result), one terminal RUN_FINISHED with the operation outcome and summed usage. - Conform RUN_FINISHED.usage to SpecTokenUsage[]; keep the rollup under metadata.tanstack.usage. createAgUiHandler uses the coalescer per request. Covers the multi-turn and resume round-trips end to end through a real @ag-ui/client HttpAgent. Generated with Claude Code. Co-Authored-By: Claude Opus 4.8 (1M context) --- .changeset/harness-ag-ui-bridge.md | 1 + packages/ai-harness/src/ag-ui.ts | 187 ++++++++++++++++++++---- packages/ai-harness/tests/ag-ui.test.ts | 156 +++++++++++++++++++- 3 files changed, 312 insertions(+), 32 deletions(-) diff --git a/.changeset/harness-ag-ui-bridge.md b/.changeset/harness-ag-ui-bridge.md index 6e77940180..d4c9c853bc 100644 --- a/.changeset/harness-ag-ui-bridge.md +++ b/.changeset/harness-ag-ui-bridge.md @@ -6,6 +6,7 @@ New subpath `@tanstack/ai-harness/ag-ui`: the AG-UI bridge. A harness session al - `sessionEventsToAgUi(events, options?)` normalizes a session's `SessionEvent` stream into a pure AG-UI event stream: run usage is surfaced under `metadata.tanstack.usage` (`normalizeUsage` handles both the AG-UI spec array and the TanStack prompt/completion shapes), interrupt/approval waits stay as `RUN_FINISHED` with `outcome.type === 'interrupt'`, and subagent attribution is preserved. It can optionally drop the harness-native `CUSTOM` control events or emit interim `tanstack.spend` `CUSTOM` ticks for a live spend meter. - `createAgUiHandler(options)` is a `fetch` handler a bare `@ag-ui/client` `HttpAgent` can point at. `POST` a `RunAgentInput` to run a prompt, or one with `resume` entries to answer the last turn's interrupts. It streams AG-UI events as SSE via `@ag-ui/encoder`'s `EventEncoder` (protobuf framing when the client's `Accept` prefers it). The stream is strict AG-UI by default (harness control events dropped so it begins with `RUN_STARTED`); the harness-native approval path and the relay dashboard are unchanged. +- `operationToAgUiRun(entries, options)` coalesces one harness operation — which may span several model turns (each a `RUN_STARTED`/`RUN_FINISHED` pair) — into a single valid AG-UI run: one `RUN_STARTED` (synthesized when a resumed run leads with a tool result), one terminal `RUN_FINISHED` carrying the operation outcome and summed usage, and `RUN_FINISHED.usage` conformed to the AG-UI `SpecTokenUsage[]` shape. This is what makes the stream consumable by a strict `@ag-ui/client` without "run already finished" / usage-shape errors. `createAgUiHandler` uses it per request. The AG-UI wire is pinned: built and tested against `@ag-ui/core@1.0.0`, `@ag-ui/encoder@1.0.0`, and `@ag-ui/client@1.0.0`. diff --git a/packages/ai-harness/src/ag-ui.ts b/packages/ai-harness/src/ag-ui.ts index c7ed7752e2..d9839d4079 100644 --- a/packages/ai-harness/src/ag-ui.ts +++ b/packages/ai-harness/src/ag-ui.ts @@ -116,8 +116,14 @@ function usageOf(event: StreamChunk): NormalizedUsage | undefined { return ns ? normalizeUsage(ns.usage) : undefined } -/** Non-destructively set `metadata.tanstack.usage` on a RUN_FINISHED event. */ -function withUsageMetadata( +const addUsage = (a: NormalizedUsage, b: NormalizedUsage): NormalizedUsage => ({ + inputTokens: a.inputTokens + b.inputTokens, + outputTokens: a.outputTokens + b.outputTokens, + totalTokens: a.totalTokens + b.totalTokens, +}) + +/** Set (overwriting) `metadata.tanstack.usage`, preserving other metadata. */ +function setUsageMetadata( event: StreamChunk, usage: NormalizedUsage, ): StreamChunk { @@ -126,17 +132,60 @@ function withUsageMetadata( const ns = isRecord(metadata[TANSTACK_METADATA_NAMESPACE]) ? { ...(metadata[TANSTACK_METADATA_NAMESPACE] as object) } : {} - if ((ns as Record).usage !== undefined) return event ;(ns as Record).usage = usage metadata[TANSTACK_METADATA_NAMESPACE] = ns return { ...(event as object), metadata } as StreamChunk } -const addUsage = (a: NormalizedUsage, b: NormalizedUsage): NormalizedUsage => ({ - inputTokens: a.inputTokens + b.inputTokens, - outputTokens: a.outputTokens + b.outputTokens, - totalTokens: a.totalTokens + b.totalTokens, -}) +/** + * Make a `RUN_FINISHED` event's top-level `usage` conform to AG-UI. The harness + * emits `usage` as a TanStack `TokenUsage` object, but the AG-UI spec (and a + * strict `@ag-ui/client`) require `SpecTokenUsage[]`. Set the spec array from + * the normalized total and record the rollup under `metadata.tanstack.usage`; + * when there is no usage, drop a non-array `usage` so it can't fail validation. + */ +function conformRunFinishedUsage( + event: StreamChunk, + usage: NormalizedUsage | undefined, +): StreamChunk { + if (usage) { + const withMeta = setUsageMetadata(event, usage) as Record + withMeta.usage = [ + { inputTokens: usage.inputTokens, outputTokens: usage.outputTokens }, + ] + return withMeta as StreamChunk + } + const record = event as Record + if (record.usage !== undefined && !Array.isArray(record.usage)) { + const clone = { ...record } + delete clone.usage + return clone as StreamChunk + } + return event +} + +const spendEvent = ( + from: StreamChunk, + usage: NormalizedUsage, + cumulative: NormalizedUsage, +): StreamChunk => { + const value: SpendSnapshot = { + ...(typeof (from as { runId?: unknown }).runId === 'string' + ? { runId: (from as { runId: string }).runId } + : {}), + ...(typeof (from as { threadId?: unknown }).threadId === 'string' + ? { threadId: (from as { threadId: string }).threadId } + : {}), + usage, + cumulative, + } + return { + type: EventType.CUSTOM, + name: TANSTACK_SPEND_EVENT, + value, + timestamp: (from as { timestamp?: number }).timestamp ?? Date.now(), + } as StreamChunk +} /** Options for {@link sessionEventsToAgUi}. */ export interface SessionEventsToAgUiOptions { @@ -186,24 +235,9 @@ export async function* sessionEventsToAgUi( // RUN_STARTED..RUN_FINISHED window a strict AG-UI consumer expects. if (usage && emitSpendEvents) { cumulative = addUsage(cumulative, usage) - const value: SpendSnapshot = { - ...(typeof (event as { runId?: unknown }).runId === 'string' - ? { runId: (event as { runId: string }).runId } - : {}), - ...(typeof (event as { threadId?: unknown }).threadId === 'string' - ? { threadId: (event as { threadId: string }).threadId } - : {}), - usage, - cumulative, - } - yield { - type: EventType.CUSTOM, - name: TANSTACK_SPEND_EVENT, - value, - timestamp: (event as { timestamp?: number }).timestamp ?? Date.now(), - } as StreamChunk + yield spendEvent(event, usage, cumulative) } - yield usage ? withUsageMetadata(event, usage) : event + yield conformRunFinishedUsage(event, usage) continue } @@ -211,6 +245,103 @@ export async function* sessionEventsToAgUi( } } +/** Options for {@link operationToAgUiRun}. */ +export interface OperationToAgUiRunOptions { + /** Keep the harness-native `CUSTOM` control events. Default: `false`. */ + includeHarnessEvents?: boolean + /** Emit interim {@link TANSTACK_SPEND_EVENT} ticks per model turn. Default: `false`. */ + emitSpendEvents?: boolean + /** + * `runId`/`threadId` for a synthesized `RUN_STARTED`. A resumed operation + * emits the resolved tool's `TOOL_CALL_RESULT` before any `RUN_STARTED`, so + * the coalescer synthesizes the run start; pass these so it is well-formed. + */ + runId?: string + threadId?: string +} + +/** + * Coalesce ONE harness operation's `SessionEvent` stream into a single, valid + * AG-UI run. + * + * A harness operation can span several model turns (a tool runs, the model is + * called again), and the harness emits a `RUN_STARTED`/`RUN_FINISHED` pair per + * turn. A strict AG-UI client treats each pair as a separate run and rejects + * the `TOOL_CALL_RESULT` that arrives between turns ("the run has already + * finished"). So this emits exactly one `RUN_STARTED` (the first) and one + * `RUN_FINISHED` (carrying the operation's terminal outcome — e.g. an + * interrupt — and the usage summed across every turn), with all content in + * between. Feed it the events of a single operation (see the handler, which + * scopes to one operation per request). + */ +export async function* operationToAgUiRun( + entries: AsyncIterable, + options: OperationToAgUiRunOptions = {}, +): AsyncGenerator { + const includeHarnessEvents = options.includeHarnessEvents ?? false + const emitSpendEvents = options.emitSpendEvents ?? false + let started = false + let terminal: StreamChunk | undefined + let cumulative = zeroUsage() + + const runStart = (seed?: StreamChunk): StreamChunk => { + const s = seed as Record | undefined + return { + type: EventType.RUN_STARTED, + runId: options.runId ?? (s?.runId as string) ?? 'run', + threadId: options.threadId ?? (s?.threadId as string) ?? 'thread', + timestamp: (s?.timestamp as number) ?? Date.now(), + } as StreamChunk + } + + for await (const entry of entries) { + const event = entry.event + if (event.type === EventType.RUN_STARTED) { + if (!started) { + started = true + yield event + } + continue + } + if (event.type === EventType.RUN_FINISHED) { + const usage = usageOf(event) + if (usage) { + cumulative = addUsage(cumulative, usage) + if (emitSpendEvents) yield spendEvent(event, usage, cumulative) + } + // Hold it; the last one becomes the single terminal RUN_FINISHED. + terminal = event + continue + } + if (isHarnessCustom(event)) { + if ((event as { name?: string }).name === HARNESS_EVENTS.operationFinished) { + break + } + if (!includeHarnessEvents) continue + } + // A resumed run streams the resolved tool's result before any RUN_STARTED; + // synthesize one so the AG-UI run always begins with RUN_STARTED. + if (!started) { + started = true + yield runStart(event) + } + yield event + } + + if (!started && !terminal) return + if (!started) { + started = true + yield runStart(terminal) + } + const base = + terminal ?? + ({ type: EventType.RUN_FINISHED, timestamp: Date.now() } as StreamChunk) + yield conformRunFinishedUsage( + base, + cumulative.totalTokens > 0 ? cumulative : undefined, + ) +} + /** Session events for one operation, up to and including its terminal event. */ async function* followOperation( events: AsyncIterable, @@ -308,9 +439,11 @@ export function createAgUiHandler( if (!principal) return json({ error: 'unauthorized' }, 401) let operationId: string + let threadId: string let session: HarnessSession try { const params = await chatParamsFromRequestBody(await request.json()) + threadId = params.threadId const opened = await openFor(principal, params.threadId) if (!opened) return json({ error: 'forbidden' }, 403) session = opened @@ -351,9 +484,11 @@ export function createAgUiHandler( session.events({ signal: reader.signal }), operationId, ) - for await (const event of sessionEventsToAgUi(entries, { + for await (const event of operationToAgUiRun(entries, { includeHarnessEvents: false, ...options.stream, + runId: operationId, + threadId, })) { controller.enqueue(encoder.encodeBinary(event as BaseEvent)) } diff --git a/packages/ai-harness/tests/ag-ui.test.ts b/packages/ai-harness/tests/ag-ui.test.ts index 3af515dd0e..876479f849 100644 --- a/packages/ai-harness/tests/ag-ui.test.ts +++ b/packages/ai-harness/tests/ag-ui.test.ts @@ -3,11 +3,12 @@ import { z } from 'zod' import { HttpAgent } from '@ag-ui/client' import { EventType, toolDefinition } from '@tanstack/ai' import { memoryPersistence } from '@tanstack/ai-persistence' -import { createHarnessHost, defineHarness } from '../src' +import { HARNESS_EVENTS, createHarnessHost, defineHarness } from '../src' import { TANSTACK_SPEND_EVENT, createAgUiHandler, normalizeUsage, + operationToAgUiRun, sessionEventsToAgUi, } from '../src/ag-ui' import { mockAdapter, text, toolCall } from './helpers' @@ -93,19 +94,28 @@ describe('sessionEventsToAgUi mapper', () => { }) }) - it('does not overwrite an existing metadata.tanstack.usage', async () => { - const preset = { inputTokens: 1, outputTokens: 1, totalTokens: 2 } + it('conforms a TokenUsage object to a spec usage array', async () => { + // The harness emits usage as a TanStack TokenUsage object; a strict AG-UI + // client requires SpecTokenUsage[]. const [event] = await collect( sessionEventsToAgUi( feed([ runFinished({ - usage: [{ inputTokens: 99, outputTokens: 99 }], - metadata: { tanstack: { usage: preset } }, + usage: { promptTokens: 8, completionTokens: 2, totalTokens: 10 }, }), ]), ), ) - expect((event as any).metadata.tanstack.usage).toEqual(preset) + expect(Array.isArray((event as any).usage)).toBe(true) + expect((event as any).usage).toEqual([{ inputTokens: 8, outputTokens: 2 }]) + expect((event as any).metadata.tanstack.usage.totalTokens).toBe(10) + }) + + it('drops a non-array usage when there is nothing to report', async () => { + const [event] = await collect( + sessionEventsToAgUi(feed([runFinished({ usage: 'n/a' })])), + ) + expect((event as any).usage).toBeUndefined() }) it('keeps interrupt outcomes intact for the AG-UI resume flow', async () => { @@ -190,6 +200,87 @@ describe('sessionEventsToAgUi mapper', () => { }) }) +describe('operationToAgUiRun coalescer', () => { + it('collapses a multi-turn operation into one RUN_STARTED/RUN_FINISHED', async () => { + // Two model turns (tool then final), as the harness emits per turn, plus + // the harness operation.finished terminator. + const stream: Array = [ + { type: EventType.RUN_STARTED, runId: 'r1', threadId: 't', timestamp: 0 } as StreamChunk, + { + type: EventType.TOOL_CALL_START, + toolCallId: 'c1', + toolCallName: 'lookup', + timestamp: 0, + } as StreamChunk, + { type: EventType.TOOL_CALL_END, toolCallId: 'c1', timestamp: 0 } as StreamChunk, + runFinished({ runId: 'r1', usage: [{ inputTokens: 10, outputTokens: 5 }] }), + { + type: EventType.TOOL_CALL_RESULT, + toolCallId: 'c1', + messageId: 'm1', + content: '{"ok":true}', + timestamp: 0, + } as StreamChunk, + { type: EventType.RUN_STARTED, runId: 'r2', threadId: 't', timestamp: 0 } as StreamChunk, + runFinished({ runId: 'r2', usage: [{ inputTokens: 2, outputTokens: 3 }] }), + custom(HARNESS_EVENTS.operationFinished, { operationId: 'op-1', status: 'completed' }), + ] + const out = await collect(operationToAgUiRun(feed(stream))) + expect(typesOf(out).filter((t) => t === EventType.RUN_STARTED)).toHaveLength(1) + expect(typesOf(out).filter((t) => t === EventType.RUN_FINISHED)).toHaveLength(1) + // The tool result survives inside the single run. + expect(out.some((e) => e.type === EventType.TOOL_CALL_RESULT)).toBe(true) + // First is RUN_STARTED, last is RUN_FINISHED (a valid AG-UI run). + expect(out[0]!.type).toBe(EventType.RUN_STARTED) + expect(out.at(-1)!.type).toBe(EventType.RUN_FINISHED) + // Usage is summed across both turns. + expect((out.at(-1) as any).metadata.tanstack.usage.totalTokens).toBe(20) + }) + + it('synthesizes a leading RUN_STARTED when a resumed run leads with a tool result', async () => { + // A resumed operation streams the resolved tool's result before RUN_STARTED. + const stream: Array = [ + { + type: EventType.TOOL_CALL_RESULT, + toolCallId: 'c1', + messageId: 'm1', + content: '{"sent":true}', + timestamp: 0, + } as StreamChunk, + { type: EventType.RUN_STARTED, runId: 'r2', threadId: 't', timestamp: 0 } as StreamChunk, + runFinished({}), + custom(HARNESS_EVENTS.operationFinished, { operationId: 'op-2', status: 'completed' }), + ] + const out = await collect( + operationToAgUiRun(feed(stream), { runId: 'op-2', threadId: 't' }), + ) + expect(out[0]!.type).toBe(EventType.RUN_STARTED) + expect((out[0] as any).runId).toBe('op-2') + expect(out[1]!.type).toBe(EventType.TOOL_CALL_RESULT) + expect(typesOf(out).filter((t) => t === EventType.RUN_STARTED)).toHaveLength(1) + expect(out.at(-1)!.type).toBe(EventType.RUN_FINISHED) + }) + + it('carries the terminal interrupt outcome onto the single RUN_FINISHED', async () => { + const stream: Array = [ + { type: EventType.RUN_STARTED, runId: 'r1', threadId: 't', timestamp: 0 } as StreamChunk, + runFinished({}), + { type: EventType.RUN_STARTED, runId: 'r2', threadId: 't', timestamp: 0 } as StreamChunk, + runFinished({ + outcome: { + type: 'interrupt', + interrupts: [{ id: 'i1', reason: 'tool_call', toolCallId: 'c1' }], + }, + }), + custom(HARNESS_EVENTS.operationFinished, { operationId: 'op-1', status: 'interrupted' }), + ] + const out = await collect(operationToAgUiRun(feed(stream))) + const finished = out.filter((e) => e.type === EventType.RUN_FINISHED) + expect(finished).toHaveLength(1) + expect((finished[0] as any).outcome.type).toBe('interrupt') + }) +}) + /** Build a POST request shaped like an AG-UI `RunAgentInput`. */ function runRequest(body: Record): Request { return new Request('http://localhost/agent', { @@ -404,4 +495,57 @@ describe('consumed by a bare @ag-ui/client HttpAgent', () => { expect(agent.pendingInterrupts[0]?.toolCallId).toBe('call_1') await host.close() }) + + it('handles a multi-turn operation (auto tool then approval) in one run', async () => { + const host = createHarnessHost({ persistence: memoryPersistence() }) + const lookup = toolDefinition({ + name: 'lookup', + description: 'Look up a ticket', + inputSchema: z.object({ id: z.string() }), + }).server(async () => ({ ok: true })) + const send = toolDefinition({ + name: 'send', + description: 'Send a reply', + needsApproval: true, + inputSchema: z.object({ to: z.string() }), + }).server(async () => ({ sent: true })) + const { adapter } = mockAdapter([ + () => toolCall('lookup', { id: 'T-1' }, 'c1'), + () => toolCall('send', { to: 'a@b.c' }, 'c2'), + () => text('done'), + ]) + const harness = defineHarness({ + name: 'test/client-multi', + adapter, + tools: [lookup, send], + }) + const handler = createAgUiHandler({ + host, + harness, + authorize: () => ({ id: 'u' }), + }) + + const agent = new HttpAgent({ + url: 'http://localhost/agent', + threadId: 't1', + fetch: (url: string, init: RequestInit) => + handler(new Request(url, init)), + }) + agent.addMessage({ id: 'u1', role: 'user', content: 'handle it' }) + // Would throw "run has already finished" without single-run coalescing. + await agent.runAgent() + + expect(agent.pendingInterrupts).toHaveLength(1) + expect(agent.pendingInterrupts[0]?.toolCallId).toBe('c2') + + // Resume over AG-UI: the continuation leads with the tool result, so the + // coalescer must synthesize a leading RUN_STARTED or the client rejects it. + const interruptId = agent.pendingInterrupts[0]!.id + await agent.runAgent({ + resume: [{ interruptId, status: 'resolved', payload: true }], + }) + expect(agent.pendingInterrupts).toHaveLength(0) + expect(JSON.stringify(agent.messages)).toContain('done') + await host.close() + }) }) From de730ecd0208d9fd1d24760865d50e55445ce2c0 Mon Sep 17 00:00:00 2001 From: Jack Herrington Date: Sat, 26 Sep 2026 08:45:09 -0700 Subject: [PATCH 03/39] feat(examples/agent-dashboard): live session view + approval queue MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A TanStack Start dashboard — mission control for TanStack AI agents. It embeds a deterministic support-triage agent (no API key needed) that looks up a ticket, drafts a reply, and pauses for approval before sending. - Session detail consumes the AG-UI SSE stream via a bare @ag-ui/client HttpAgent (@tanstack/ai-harness/ag-ui), projected into TanStack DB localOnly collections read with useLiveQuery — the UI is a projection of the stream. - Inline approval queue: approve (AG-UI resume), deny (harness control tier), or edit the tool args, mid-run. - Live spend meter from tanstack.spend events; hosts/sessions lists via TanStack Query. - Server routes host the agent: /api/agent (AG-UI run stream), /api/harness/* (snapshot/events/control), /api/hosts, /api/sessions. - Playwright e2e drives the full stream + mid-run approval. Generated with Claude Code. Co-Authored-By: Claude Opus 4.8 (1M context) --- examples/agent-dashboard/.gitignore | 14 + examples/agent-dashboard/README.md | 50 ++++ .../agent-dashboard/e2e/dashboard.spec.ts | 32 +++ examples/agent-dashboard/package.json | 44 +++ examples/agent-dashboard/playwright.config.ts | 24 ++ .../agent-dashboard/src/db/collections.ts | 87 ++++++ examples/agent-dashboard/src/lib/query.ts | 12 + .../src/lib/session-controller.ts | 212 ++++++++++++++ examples/agent-dashboard/src/routeTree.gen.ts | 177 ++++++++++++ examples/agent-dashboard/src/router.tsx | 10 + .../agent-dashboard/src/routes/__root.tsx | 59 ++++ .../agent-dashboard/src/routes/api.agent.ts | 22 ++ .../src/routes/api.harness.$.ts | 22 ++ .../agent-dashboard/src/routes/api.hosts.ts | 21 ++ .../src/routes/api.sessions.ts | 27 ++ examples/agent-dashboard/src/routes/index.tsx | 111 +++++++ .../src/routes/sessions.$threadId.tsx | 270 ++++++++++++++++++ .../agent-dashboard/src/server/harness.ts | 204 +++++++++++++ examples/agent-dashboard/src/styles.css | 11 + examples/agent-dashboard/tsconfig.json | 26 ++ examples/agent-dashboard/vite.config.ts | 29 ++ 21 files changed, 1464 insertions(+) create mode 100644 examples/agent-dashboard/.gitignore create mode 100644 examples/agent-dashboard/README.md create mode 100644 examples/agent-dashboard/e2e/dashboard.spec.ts create mode 100644 examples/agent-dashboard/package.json create mode 100644 examples/agent-dashboard/playwright.config.ts create mode 100644 examples/agent-dashboard/src/db/collections.ts create mode 100644 examples/agent-dashboard/src/lib/query.ts create mode 100644 examples/agent-dashboard/src/lib/session-controller.ts create mode 100644 examples/agent-dashboard/src/routeTree.gen.ts create mode 100644 examples/agent-dashboard/src/router.tsx create mode 100644 examples/agent-dashboard/src/routes/__root.tsx create mode 100644 examples/agent-dashboard/src/routes/api.agent.ts create mode 100644 examples/agent-dashboard/src/routes/api.harness.$.ts create mode 100644 examples/agent-dashboard/src/routes/api.hosts.ts create mode 100644 examples/agent-dashboard/src/routes/api.sessions.ts create mode 100644 examples/agent-dashboard/src/routes/index.tsx create mode 100644 examples/agent-dashboard/src/routes/sessions.$threadId.tsx create mode 100644 examples/agent-dashboard/src/server/harness.ts create mode 100644 examples/agent-dashboard/src/styles.css create mode 100644 examples/agent-dashboard/tsconfig.json create mode 100644 examples/agent-dashboard/vite.config.ts diff --git a/examples/agent-dashboard/.gitignore b/examples/agent-dashboard/.gitignore new file mode 100644 index 0000000000..a0c89c8e56 --- /dev/null +++ b/examples/agent-dashboard/.gitignore @@ -0,0 +1,14 @@ +node_modules +.DS_Store +dist +dist-ssr +*.local +.env +.env.local +.nitro +.tanstack +.output +.vinxi +test-results +playwright-report +*.log diff --git a/examples/agent-dashboard/README.md b/examples/agent-dashboard/README.md new file mode 100644 index 0000000000..cfd1b9d47a --- /dev/null +++ b/examples/agent-dashboard/README.md @@ -0,0 +1,50 @@ +# Agent Dashboard + +Mission control for TanStack AI agents — a TanStack Start app that watches live +agent sessions, approves tool calls mid-run, and tracks spend, all as a live +projection of the [AG-UI](../../docs/harness/ag-ui.md) event stream. + +It embeds a deterministic **support-triage** agent, so it runs with **no API +key**: the agent looks up a ticket (an auto tool), drafts a customer reply (an +approval-gated tool that pauses the run), and sends it once you approve. + +```bash +pnpm --filter agent-dashboard dev # http://localhost:3002 +``` + +## What it exercises + +- **TanStack Start** — app shell + routing: hosts → sessions → session detail + (`src/routes`), and server API routes that host the agent (`src/routes/api.*`). +- **`@tanstack/ai-harness/ag-ui`** — the session view consumes the AG-UI SSE + stream with a bare `@ag-ui/client` `HttpAgent` (`src/lib/session-controller.ts`); + no bespoke protocol code. +- **TanStack DB** — run state (messages, tool calls, approvals, spend) lives in + `localOnly` collections written from the stream and read with `useLiveQuery` + (`src/db/collections.ts`). The UI is a projection of the stream, not a poller. +- **TanStack Query** — server state: the host and session lists. + +## The approval queue + +When a run pauses on an approval, an approval card appears. It resolves the +interrupt through **both** paths the interrupt supports: + +- **Approve** → the AG-UI resume flow (`runAgent({ resume })` over `/api/agent`), + so the continuation streams back into the view. +- **Deny** → the harness-native control endpoint (`/api/harness/control`). +- **Edit** → approve with edited tool arguments. + +## Server wiring + +- `POST /api/agent` — the AG-UI run stream (`createAgUiHandler`, spend ticks on). +- `GET|POST /api/harness/*` — the harness control tier (`createHarnessHandler`): + `snapshot`, `events`, `control`. +- `GET /api/hosts`, `GET /api/sessions` — host and session lists. + +## Test + +```bash +pnpm --filter agent-dashboard test:e2e # Playwright: stream + approve mid-run +``` + +Generated with Claude Code. diff --git a/examples/agent-dashboard/e2e/dashboard.spec.ts b/examples/agent-dashboard/e2e/dashboard.spec.ts new file mode 100644 index 0000000000..fb309066b0 --- /dev/null +++ b/examples/agent-dashboard/e2e/dashboard.spec.ts @@ -0,0 +1,32 @@ +import { expect, test } from '@playwright/test' + +// The dashboard embeds a deterministic support-triage agent, so this runs with +// no API key: send a prompt, watch the AG-UI stream, approve a tool call +// mid-run, and see the run resume — the north-star approval demo. +test('streams a session and approves a tool call mid-run', async ({ page }) => { + const threadId = `e2e-${Date.now()}` + await page.goto(`/sessions/${threadId}`) + + const startButton = page.getByRole('button', { name: 'Start triage demo' }) + await expect(startButton).toBeVisible() + await startButton.click() + + // The agent looks up the ticket (auto tool) then drafts a reply for approval. + await expect(page.getByText('lookup_ticket').first()).toBeVisible() + await expect( + page.getByText("Here's a draft reply for your approval."), + ).toBeVisible() + + // The approval card appears (the run paused on the interrupt). + await expect(page.getByText('Approval required', { exact: true })).toBeVisible() + await expect(page.getByText('send_reply').first()).toBeVisible() + + // Spend meter is live (tokens accrued from the stream). + await expect(page.getByText(/[1-9][0-9,]* tokens/)).toBeVisible() + + // Approve via the AG-UI resume flow; the run continues and finishes. + await page.getByRole('button', { name: 'Approve', exact: true }).click() + + await expect(page.getByText(/Sent ✅/)).toBeVisible() + await expect(page.getByText('Approval required', { exact: true })).toHaveCount(0) +}) diff --git a/examples/agent-dashboard/package.json b/examples/agent-dashboard/package.json new file mode 100644 index 0000000000..2309dc6b00 --- /dev/null +++ b/examples/agent-dashboard/package.json @@ -0,0 +1,44 @@ +{ + "name": "agent-dashboard", + "private": true, + "type": "module", + "scripts": { + "dev": "vite dev --port 3002", + "build": "vite build", + "serve": "vite preview", + "start": "node .output/server/index.mjs", + "test": "exit 0", + "test:e2e": "playwright test", + "test:types": "tsc --noEmit" + }, + "dependencies": { + "@ag-ui/client": "1.0.0", + "@tailwindcss/vite": "^4.1.18", + "@tanstack/ai": "workspace:*", + "@tanstack/ai-anthropic": "workspace:*", + "@tanstack/ai-harness": "workspace:*", + "@tanstack/ai-persistence": "workspace:*", + "@tanstack/react-db": "^0.1.55", + "@tanstack/react-devtools": "^0.9.10", + "@tanstack/react-query": "^5.90.12", + "@tanstack/react-router": "^1.158.4", + "@tanstack/react-router-devtools": "^1.158.4", + "@tanstack/react-start": "^1.159.0", + "@tanstack/router-plugin": "^1.158.4", + "nitro": "3.0.260610-beta", + "react": "^19.2.3", + "react-dom": "^19.2.3", + "tailwindcss": "^4.1.18", + "zod": "^4.2.0" + }, + "devDependencies": { + "@playwright/test": "^1.57.0", + "@tanstack/devtools-vite": "^0.5.3", + "@types/node": "^24.10.1", + "@types/react": "^19.2.7", + "@types/react-dom": "^19.2.3", + "@vitejs/plugin-react": "^5.2.0", + "typescript": "5.9.3", + "vite": "^8.2.1" + } +} diff --git a/examples/agent-dashboard/playwright.config.ts b/examples/agent-dashboard/playwright.config.ts new file mode 100644 index 0000000000..39f5a08d4e --- /dev/null +++ b/examples/agent-dashboard/playwright.config.ts @@ -0,0 +1,24 @@ +import { defineConfig, devices } from '@playwright/test' + +export default defineConfig({ + testDir: './e2e', + fullyParallel: false, + forbidOnly: !!process.env.CI, + retries: 0, + workers: 1, + reporter: [['list']], + timeout: 30_000, + expect: { timeout: 15_000 }, + use: { + baseURL: 'http://localhost:3002', + screenshot: 'only-on-failure', + trace: 'on-first-retry', + }, + projects: [{ name: 'chromium', use: { ...devices['Desktop Chrome'] } }], + webServer: { + command: 'pnpm run dev', + url: 'http://localhost:3002', + reuseExistingServer: !process.env.CI, + timeout: 120_000, + }, +}) diff --git a/examples/agent-dashboard/src/db/collections.ts b/examples/agent-dashboard/src/db/collections.ts new file mode 100644 index 0000000000..cb10c5b64e --- /dev/null +++ b/examples/agent-dashboard/src/db/collections.ts @@ -0,0 +1,87 @@ +/** + * TanStack DB collections. The dashboard's run state is a projection of the + * AG-UI event stream: the session controller writes events here, and the UI + * reads them with live queries — no polling. + */ +import { createCollection, localOnlyCollectionOptions } from '@tanstack/react-db' + +export interface MessageRow { + id: string + threadId: string + role: 'user' | 'assistant' + text: string + /** Set for text produced by a subagent, for attribution. */ + subagentRunId?: string + createdAt: number +} + +export interface ToolCallRow { + id: string + threadId: string + name: string + args: string + result?: string + status: 'running' | 'done' + subagentRunId?: string + createdAt: number +} + +export interface ApprovalRow { + id: string + threadId: string + toolCallId?: string + reason: string + message: string + responseSchema?: Record + status: 'pending' | 'approved' | 'denied' + createdAt: number +} + +export interface SpendRow { + id: string + threadId: string + inputTokens: number + outputTokens: number + totalTokens: number +} + +export interface SessionRow { + id: string + threadId: string + status: 'idle' | 'running' | 'requires_action' + createdAt: number +} + +export const messages = createCollection( + localOnlyCollectionOptions({ getKey: (row: MessageRow) => row.id }), +) +export const toolCalls = createCollection( + localOnlyCollectionOptions({ getKey: (row: ToolCallRow) => row.id }), +) +export const approvals = createCollection( + localOnlyCollectionOptions({ getKey: (row: ApprovalRow) => row.id }), +) +export const spend = createCollection( + localOnlyCollectionOptions({ getKey: (row: SpendRow) => row.id }), +) +export const sessions = createCollection( + localOnlyCollectionOptions({ getKey: (row: SessionRow) => row.id }), +) + +/** + * Insert if absent, else apply the updater. localOnly writes are synchronous. + * `collection` is a TanStack DB collection; it is typed loosely here so `T` is + * pinned by the row, not the collection's overloaded `insert`/`update`. + */ +export function upsert( + // eslint-disable-next-line @typescript-eslint/no-explicit-any + collection: any, + row: T, + updater?: (draft: T) => void, +): void { + if (collection.has(row.id)) { + collection.update(row.id, (draft: T) => (updater ?? (() => {}))(draft)) + } else { + collection.insert(row) + } +} diff --git a/examples/agent-dashboard/src/lib/query.ts b/examples/agent-dashboard/src/lib/query.ts new file mode 100644 index 0000000000..adb047c44d --- /dev/null +++ b/examples/agent-dashboard/src/lib/query.ts @@ -0,0 +1,12 @@ +import { QueryClient } from '@tanstack/react-query' + +let client: QueryClient | undefined + +export function getQueryClient(): QueryClient { + client ??= new QueryClient({ + defaultOptions: { + queries: { staleTime: 2000, refetchOnWindowFocus: false }, + }, + }) + return client +} diff --git a/examples/agent-dashboard/src/lib/session-controller.ts b/examples/agent-dashboard/src/lib/session-controller.ts new file mode 100644 index 0000000000..d5484dfc6f --- /dev/null +++ b/examples/agent-dashboard/src/lib/session-controller.ts @@ -0,0 +1,212 @@ +/** + * The live session controller. One `@ag-ui/client` HttpAgent per thread, whose + * AG-UI events are projected into TanStack DB collections. The UI is then a live + * query over those collections — a projection of the event stream, not a poller. + */ +import { HttpAgent } from '@ag-ui/client' +import { + approvals, + messages, + sessions, + spend, + toolCalls, + upsert, +} from '@/db/collections' + +interface Entry { + agent: HttpAgent + subscribed: boolean +} +const registry = new Map() + +function agentUrl() { + const origin = + typeof window !== 'undefined' ? window.location.origin : 'http://localhost' + return `${origin}/api/agent` +} + +function setStatus(threadId: string, status: 'idle' | 'running' | 'requires_action') { + upsert( + sessions, + { id: threadId, threadId, status, createdAt: Date.now() }, + (draft) => { + draft.status = status + }, + ) +} + +/** Project one AG-UI event into the collections. */ +function project(threadId: string, event: any) { + switch (event.type) { + case 'RUN_STARTED': + setStatus(threadId, 'running') + break + case 'TEXT_MESSAGE_START': + upsert(messages, { + id: event.messageId, + threadId, + role: event.role === 'user' ? 'user' : 'assistant', + text: '', + ...(event.subagentRunId ? { subagentRunId: event.subagentRunId } : {}), + createdAt: Date.now(), + }) + break + case 'TEXT_MESSAGE_CONTENT': + if (messages.has(event.messageId)) { + messages.update(event.messageId, (draft) => { + draft.text += event.delta ?? '' + }) + } + break + case 'TOOL_CALL_START': + upsert(toolCalls, { + id: event.toolCallId, + threadId, + name: event.toolCallName ?? 'tool', + args: '', + status: 'running', + ...(event.subagentRunId ? { subagentRunId: event.subagentRunId } : {}), + createdAt: Date.now(), + }) + break + case 'TOOL_CALL_ARGS': + if (toolCalls.has(event.toolCallId)) { + toolCalls.update(event.toolCallId, (draft) => { + draft.args += event.delta ?? '' + }) + } + break + case 'TOOL_CALL_RESULT': + if (toolCalls.has(event.toolCallId)) { + toolCalls.update(event.toolCallId, (draft) => { + draft.result = + typeof event.content === 'string' + ? event.content + : JSON.stringify(event.content) + draft.status = 'done' + }) + } + break + case 'CUSTOM': + if (event.name === 'tanstack.spend') { + const c = event.value?.cumulative ?? {} + upsert( + spend, + { + id: threadId, + threadId, + inputTokens: c.inputTokens ?? 0, + outputTokens: c.outputTokens ?? 0, + totalTokens: c.totalTokens ?? 0, + }, + (draft) => { + draft.inputTokens = c.inputTokens ?? draft.inputTokens + draft.outputTokens = c.outputTokens ?? draft.outputTokens + draft.totalTokens = c.totalTokens ?? draft.totalTokens + }, + ) + } + break + case 'RUN_FINISHED': + if (event.outcome?.type === 'interrupt') { + for (const interrupt of event.outcome.interrupts ?? []) { + upsert(approvals, { + id: interrupt.id, + threadId, + toolCallId: interrupt.toolCallId, + reason: interrupt.reason ?? 'tool_call', + message: interrupt.message ?? 'Approval required', + responseSchema: interrupt.responseSchema, + status: 'pending', + createdAt: Date.now(), + }) + } + setStatus(threadId, 'requires_action') + } else { + setStatus(threadId, 'idle') + } + break + default: + break + } +} + +function getEntry(threadId: string): Entry { + let entry = registry.get(threadId) + if (!entry) { + const agent = new HttpAgent({ url: agentUrl(), threadId }) + entry = { agent, subscribed: false } + registry.set(threadId, entry) + } + if (!entry.subscribed) { + entry.agent.subscribe({ + onEvent: ({ event }) => { + project(threadId, event as any) + }, + }) + entry.subscribed = true + } + return entry +} + +/** Create + subscribe the agent for a thread so its events start projecting. */ +export function ensureSession(threadId: string): void { + getEntry(threadId) +} + +/** Send a user prompt and stream the reply into the collections. */ +export async function sendPrompt(threadId: string, text: string): Promise { + const { agent } = getEntry(threadId) + upsert(messages, { + id: `user-${Date.now()}`, + threadId, + role: 'user', + text, + createdAt: Date.now(), + }) + agent.addMessage({ id: `u-${Date.now()}`, role: 'user', content: text }) + await agent.runAgent() +} + +export type ApprovalDecision = 'approve' | 'deny' + +/** + * Resolve an approval. `approve` goes through the AG-UI resume flow (a new run + * over /api/agent, so the continuation streams back). `deny` goes through the + * harness-native control endpoint — exercising both paths the spec asks for. + */ +export async function resolveApproval( + threadId: string, + interruptId: string, + decision: ApprovalDecision, + editedArgs?: Record, +): Promise { + if (approvals.has(interruptId)) { + approvals.update(interruptId, (draft) => { + draft.status = decision === 'approve' ? 'approved' : 'denied' + }) + } + + if (decision === 'approve') { + const payload = editedArgs ? { approved: true, editedArgs } : true + const { agent } = getEntry(threadId) + await agent.runAgent({ + resume: [{ interruptId, status: 'resolved', payload }], + }) + return + } + + // Harness-native path. + await fetch(`${window.location.origin}/api/harness/control`, { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ + threadId, + input: { + op: 'resolve', + resume: [{ interruptId, status: 'cancelled', payload: false }], + }, + }), + }) + setStatus(threadId, 'idle') +} diff --git a/examples/agent-dashboard/src/routeTree.gen.ts b/examples/agent-dashboard/src/routeTree.gen.ts new file mode 100644 index 0000000000..a5b4c17057 --- /dev/null +++ b/examples/agent-dashboard/src/routeTree.gen.ts @@ -0,0 +1,177 @@ +/* eslint-disable */ + +// @ts-nocheck + +// noinspection JSUnusedGlobalSymbols + +// This file was automatically generated by TanStack Router. +// You should NOT make any changes in this file as it will be overwritten. +// Additionally, you should also exclude this file from your linter and/or formatter to prevent it from being checked or modified. + +import { Route as rootRouteImport } from './routes/__root' +import { Route as IndexRouteImport } from './routes/index' +import { Route as SessionsThreadIdRouteImport } from './routes/sessions.$threadId' +import { Route as ApiSessionsRouteImport } from './routes/api.sessions' +import { Route as ApiHostsRouteImport } from './routes/api.hosts' +import { Route as ApiAgentRouteImport } from './routes/api.agent' +import { Route as ApiHarnessSplatRouteImport } from './routes/api.harness.$' + +const IndexRoute = IndexRouteImport.update({ + id: '/', + path: '/', + getParentRoute: () => rootRouteImport, +} as any) +const SessionsThreadIdRoute = SessionsThreadIdRouteImport.update({ + id: '/sessions/$threadId', + path: '/sessions/$threadId', + getParentRoute: () => rootRouteImport, +} as any) +const ApiSessionsRoute = ApiSessionsRouteImport.update({ + id: '/api/sessions', + path: '/api/sessions', + getParentRoute: () => rootRouteImport, +} as any) +const ApiHostsRoute = ApiHostsRouteImport.update({ + id: '/api/hosts', + path: '/api/hosts', + getParentRoute: () => rootRouteImport, +} as any) +const ApiAgentRoute = ApiAgentRouteImport.update({ + id: '/api/agent', + path: '/api/agent', + getParentRoute: () => rootRouteImport, +} as any) +const ApiHarnessSplatRoute = ApiHarnessSplatRouteImport.update({ + id: '/api/harness/$', + path: '/api/harness/$', + getParentRoute: () => rootRouteImport, +} as any) + +export interface FileRoutesByFullPath { + '/': typeof IndexRoute + '/api/agent': typeof ApiAgentRoute + '/api/hosts': typeof ApiHostsRoute + '/api/sessions': typeof ApiSessionsRoute + '/sessions/$threadId': typeof SessionsThreadIdRoute + '/api/harness/$': typeof ApiHarnessSplatRoute +} +export interface FileRoutesByTo { + '/': typeof IndexRoute + '/api/agent': typeof ApiAgentRoute + '/api/hosts': typeof ApiHostsRoute + '/api/sessions': typeof ApiSessionsRoute + '/sessions/$threadId': typeof SessionsThreadIdRoute + '/api/harness/$': typeof ApiHarnessSplatRoute +} +export interface FileRoutesById { + __root__: typeof rootRouteImport + '/': typeof IndexRoute + '/api/agent': typeof ApiAgentRoute + '/api/hosts': typeof ApiHostsRoute + '/api/sessions': typeof ApiSessionsRoute + '/sessions/$threadId': typeof SessionsThreadIdRoute + '/api/harness/$': typeof ApiHarnessSplatRoute +} +export interface FileRouteTypes { + fileRoutesByFullPath: FileRoutesByFullPath + fullPaths: + | '/' + | '/api/agent' + | '/api/hosts' + | '/api/sessions' + | '/sessions/$threadId' + | '/api/harness/$' + fileRoutesByTo: FileRoutesByTo + to: + | '/' + | '/api/agent' + | '/api/hosts' + | '/api/sessions' + | '/sessions/$threadId' + | '/api/harness/$' + id: + | '__root__' + | '/' + | '/api/agent' + | '/api/hosts' + | '/api/sessions' + | '/sessions/$threadId' + | '/api/harness/$' + fileRoutesById: FileRoutesById +} +export interface RootRouteChildren { + IndexRoute: typeof IndexRoute + ApiAgentRoute: typeof ApiAgentRoute + ApiHostsRoute: typeof ApiHostsRoute + ApiSessionsRoute: typeof ApiSessionsRoute + SessionsThreadIdRoute: typeof SessionsThreadIdRoute + ApiHarnessSplatRoute: typeof ApiHarnessSplatRoute +} + +declare module '@tanstack/react-router' { + interface FileRoutesByPath { + '/': { + id: '/' + path: '/' + fullPath: '/' + preLoaderRoute: typeof IndexRouteImport + parentRoute: typeof rootRouteImport + } + '/sessions/$threadId': { + id: '/sessions/$threadId' + path: '/sessions/$threadId' + fullPath: '/sessions/$threadId' + preLoaderRoute: typeof SessionsThreadIdRouteImport + parentRoute: typeof rootRouteImport + } + '/api/sessions': { + id: '/api/sessions' + path: '/api/sessions' + fullPath: '/api/sessions' + preLoaderRoute: typeof ApiSessionsRouteImport + parentRoute: typeof rootRouteImport + } + '/api/hosts': { + id: '/api/hosts' + path: '/api/hosts' + fullPath: '/api/hosts' + preLoaderRoute: typeof ApiHostsRouteImport + parentRoute: typeof rootRouteImport + } + '/api/agent': { + id: '/api/agent' + path: '/api/agent' + fullPath: '/api/agent' + preLoaderRoute: typeof ApiAgentRouteImport + parentRoute: typeof rootRouteImport + } + '/api/harness/$': { + id: '/api/harness/$' + path: '/api/harness/$' + fullPath: '/api/harness/$' + preLoaderRoute: typeof ApiHarnessSplatRouteImport + parentRoute: typeof rootRouteImport + } + } +} + +const rootRouteChildren: RootRouteChildren = { + IndexRoute: IndexRoute, + ApiAgentRoute: ApiAgentRoute, + ApiHostsRoute: ApiHostsRoute, + ApiSessionsRoute: ApiSessionsRoute, + SessionsThreadIdRoute: SessionsThreadIdRoute, + ApiHarnessSplatRoute: ApiHarnessSplatRoute, +} +export const routeTree = rootRouteImport + ._addFileChildren(rootRouteChildren) + ._addFileTypes() + +import type { getRouter } from './router.tsx' +import type { createStart } from '@tanstack/react-start' +declare module '@tanstack/react-start' { + interface Register { + ssr: true + router: Awaited> + } +} diff --git a/examples/agent-dashboard/src/router.tsx b/examples/agent-dashboard/src/router.tsx new file mode 100644 index 0000000000..a595444648 --- /dev/null +++ b/examples/agent-dashboard/src/router.tsx @@ -0,0 +1,10 @@ +import { createRouter } from '@tanstack/react-router' +import { routeTree } from './routeTree.gen' + +export const getRouter = () => { + return createRouter({ + routeTree, + scrollRestoration: true, + defaultPreloadStaleTime: 0, + }) +} diff --git a/examples/agent-dashboard/src/routes/__root.tsx b/examples/agent-dashboard/src/routes/__root.tsx new file mode 100644 index 0000000000..a919f71df2 --- /dev/null +++ b/examples/agent-dashboard/src/routes/__root.tsx @@ -0,0 +1,59 @@ +import { + HeadContent, + Link, + Scripts, + createRootRoute, +} from '@tanstack/react-router' +import { TanStackRouterDevtoolsPanel } from '@tanstack/react-router-devtools' +import { TanStackDevtools } from '@tanstack/react-devtools' +import { QueryClientProvider } from '@tanstack/react-query' +import { getQueryClient } from '@/lib/query' +import appCss from '../styles.css?url' + +export const Route = createRootRoute({ + head: () => ({ + meta: [ + { charSet: 'utf-8' }, + { name: 'viewport', content: 'width=device-width, initial-scale=1' }, + { title: 'Agent Dashboard' }, + ], + links: [{ rel: 'stylesheet', href: appCss }], + }), + shellComponent: RootDocument, +}) + +function RootDocument({ children }: { children: React.ReactNode }) { + return ( + + + + + + +
+
+ + 🛰️ Agent Dashboard + + + mission control for TanStack AI agents + +
+
{children}
+
+
+ , + }, + ]} + eventBusConfig={{ connectToServerBus: false }} + /> + + + + ) +} diff --git a/examples/agent-dashboard/src/routes/api.agent.ts b/examples/agent-dashboard/src/routes/api.agent.ts new file mode 100644 index 0000000000..6ea1d1e414 --- /dev/null +++ b/examples/agent-dashboard/src/routes/api.agent.ts @@ -0,0 +1,22 @@ +import { createFileRoute } from '@tanstack/react-router' +import { createAgUiHandler } from '@tanstack/ai-harness/ag-ui' +import { authorize, canAccess, getHost, triage } from '@/server/harness' + +// The AG-UI run stream. A bare @ag-ui/client HttpAgent POSTs a RunAgentInput +// (prompt or resume) and gets AG-UI events back as SSE. Spend ticks are on so +// the dashboard can meter tokens live. +const handler = createAgUiHandler({ + host: getHost(), + harness: triage, + authorize, + canAccess, + stream: { emitSpendEvents: true }, +}) + +export const Route = createFileRoute('/api/agent')({ + server: { + handlers: { + POST: ({ request }) => handler(request), + }, + }, +}) diff --git a/examples/agent-dashboard/src/routes/api.harness.$.ts b/examples/agent-dashboard/src/routes/api.harness.$.ts new file mode 100644 index 0000000000..47e6a87ef1 --- /dev/null +++ b/examples/agent-dashboard/src/routes/api.harness.$.ts @@ -0,0 +1,22 @@ +import { createFileRoute } from '@tanstack/react-router' +import { createHarnessHandler } from '@tanstack/ai-harness' +import { authorize, canAccess, getHost, triage } from '@/server/harness' + +// The harness protocol control tier: `.../snapshot`, `.../events`, `.../control`. +// The dashboard uses this for the control plane (and the harness-native approval +// path) alongside the AG-UI run stream at /api/agent. +const handler = createHarnessHandler({ + host: getHost(), + harness: triage, + authorize, + canAccess, +}) + +export const Route = createFileRoute('/api/harness/$')({ + server: { + handlers: { + GET: ({ request }) => handler(request), + POST: ({ request }) => handler(request), + }, + }, +}) diff --git a/examples/agent-dashboard/src/routes/api.hosts.ts b/examples/agent-dashboard/src/routes/api.hosts.ts new file mode 100644 index 0000000000..ce304257ff --- /dev/null +++ b/examples/agent-dashboard/src/routes/api.hosts.ts @@ -0,0 +1,21 @@ +import { createFileRoute } from '@tanstack/react-router' +import { listThreads, triage } from '@/server/harness' + +// One embedded host for the demo. Shaped as a list so the UI can grow to the +// relay's multi-host model later. +export const Route = createFileRoute('/api/hosts')({ + server: { + handlers: { + GET: () => + Response.json([ + { + id: 'local', + name: 'Local host', + harness: triage.name, + description: triage.description ?? '', + sessions: listThreads().length, + }, + ]), + }, + }, +}) diff --git a/examples/agent-dashboard/src/routes/api.sessions.ts b/examples/agent-dashboard/src/routes/api.sessions.ts new file mode 100644 index 0000000000..4ed83cbf21 --- /dev/null +++ b/examples/agent-dashboard/src/routes/api.sessions.ts @@ -0,0 +1,27 @@ +import { createFileRoute } from '@tanstack/react-router' +import { getHost, listThreads, triage } from '@/server/harness' + +// List the threads this host has touched, with a live status from each session +// snapshot (running / requires_action / idle). +export const Route = createFileRoute('/api/sessions')({ + server: { + handlers: { + GET: async () => { + const host = getHost() + const sessions = await Promise.all( + listThreads().map(async (thread) => { + const session = await host.open(triage, { threadId: thread.id }) + const snapshot = session.snapshot() + return { + ...thread, + status: snapshot.status, + pendingInterrupts: snapshot.pendingInterrupts.length, + pendingQuestions: snapshot.pendingQuestions.length, + } + }), + ) + return Response.json(sessions) + }, + }, + }, +}) diff --git a/examples/agent-dashboard/src/routes/index.tsx b/examples/agent-dashboard/src/routes/index.tsx new file mode 100644 index 0000000000..ac70425779 --- /dev/null +++ b/examples/agent-dashboard/src/routes/index.tsx @@ -0,0 +1,111 @@ +import { createFileRoute, useNavigate } from '@tanstack/react-router' +import { useQuery } from '@tanstack/react-query' + +export const Route = createFileRoute('/')({ + component: Home, +}) + +interface Host { + id: string + name: string + harness: string + description: string + sessions: number +} + +interface SessionSummary { + id: string + status: 'idle' | 'running' | 'requires_action' + pendingInterrupts: number + lastActivity: number +} + +const statusStyle: Record = { + idle: 'bg-white/10 text-white/60', + running: 'bg-sky-500/20 text-sky-300', + requires_action: 'bg-amber-500/20 text-amber-300', +} + +function Home() { + const navigate = useNavigate() + const hosts = useQuery>({ + queryKey: ['hosts'], + queryFn: () => fetch('/api/hosts').then((r) => r.json()), + }) + const sessions = useQuery>({ + queryKey: ['sessions'], + queryFn: () => fetch('/api/sessions').then((r) => r.json()), + refetchInterval: 2000, + }) + + const newSession = () => { + const id = `triage-${Math.random().toString(36).slice(2, 8)}` + navigate({ to: '/sessions/$threadId', params: { threadId: id } }) + } + + return ( +
+
+
+

Hosts

+ +
+
+ {(hosts.data ?? []).map((host) => ( +
+
+ + {host.name} + + {host.sessions} session{host.sessions === 1 ? '' : 's'} + +
+

{host.harness}

+

{host.description}

+
+ ))} +
+
+ +
+

Sessions

+ {(sessions.data ?? []).length === 0 ? ( +

+ No sessions yet. Start one to watch it live. +

+ ) : ( + + )} +
+
+ ) +} diff --git a/examples/agent-dashboard/src/routes/sessions.$threadId.tsx b/examples/agent-dashboard/src/routes/sessions.$threadId.tsx new file mode 100644 index 0000000000..0821d5203a --- /dev/null +++ b/examples/agent-dashboard/src/routes/sessions.$threadId.tsx @@ -0,0 +1,270 @@ +import { createFileRoute } from '@tanstack/react-router' +import { eq, useLiveQuery } from '@tanstack/react-db' +import { useEffect, useState } from 'react' +import { + approvals, + messages, + sessions, + spend, + toolCalls, +} from '@/db/collections' +import { + ensureSession, + resolveApproval, + sendPrompt, +} from '@/lib/session-controller' +import type { ApprovalRow, MessageRow, ToolCallRow } from '@/db/collections' + +export const Route = createFileRoute('/sessions/$threadId')({ + component: SessionDetail, +}) + +function SessionDetail() { + const { threadId } = Route.useParams() + const [input, setInput] = useState('') + + useEffect(() => { + ensureSession(threadId) + }, [threadId]) + + const { data: msgs = [] } = useLiveQuery( + (q) => q.from({ m: messages }).where(({ m }) => eq(m.threadId, threadId)), + [threadId], + ) + const { data: tools = [] } = useLiveQuery( + (q) => q.from({ t: toolCalls }).where(({ t }) => eq(t.threadId, threadId)), + [threadId], + ) + const { data: apprs = [] } = useLiveQuery( + (q) => q.from({ a: approvals }).where(({ a }) => eq(a.threadId, threadId)), + [threadId], + ) + const { data: spendRows = [] } = useLiveQuery( + (q) => q.from({ s: spend }).where(({ s }) => eq(s.threadId, threadId)), + [threadId], + ) + const { data: sess = [] } = useLiveQuery( + (q) => q.from({ s: sessions }).where(({ s }) => eq(s.threadId, threadId)), + [threadId], + ) + + const status = (sess as Array<{ status: string }>)[0]?.status ?? 'idle' + const tokens = (spendRows as Array<{ totalTokens: number }>)[0]?.totalTokens ?? 0 + const pending = (apprs as Array).filter( + (a) => a.status === 'pending', + ) + + const timeline = [ + ...(msgs as Array).map((m) => ({ kind: 'message' as const, at: m.createdAt, m })), + ...(tools as Array).map((t) => ({ kind: 'tool' as const, at: t.createdAt, t })), + ].sort((a, b) => a.at - b.at) + + const send = async () => { + const text = input.trim() + if (!text) return + setInput('') + await sendPrompt(threadId, text) + } + + return ( +
+
+ + ← hosts + +

{threadId}

+ + {status.replace('_', ' ')} + + + {tokens.toLocaleString()} tokens + +
+ + {pending.map((approval) => ( + ).find( + (t) => t.id === approval.toolCallId, + )} + threadId={threadId} + /> + ))} + +
+ {timeline.length === 0 && ( +

+ No activity yet. Send a prompt or start the triage demo below. +

+ )} + {timeline.map((entry) => + entry.kind === 'message' ? ( + + ) : ( + + ), + )} +
+ +
+ + setInput(e.target.value)} + onKeyDown={(e) => e.key === 'Enter' && send()} + placeholder="Send a message…" + className="flex-1 rounded-md border border-white/15 bg-transparent px-3 py-2 text-sm outline-none focus:border-white/30" + /> + +
+
+ ) +} + +function MessageBubble({ message }: { message: MessageRow }) { + const isUser = message.role === 'user' + return ( +
+
+ {message.subagentRunId && ( + + subagent + + )} + {message.text || …} +
+
+ ) +} + +function ToolCard({ tool }: { tool: ToolCallRow }) { + return ( +
+
+ ⚙ {tool.name} + + {tool.status} + +
+ {tool.args &&
{tool.args}
} + {tool.result && ( +
→ {tool.result}
+ )} +
+ ) +} + +function ApprovalCard({ + approval, + tool, + threadId, +}: { + approval: ApprovalRow + tool?: ToolCallRow + threadId: string +}) { + const [editing, setEditing] = useState(false) + const [draft, setDraft] = useState(tool?.args ?? '{}') + const [busy, setBusy] = useState(false) + + const act = async (decision: 'approve' | 'deny', edited?: boolean) => { + setBusy(true) + let editedArgs: Record | undefined + if (edited) { + try { + editedArgs = JSON.parse(draft) + } catch { + setBusy(false) + return + } + } + await resolveApproval(threadId, approval.id, decision, editedArgs) + setBusy(false) + } + + return ( +
+
+ 🔔 + Approval required + {tool && ( + + {tool.name} + + )} +
+

{approval.message}

+ {editing ? ( +