diff --git a/.changeset/anthropic-thinking-replay.md b/.changeset/anthropic-thinking-replay.md new file mode 100644 index 0000000000..e5c18def33 --- /dev/null +++ b/.changeset/anthropic-thinking-replay.md @@ -0,0 +1,12 @@ +--- +'@tanstack/ai': patch +'@tanstack/ai-anthropic': patch +--- + +Send Claude's thinking and tool errors back the way Claude sent them. + +- A tool message with `error` now sends `tool_result.is_error: true`, so Claude sees that the tool failed. +- A `redacted_thinking` block is no longer dropped. It becomes a thinking part with `redacted: true`, an empty `content`, and the encrypted data in `signature`. The flag survives the stream, the UI messages, the wire, and stored threads, and the next request sends the block back as `{ type: 'redacted_thinking', data }`. +- On the AG-UI wire, a redacted block is its own reasoning message. Its id starts with `redacted_thinking-`, and the `REASONING_ENCRYPTED_VALUE` event's `entityId` points to that id. An AG-UI client keeps message ids, so it sends the block back as redacted data, not as a signature. +- A thinking block's signature now names its reasoning message in `entityId`, not the step. An AG-UI client attaches the signature to that message, so it can send it back. +- `ThinkingPart` and `ModelMessage['thinking']` have the new optional `redacted` field. diff --git a/docs/chat/thinking-content.md b/docs/chat/thinking-content.md index b610f08b8e..b5af832a01 100644 --- a/docs/chat/thinking-content.md +++ b/docs/chat/thinking-content.md @@ -28,11 +28,14 @@ interface ThinkingPart { content: string; stepId?: string; signature?: string; + redacted?: boolean; } ``` The `ThinkingPart` appears in `UIMessage.parts` alongside `TextPart` and `ToolCallPart` entries. As reasoning tokens arrive, its `content` accumulates token by token. +Claude can also send a redacted thinking block. It is encrypted, so it has no text. It arrives as a `ThinkingPart` with `redacted: true`, an empty `content`, and the encrypted data in `signature`. Keep the part in your stored messages: the next request sends it back to Claude unchanged. In a UI, show a short placeholder such as "Thinking hidden" instead of the empty text. + ## Enabling Thinking How you enable thinking depends on the provider. diff --git a/docs/config.json b/docs/config.json index 0bd1fa3672..bb512240ac 100644 --- a/docs/config.json +++ b/docs/config.json @@ -177,7 +177,7 @@ "label": "Thinking & Reasoning", "to": "chat/thinking-content", "addedAt": "2026-04-15", - "updatedAt": "2026-08-21" + "updatedAt": "2026-09-30" }, { "label": "Message Queue", diff --git a/packages/ai-anthropic/src/adapters/text.ts b/packages/ai-anthropic/src/adapters/text.ts index b84a7a65c4..9682c224ce 100644 --- a/packages/ai-anthropic/src/adapters/text.ts +++ b/packages/ai-anthropic/src/adapters/text.ts @@ -4,7 +4,10 @@ import { isFileSource, normalizeSystemPrompts, } from '@tanstack/ai' -import { toRunErrorRawEvent } from '@tanstack/ai/adapter-internals' +import { + REDACTED_THINKING_ID_PREFIX, + toRunErrorRawEvent, +} from '@tanstack/ai/adapter-internals' import { BaseTextAdapter } from '@tanstack/ai/adapters' import { convertToolsToProviderFormat } from '../tools/tool-converter' import { getAnthropicProviderToolKind } from '../tools/anthropic-provider-tool' @@ -824,6 +827,7 @@ export class AnthropicTextAdapter< : typeof toolContent === 'string' ? toolContent : '', + ...(message.error !== undefined && { is_error: true }), }, ], }) @@ -952,6 +956,13 @@ export class AnthropicTextAdapter< for (const thinking of thinkingParts) { if (!thinking.signature) continue + if (thinking.redacted) { + contentBlocks.push({ + type: 'redacted_thinking', + data: thinking.signature, + }) + continue + } const block: ThinkingBlockParam = { type: 'thinking', thinking: thinking.content, @@ -1214,6 +1225,63 @@ export class AnthropicTextAdapter< timestamp: Date.now(), stepType: 'thinking', } + } else if (event.content_block.type === 'redacted_thinking') { + // Encrypted thinking: no text, and its data must go back to + // Anthropic unchanged. It gets its own reasoning message, and the + // encrypted value names that message. The id prefix marks the + // data as a redacted block, not a signature. The step reuses the + // id so the stream processor sees the prefix too. + const redactedId = `${REDACTED_THINKING_ID_PREFIX}${genId()}` + yield { + type: EventType.REASONING_START, + messageId: redactedId, + model, + timestamp: Date.now(), + } + yield { + type: EventType.REASONING_MESSAGE_START, + messageId: redactedId, + role: 'reasoning' as const, + model, + timestamp: Date.now(), + } + yield { + type: EventType.STEP_STARTED, + stepName: redactedId, + stepId: redactedId, + model, + timestamp: Date.now(), + stepType: 'thinking', + } + yield { + type: EventType.STEP_FINISHED, + stepName: redactedId, + stepId: redactedId, + model, + timestamp: Date.now(), + delta: '', + content: '', + } + yield { + type: EventType.REASONING_ENCRYPTED_VALUE, + subtype: 'message' as const, + entityId: redactedId, + encryptedValue: event.content_block.data, + model, + timestamp: Date.now(), + } + yield { + type: EventType.REASONING_MESSAGE_END, + messageId: redactedId, + model, + timestamp: Date.now(), + } + yield { + type: EventType.REASONING_END, + messageId: redactedId, + model, + timestamp: Date.now(), + } } } else if (event.type === 'content_block_delta') { if (event.delta.type === 'text_delta') { @@ -1332,8 +1400,10 @@ export class AnthropicTextAdapter< } } else if (event.type === 'content_block_stop') { if (currentBlockType === 'thinking') { - // Emit signature so it can be replayed in multi-turn context - if (accumulatedSignature && stepId) { + // Emit signature so it can be replayed in multi-turn context. + // It belongs to the reasoning message, not the step, so an AG-UI + // client finds that message by `entityId`. + if (accumulatedSignature && stepId && reasoningMessageId) { yield { type: EventType.STEP_FINISHED, stepName: stepId, @@ -1342,7 +1412,14 @@ export class AnthropicTextAdapter< timestamp: Date.now(), delta: '', content: accumulatedThinking, - signature: accumulatedSignature, + } + yield { + type: EventType.REASONING_ENCRYPTED_VALUE, + subtype: 'message' as const, + entityId: reasoningMessageId, + encryptedValue: accumulatedSignature, + model, + timestamp: Date.now(), } } } else if (currentBlockType === 'tool_use') { diff --git a/packages/ai-anthropic/tests/thinking-replay.test.ts b/packages/ai-anthropic/tests/thinking-replay.test.ts new file mode 100644 index 0000000000..282e5b5a6e --- /dev/null +++ b/packages/ai-anthropic/tests/thinking-replay.test.ts @@ -0,0 +1,272 @@ +import { beforeEach, describe, expect, it, vi } from 'vitest' +import { chat, StreamProcessor } from '@tanstack/ai' +import { z } from 'zod' +import { AnthropicTextAdapter } from '../src/adapters/text' +import type { ModelMessage, StreamChunk, Tool } from '@tanstack/ai' + +const mocks = vi.hoisted(() => { + const betaMessagesCreate = vi.fn() + const client = { + beta: { messages: { create: betaMessagesCreate } }, + messages: { create: vi.fn() }, + } + return { betaMessagesCreate, client } +}) + +vi.mock('@anthropic-ai/sdk', () => { + class MockAnthropic { + beta = mocks.client.beta + messages = mocks.client.messages + constructor(_: { apiKey: string }) {} + } + return { default: MockAnthropic } +}) + +type Block = { type: string } & Record + +const adapter = () => + new AnthropicTextAdapter({ apiKey: 'test-key' }, 'claude-opus-4-1') + +const weather: Tool = { + name: 'lookup_weather', + description: 'Return the weather for a city', + inputSchema: z.object({ location: z.string() }), + execute: () => 'sunny', +} + +/** A stream of raw Anthropic events. */ +function stream(events: Array>) { + return (async function* () { + for (const event of events) yield event + })() +} + +const redactedBlock = [ + { + type: 'content_block_start', + index: 0, + content_block: { type: 'redacted_thinking', data: 'opaque-1' }, + }, + { type: 'content_block_stop', index: 0 }, +] + +function textEvents(index: number, text: string) { + return [ + { + type: 'content_block_start', + index, + content_block: { type: 'text', text: '' }, + }, + { + type: 'content_block_delta', + index, + delta: { type: 'text_delta', text }, + }, + { type: 'content_block_stop', index }, + ] +} + +function end(stopReason: string) { + return [ + { + type: 'message_delta', + delta: { stop_reason: stopReason }, + usage: { output_tokens: 5 }, + }, + { type: 'message_stop' }, + ] +} + +/** The content blocks of the last assistant message in a request. */ +function assistantBlocks(call: number) { + const [payload] = mocks.betaMessagesCreate.mock.calls[call]! + const messages: Array<{ role: string; content: unknown }> = payload.messages + const assistant = messages.filter((m) => m.role === 'assistant').at(-1) + return Array.isArray(assistant?.content) + ? (assistant.content as Array) + : [] +} + +async function drain(iterable: AsyncIterable) { + for await (const _ of iterable) { + // consume + } +} + +describe('Anthropic replay', () => { + beforeEach(() => { + mocks.betaMessagesCreate.mockReset() + }) + + it('marks a failed tool result with is_error', async () => { + mocks.betaMessagesCreate.mockResolvedValueOnce( + stream([...textEvents(0, 'Sorry.'), ...end('end_turn')]), + ) + const messages: Array = [ + { role: 'user', content: 'Weather in Atlantis?' }, + { + role: 'assistant', + content: null, + toolCalls: [ + { + id: 'call-1', + type: 'function', + function: { + name: 'lookup_weather', + arguments: '{"location":"Atlantis"}', + }, + }, + ], + }, + { + role: 'tool', + toolCallId: 'call-1', + content: '{"error":"City not found"}', + error: 'City not found', + }, + ] + + await drain(chat({ adapter: adapter(), messages })) + + const [payload] = mocks.betaMessagesCreate.mock.calls[0]! + const blocks: Array = payload.messages.flatMap( + (m: { content: unknown }) => (Array.isArray(m.content) ? m.content : []), + ) + expect(blocks.find((b) => b.type === 'tool_result')).toMatchObject({ + tool_use_id: 'call-1', + is_error: true, + }) + }) + + it('sends a redacted thinking block back in the same run', async () => { + mocks.betaMessagesCreate + .mockResolvedValueOnce( + stream([ + ...redactedBlock, + { + type: 'content_block_start', + index: 1, + content_block: { + type: 'tool_use', + id: 'call-1', + name: 'lookup_weather', + input: {}, + }, + }, + { + type: 'content_block_delta', + index: 1, + delta: { + type: 'input_json_delta', + partial_json: '{"location":"Berlin"}', + }, + }, + { type: 'content_block_stop', index: 1 }, + ...end('tool_use'), + ]), + ) + .mockResolvedValueOnce( + stream([...textEvents(0, 'It is sunny.'), ...end('end_turn')]), + ) + + await drain( + chat({ + adapter: adapter(), + messages: [{ role: 'user', content: 'Weather in Berlin?' }], + tools: [weather], + }), + ) + + expect(mocks.betaMessagesCreate).toHaveBeenCalledTimes(2) + expect(assistantBlocks(1).map((b) => b.type)).toEqual([ + 'redacted_thinking', + 'tool_use', + ]) + expect(assistantBlocks(1)[0]).toEqual({ + type: 'redacted_thinking', + data: 'opaque-1', + }) + }) + + it('ties each encrypted value to its reasoning message by id', async () => { + mocks.betaMessagesCreate.mockResolvedValueOnce( + stream([ + // A signed block with omitted text, the default on newer models. + { + type: 'content_block_start', + index: 0, + content_block: { type: 'thinking', thinking: '' }, + }, + { + type: 'content_block_delta', + index: 0, + delta: { type: 'signature_delta', signature: 'sig-1' }, + }, + { type: 'content_block_stop', index: 0 }, + { + type: 'content_block_start', + index: 1, + content_block: { type: 'redacted_thinking', data: 'opaque-1' }, + }, + { type: 'content_block_stop', index: 1 }, + ...textEvents(2, 'Hello.'), + ...end('end_turn'), + ]), + ) + const chunks: Array = [] + for await (const chunk of chat({ + adapter: adapter(), + messages: [{ role: 'user', content: 'Hi' }], + })) { + chunks.push(chunk) + } + + const reasoningIds = chunks.flatMap((c) => + c.type === 'REASONING_MESSAGE_START' ? [c.messageId] : [], + ) + const values = chunks.flatMap((c) => + c.type === 'REASONING_ENCRYPTED_VALUE' ? [c] : [], + ) + expect(values.map((v) => v.encryptedValue)).toEqual(['sig-1', 'opaque-1']) + // An AG-UI client attaches each value to the message with this id. + expect(values.map((v) => v.entityId)).toEqual(reasoningIds) + expect(reasoningIds[0]).not.toMatch(/^redacted_thinking-/) + expect(reasoningIds[1]).toMatch(/^redacted_thinking-/) + }) + + it('sends a redacted thinking block back on the next turn', async () => { + mocks.betaMessagesCreate + .mockResolvedValueOnce( + stream([ + ...redactedBlock, + ...textEvents(1, 'Hello.'), + ...end('end_turn'), + ]), + ) + .mockResolvedValueOnce( + stream([...textEvents(0, 'Again.'), ...end('end_turn')]), + ) + const first: Array = [{ role: 'user', content: 'Hi' }] + const processor = new StreamProcessor() + processor.addUserMessage('Hi') + for await (const chunk of chat({ adapter: adapter(), messages: first })) { + processor.processChunk(chunk) + } + processor.finalizeStream() + + await drain( + chat({ + adapter: adapter(), + messages: [ + ...processor.getMessages(), + { role: 'user', content: 'Say it again.' }, + ], + }), + ) + + expect(assistantBlocks(1)).toEqual([ + { type: 'redacted_thinking', data: 'opaque-1' }, + { type: 'text', text: 'Hello.' }, + ]) + }) +}) diff --git a/packages/ai/src/activities/chat/index.ts b/packages/ai/src/activities/chat/index.ts index 5b52331b0b..e4ce88e7a7 100644 --- a/packages/ai/src/activities/chat/index.ts +++ b/packages/ai/src/activities/chat/index.ts @@ -41,6 +41,7 @@ import { import { subagentHostMessageId } from '../../utilities/subagent-wire' import { withDurabilityBatchHint } from '../../utilities/durability-batch' import { normalizeStreamChunk } from '../../utilities/normalize-stream-chunk' +import { isRedactedThinkingId } from '../../utilities/reasoning-encrypted-value' import { restorePublicUsage } from '../../utilities/restore-inbound-chunk' import type { AdapterYieldChunk } from '../../utilities/adapter-yield-chunk' import { @@ -855,8 +856,7 @@ class TextEngine< private currentMessageCreatedAt: Date | null = null private streamIdentityCaptured = false private accumulatedContent = '' - private accumulatedThinking: Array<{ content: string; signature?: string }> = - [] + private accumulatedThinking: NonNullable = [] /** * Arrival order of this iteration's thinking steps, text and tool calls. * A ModelMessage keeps `thinking` apart from `content`/`toolCalls`, so a @@ -868,6 +868,7 @@ class TextEngine< private turnParts: Array | null = [] private currentThinkingContent = '' private currentThinkingSignature = '' + private currentThinkingRedacted = false private eventOptions?: Record | undefined private eventToolNames?: Array private finishedEvent: RunFinishedEvent | null = null @@ -1526,6 +1527,7 @@ class TextEngine< this.turnParts = [] this.currentThinkingContent = '' this.currentThinkingSignature = '' + this.currentThinkingRedacted = false this.finishedEvent = null this.streamedToolErrorResults.clear() @@ -1994,6 +1996,7 @@ class TextEngine< ...(this.currentThinkingSignature && { signature: this.currentThinkingSignature, }), + ...(this.currentThinkingRedacted && { redacted: true }), }) if (this.turnParts) { const placeholder = [...this.turnParts] @@ -2011,6 +2014,7 @@ class TextEngine< } this.currentThinkingContent = '' this.currentThinkingSignature = '' + this.currentThinkingRedacted = false } } @@ -2038,6 +2042,7 @@ class TextEngine< if (typeof chunk.signature === 'string' && chunk.signature !== '') { this.noteThinkingStepPosition() this.currentThinkingSignature = chunk.signature + this.currentThinkingRedacted = isRedactedThinkingId(chunk.stepId) } } @@ -2067,6 +2072,7 @@ class TextEngine< } this.noteThinkingStepPosition() this.currentThinkingSignature = chunk.encryptedValue + this.currentThinkingRedacted = isRedactedThinkingId(chunk.entityId) } /** @@ -2583,7 +2589,7 @@ class TextEngine< ), ) type Segment = { - thinking: Array<{ content: string; signature?: string }> + thinking: NonNullable text: string callIds: Array } diff --git a/packages/ai/src/activities/chat/messages.ts b/packages/ai/src/activities/chat/messages.ts index 4a4bac88e1..0c4759e70d 100644 --- a/packages/ai/src/activities/chat/messages.ts +++ b/packages/ai/src/activities/chat/messages.ts @@ -11,6 +11,7 @@ import { tanstackMetadata, withTanstackMetadata, } from '../../utilities/merge-metadata' +import { isRedactedThinkingId } from '../../utilities/reasoning-encrypted-value' import { splitSubagentWire, subagentWireText, @@ -72,6 +73,13 @@ function encryptedValueFrom(value: object): string | undefined { return nonEmptyString(tanstackMetadata(value)?.signature) } +/** `{ redacted: true }` when a reasoning message's id marks a redacted block. */ +function redactedFrom(value: object) { + return 'id' in value && isRedactedThinkingId(value.id) + ? { redacted: true } + : {} +} + function toolCallFromWire(toolCall: ToolCall, bag: unknown): ToolCall { const fromBag = bag != null && typeof bag === 'object' && !Array.isArray(bag) @@ -246,7 +254,7 @@ function convertOwnMessages( } const modelMessages: Array = [] - let pendingThinking: Array<{ content: string; signature?: string }> = [] + let pendingThinking: NonNullable = [] for (const msg of messages) { if ('parts' in msg) { modelMessages.push(...uiMessageToModelMessages(msg)) @@ -273,6 +281,7 @@ function convertOwnMessages( pendingThinking.push({ content: typeof content === 'string' ? content : '', ...(signature !== undefined ? { signature } : {}), + ...redactedFrom(msg), }) } continue @@ -663,7 +672,7 @@ function buildAssistantMessages(uiMessage: UIMessage): Array { // shared UI id on each one so persistence can retain the original identity. const messageList: Array = [] let current = createSegment() - let pendingThinking: Array<{ content: string; signature?: string }> = [] + let pendingThinking: NonNullable = [] // Track emitted tool result IDs to avoid duplicates. // A tool call can have BOTH an explicit tool-result part AND an output @@ -764,6 +773,7 @@ function buildAssistantMessages(uiMessage: UIMessage): Array { pendingThinking.push({ content: part.content, ...(part.signature && { signature: part.signature }), + ...(part.redacted && { redacted: true }), }) } break @@ -898,6 +908,7 @@ export function modelMessageToUIMessage( type: 'thinking', content: thinking.content, ...(thinking.signature && { signature: thinking.signature }), + ...(thinking.redacted && { redacted: true }), }) } } @@ -1121,6 +1132,7 @@ export function aguiSnapshotMessageToUIMessage( type: 'thinking' as const, content, ...(signature !== undefined ? { signature } : {}), + ...redactedFrom(message), }, ] : [], diff --git a/packages/ai/src/activities/chat/stream/message-updaters.ts b/packages/ai/src/activities/chat/stream/message-updaters.ts index 0b999e556c..b7e936beab 100644 --- a/packages/ai/src/activities/chat/stream/message-updaters.ts +++ b/packages/ai/src/activities/chat/stream/message-updaters.ts @@ -5,6 +5,7 @@ * These are used by StreamProcessor to manage the message array. */ +import { isRedactedThinkingId } from '../../../utilities/reasoning-encrypted-value' import { parsePartialJSON } from './json-parser' import type { ContentPart, @@ -483,12 +484,16 @@ export function updateThinkingPart( // not carry one; losing it would strip the provider's encrypted reasoning // from a message that is about to be sent back. const nextSignature = signature ?? adopted?.signature + // A hydrated part has no stepId to carry the redacted marker, so it keeps + // its own flag. + const redacted = isRedactedThinkingId(stepId) || adopted?.redacted === true const thinkingPart: ThinkingPart = { type: 'thinking', content, stepId, ...(nextSignature && { signature: nextSignature }), + ...(redacted && { redacted: true }), } if (thinkingPartIndex >= 0) { diff --git a/packages/ai/src/adapter-internals.ts b/packages/ai/src/adapter-internals.ts index fb274dfce6..7005b72584 100644 --- a/packages/ai/src/adapter-internals.ts +++ b/packages/ai/src/adapter-internals.ts @@ -59,3 +59,4 @@ export { } from './utilities/structured-output-events' export { tanstackMetadata } from './utilities/merge-metadata' export { isSpecTopLevelKey } from './utilities/spec-event-keys' +export { REDACTED_THINKING_ID_PREFIX } from './utilities/reasoning-encrypted-value' diff --git a/packages/ai/src/types.ts b/packages/ai/src/types.ts index fd96f5c78a..4fbd5b89f1 100644 --- a/packages/ai/src/types.ts +++ b/packages/ai/src/types.ts @@ -381,7 +381,12 @@ export interface ModelMessage< name?: string toolCalls?: Array toolCallId?: string - thinking?: Array<{ content: string; signature?: string }> + /** + * Signed thinking to send back to the provider. `redacted: true` marks a + * block the provider encrypted: `content` is empty and `signature` holds its + * opaque data. See `ThinkingPart.signature` for the planned rename. + */ + thinking?: Array<{ content: string; signature?: string; redacted?: boolean }> /** Error reported by an AG-UI tool message. */ error?: string /** Optional AG-UI message metadata. TanStack-owned fields live under `tanstack`. */ @@ -465,7 +470,20 @@ export interface ThinkingPart { type: 'thinking' content: string stepId?: string + /** + * The provider's opaque reasoning artefact, sent back unchanged: an + * Anthropic signature, Anthropic redacted data, or OpenAI encrypted content. + * TODO(#1581): rename to `encryptedValue` to match AG-UI's `ReasoningMessage`. + * Renaming breaks stored messages, so it needs a read shim for `signature`. + */ signature?: string + /** + * The provider encrypted this thinking block (Anthropic `redacted_thinking`). + * `content` is empty, and `signature` holds the opaque data that goes back + * to the provider unchanged. On the AG-UI wire, the reasoning message id + * starts with `redacted_thinking-` instead. + */ + redacted?: boolean } /** diff --git a/packages/ai/src/utilities/ag-ui-wire.ts b/packages/ai/src/utilities/ag-ui-wire.ts index da2aac2bcc..0edeb0c37f 100644 --- a/packages/ai/src/utilities/ag-ui-wire.ts +++ b/packages/ai/src/utilities/ag-ui-wire.ts @@ -21,6 +21,7 @@ import type { import type { MetadataRecord } from './merge-metadata' import { tanstackMetadata } from './merge-metadata' import { isProviderExecutedToolCall } from './provider-executed' +import { REDACTED_THINKING_ID_PREFIX } from './reasoning-encrypted-value' import { normalizeToolResult } from './tool-result' import { wireSubagentInfo, wireSubagentRunId } from './subagent-wire' import type { SubagentWireInfo } from './subagent-wire' @@ -614,7 +615,10 @@ function collectToolCalls( } function deriveReasoningId(messageId: string, part: MessagePart): string { - return `${messageId}-reasoning-${(part as { id?: string }).id ?? hashContent((part as { content: string }).content)}` + const id = `${messageId}-reasoning-${(part as { id?: string }).id ?? hashContent((part as { content: string }).content)}` + return part.type === 'thinking' && part.redacted + ? `${REDACTED_THINKING_ID_PREFIX}${id}` + : id } function deriveToolMessageId(toolCallId: string): string { diff --git a/packages/ai/src/utilities/reasoning-encrypted-value.ts b/packages/ai/src/utilities/reasoning-encrypted-value.ts index 1236a30234..4b2de205b8 100644 --- a/packages/ai/src/utilities/reasoning-encrypted-value.ts +++ b/packages/ai/src/utilities/reasoning-encrypted-value.ts @@ -16,3 +16,15 @@ export function reasoningEncryptedValue(opts: { ...(opts.timestamp !== undefined ? { timestamp: opts.timestamp } : {}), } } + +/** + * Id prefix for the reasoning message of a redacted thinking block (Anthropic + * `redacted_thinking`). AG-UI's `encryptedValue` does not say what kind of + * bytes it holds (ag-ui-protocol/ag-ui#2884), so the id that `entityId` points + * to says it. AG-UI clients keep message ids, but they drop event metadata. + */ +export const REDACTED_THINKING_ID_PREFIX = 'redacted_thinking-' + +export function isRedactedThinkingId(id: unknown): boolean { + return typeof id === 'string' && id.startsWith(REDACTED_THINKING_ID_PREFIX) +} diff --git a/packages/ai/tests/redacted-thinking.test.ts b/packages/ai/tests/redacted-thinking.test.ts new file mode 100644 index 0000000000..2794b9410d --- /dev/null +++ b/packages/ai/tests/redacted-thinking.test.ts @@ -0,0 +1,97 @@ +import { describe, expect, it } from 'vitest' +import { + aguiSnapshotMessageToUIMessage, + convertMessagesToModelMessages, + modelMessagesToUIMessages, +} from '../src/activities/chat/messages' +import { uiMessagesToWire } from '../src/utilities/ag-ui-wire' +import { chatParamsFromRequestBody } from '../src/utilities/chat-params' +import type { ModelMessage } from '../src/types' + +const stored: Array = [ + { role: 'user', content: 'Hi' }, + { + role: 'assistant', + content: 'Hello.', + thinking: [ + { content: '', signature: 'opaque-1', redacted: true }, + { content: 'I greet back.', signature: 'sig-2' }, + ], + }, +] + +function thinkingOf(messages: Array) { + return messages.find((message) => message.role === 'assistant')?.thinking +} + +describe('redacted thinking', () => { + it('survives a stored thread that is loaded into the UI and sent back', () => { + const ui = modelMessagesToUIMessages(stored) + + expect(thinkingOf(convertMessagesToModelMessages(ui))).toEqual([ + { content: '', signature: 'opaque-1', redacted: true }, + { content: 'I greet back.', signature: 'sig-2' }, + ]) + }) + + it('survives the wire from the client to the server', async () => { + const wire = uiMessagesToWire(modelMessagesToUIMessages(stored)) + // `JSON.parse(JSON.stringify(...))` stands in for the HTTP hop. + const params = await chatParamsFromRequestBody({ + threadId: 'thread-1', + runId: 'run-1', + messages: JSON.parse(JSON.stringify(wire)), + tools: [], + context: [], + }) + + expect(thinkingOf(convertMessagesToModelMessages(params.messages))).toEqual( + [ + { content: '', signature: 'opaque-1', redacted: true }, + { content: 'I greet back.', signature: 'sig-2' }, + ], + ) + }) + + it('reads the kind from the reasoning message id, without metadata', async () => { + // What any AG-UI client sends back: the spec fields only. The second + // message is a signed block with omitted text, so `content` is empty too. + const params = await chatParamsFromRequestBody({ + threadId: 'thread-1', + runId: 'run-1', + messages: [ + { id: 'user-1', role: 'user', content: 'Hi' }, + { + id: 'redacted_thinking-r1', + role: 'reasoning', + content: '', + encryptedValue: 'opaque-1', + }, + { id: 'r2', role: 'reasoning', content: '', encryptedValue: 'sig-2' }, + { id: 'assistant-1', role: 'assistant', content: 'Hello.' }, + ], + tools: [], + context: [], + }) + + expect(thinkingOf(convertMessagesToModelMessages(params.messages))).toEqual( + [ + { content: '', signature: 'opaque-1', redacted: true }, + { content: '', signature: 'sig-2' }, + ], + ) + }) + + it('survives an interrupt snapshot that the client loads', () => { + const wire = uiMessagesToWire(modelMessagesToUIMessages(stored)) + + const parts = wire + .filter((message) => message.role === 'reasoning') + .flatMap((message) => aguiSnapshotMessageToUIMessage(message).parts) + + expect(parts).toEqual([ + { type: 'thinking', content: '', signature: 'opaque-1', redacted: true }, + { type: 'thinking', content: 'I greet back.', signature: 'sig-2' }, + ]) + }) +}) diff --git a/testing/e2e/src/routeTree.gen.ts b/testing/e2e/src/routeTree.gen.ts index 013e1601a1..439b88d192 100644 --- a/testing/e2e/src/routeTree.gen.ts +++ b/testing/e2e/src/routeTree.gen.ts @@ -123,6 +123,7 @@ import { Route as ApiArktypeToolWireRouteImport } from './routes/api.arktype-too import { Route as ApiAnthropicThinkingOrderWireRouteImport } from './routes/api.anthropic-thinking-order-wire' import { Route as ApiAnthropicStructuredUsageRouteImport } from './routes/api.anthropic-structured-usage' import { Route as ApiAnthropicSkillsWireRouteImport } from './routes/api.anthropic-skills-wire' +import { Route as ApiAnthropicRedactedThinkingWireRouteImport } from './routes/api.anthropic-redacted-thinking-wire' import { Route as ApiAnthropicOpus5CombinedWireRouteImport } from './routes/api.anthropic-opus-5-combined-wire' import { Route as ApiAnthropicMultiTurnStructuredWireRouteImport } from './routes/api.anthropic-multi-turn-structured-wire' import { Route as ApiAnthropicBugTestRouteImport } from './routes/api.anthropic-bug-test' @@ -727,6 +728,12 @@ const ApiAnthropicSkillsWireRoute = ApiAnthropicSkillsWireRouteImport.update({ path: '/api/anthropic-skills-wire', getParentRoute: () => rootRouteImport, } as any) +const ApiAnthropicRedactedThinkingWireRoute = + ApiAnthropicRedactedThinkingWireRouteImport.update({ + id: '/api/anthropic-redacted-thinking-wire', + path: '/api/anthropic-redacted-thinking-wire', + getParentRoute: () => rootRouteImport, + } as any) const ApiAnthropicOpus5CombinedWireRoute = ApiAnthropicOpus5CombinedWireRouteImport.update({ id: '/api/anthropic-opus-5-combined-wire', @@ -811,6 +818,7 @@ export interface FileRoutesByFullPath { '/api/anthropic-bug-test': typeof ApiAnthropicBugTestRoute '/api/anthropic-multi-turn-structured-wire': typeof ApiAnthropicMultiTurnStructuredWireRoute '/api/anthropic-opus-5-combined-wire': typeof ApiAnthropicOpus5CombinedWireRoute + '/api/anthropic-redacted-thinking-wire': typeof ApiAnthropicRedactedThinkingWireRoute '/api/anthropic-skills-wire': typeof ApiAnthropicSkillsWireRoute '/api/anthropic-structured-usage': typeof ApiAnthropicStructuredUsageRoute '/api/anthropic-thinking-order-wire': typeof ApiAnthropicThinkingOrderWireRoute @@ -936,6 +944,7 @@ export interface FileRoutesByTo { '/api/anthropic-bug-test': typeof ApiAnthropicBugTestRoute '/api/anthropic-multi-turn-structured-wire': typeof ApiAnthropicMultiTurnStructuredWireRoute '/api/anthropic-opus-5-combined-wire': typeof ApiAnthropicOpus5CombinedWireRoute + '/api/anthropic-redacted-thinking-wire': typeof ApiAnthropicRedactedThinkingWireRoute '/api/anthropic-skills-wire': typeof ApiAnthropicSkillsWireRoute '/api/anthropic-structured-usage': typeof ApiAnthropicStructuredUsageRoute '/api/anthropic-thinking-order-wire': typeof ApiAnthropicThinkingOrderWireRoute @@ -1062,6 +1071,7 @@ export interface FileRoutesById { '/api/anthropic-bug-test': typeof ApiAnthropicBugTestRoute '/api/anthropic-multi-turn-structured-wire': typeof ApiAnthropicMultiTurnStructuredWireRoute '/api/anthropic-opus-5-combined-wire': typeof ApiAnthropicOpus5CombinedWireRoute + '/api/anthropic-redacted-thinking-wire': typeof ApiAnthropicRedactedThinkingWireRoute '/api/anthropic-skills-wire': typeof ApiAnthropicSkillsWireRoute '/api/anthropic-structured-usage': typeof ApiAnthropicStructuredUsageRoute '/api/anthropic-thinking-order-wire': typeof ApiAnthropicThinkingOrderWireRoute @@ -1189,6 +1199,7 @@ export interface FileRouteTypes { | '/api/anthropic-bug-test' | '/api/anthropic-multi-turn-structured-wire' | '/api/anthropic-opus-5-combined-wire' + | '/api/anthropic-redacted-thinking-wire' | '/api/anthropic-skills-wire' | '/api/anthropic-structured-usage' | '/api/anthropic-thinking-order-wire' @@ -1314,6 +1325,7 @@ export interface FileRouteTypes { | '/api/anthropic-bug-test' | '/api/anthropic-multi-turn-structured-wire' | '/api/anthropic-opus-5-combined-wire' + | '/api/anthropic-redacted-thinking-wire' | '/api/anthropic-skills-wire' | '/api/anthropic-structured-usage' | '/api/anthropic-thinking-order-wire' @@ -1439,6 +1451,7 @@ export interface FileRouteTypes { | '/api/anthropic-bug-test' | '/api/anthropic-multi-turn-structured-wire' | '/api/anthropic-opus-5-combined-wire' + | '/api/anthropic-redacted-thinking-wire' | '/api/anthropic-skills-wire' | '/api/anthropic-structured-usage' | '/api/anthropic-thinking-order-wire' @@ -1565,6 +1578,7 @@ export interface RootRouteChildren { ApiAnthropicBugTestRoute: typeof ApiAnthropicBugTestRoute ApiAnthropicMultiTurnStructuredWireRoute: typeof ApiAnthropicMultiTurnStructuredWireRoute ApiAnthropicOpus5CombinedWireRoute: typeof ApiAnthropicOpus5CombinedWireRoute + ApiAnthropicRedactedThinkingWireRoute: typeof ApiAnthropicRedactedThinkingWireRoute ApiAnthropicSkillsWireRoute: typeof ApiAnthropicSkillsWireRoute ApiAnthropicStructuredUsageRoute: typeof ApiAnthropicStructuredUsageRoute ApiAnthropicThinkingOrderWireRoute: typeof ApiAnthropicThinkingOrderWireRoute @@ -2450,6 +2464,13 @@ declare module '@tanstack/react-router' { preLoaderRoute: typeof ApiAnthropicSkillsWireRouteImport parentRoute: typeof rootRouteImport } + '/api/anthropic-redacted-thinking-wire': { + id: '/api/anthropic-redacted-thinking-wire' + path: '/api/anthropic-redacted-thinking-wire' + fullPath: '/api/anthropic-redacted-thinking-wire' + preLoaderRoute: typeof ApiAnthropicRedactedThinkingWireRouteImport + parentRoute: typeof rootRouteImport + } '/api/anthropic-opus-5-combined-wire': { id: '/api/anthropic-opus-5-combined-wire' path: '/api/anthropic-opus-5-combined-wire' @@ -2611,6 +2632,7 @@ const rootRouteChildren: RootRouteChildren = { ApiAnthropicMultiTurnStructuredWireRoute: ApiAnthropicMultiTurnStructuredWireRoute, ApiAnthropicOpus5CombinedWireRoute: ApiAnthropicOpus5CombinedWireRoute, + ApiAnthropicRedactedThinkingWireRoute: ApiAnthropicRedactedThinkingWireRoute, ApiAnthropicSkillsWireRoute: ApiAnthropicSkillsWireRoute, ApiAnthropicStructuredUsageRoute: ApiAnthropicStructuredUsageRoute, ApiAnthropicThinkingOrderWireRoute: ApiAnthropicThinkingOrderWireRoute, diff --git a/testing/e2e/src/routes/api.anthropic-redacted-thinking-wire.ts b/testing/e2e/src/routes/api.anthropic-redacted-thinking-wire.ts new file mode 100644 index 0000000000..d98acda142 --- /dev/null +++ b/testing/e2e/src/routes/api.anthropic-redacted-thinking-wire.ts @@ -0,0 +1,241 @@ +import { createFileRoute } from '@tanstack/react-router' +import { + chat, + chatParamsFromRequestBody, + createChatOptions, +} from '@tanstack/ai' +import { StreamProcessor, uiMessagesToWire } from '@tanstack/ai/client' +import { createAnthropicChat } from '@tanstack/ai-anthropic' + +const DUMMY_KEY = 'sk-ant-e2e-test-dummy-key' + +type SseEvent = Record & { type: string } + +const messageStart: SseEvent = { + type: 'message_start', + message: { + id: 'msg_redacted', + type: 'message', + role: 'assistant', + content: [], + model: 'claude-sonnet-4-5', + stop_reason: null, + stop_sequence: null, + usage: { input_tokens: 5, output_tokens: 0 }, + }, +} + +/** Claude sends a redacted block whole, on its start event. */ +function redacted(index: number, data: string) { + return [ + { + type: 'content_block_start', + index, + content_block: { type: 'redacted_thinking', data }, + }, + { type: 'content_block_stop', index }, + ] +} + +function text(index: number, value: string) { + return [ + { + type: 'content_block_start', + index, + content_block: { type: 'text', text: '' }, + }, + { + type: 'content_block_delta', + index, + delta: { type: 'text_delta', text: value }, + }, + { type: 'content_block_stop', index }, + ] +} + +function toolUse(index: number, name: string) { + return [ + { + type: 'content_block_start', + index, + content_block: { type: 'tool_use', id: 'toolu_1', name, input: {} }, + }, + { + type: 'content_block_delta', + index, + delta: { type: 'input_json_delta', partial_json: '{}' }, + }, + { type: 'content_block_stop', index }, + ] +} + +function end(stopReason: string) { + return [ + { + type: 'message_delta', + delta: { stop_reason: stopReason, stop_sequence: null }, + usage: { output_tokens: 20 }, + }, + { type: 'message_stop' }, + ] +} + +function sse(events: Array): Response { + return new Response( + events + .map( + (event) => `event: ${event.type}\ndata: ${JSON.stringify(event)}\n\n`, + ) + .join(''), + { status: 200, headers: { 'Content-Type': 'text/event-stream' } }, + ) +} + +type Body = { + messages: Array<{ role: string; content: unknown }> +} + +/** Answers each call with the next scripted stream and records its body. */ +function scriptedFetch(streams: Array>) { + const bodies: Array = [] + const fetchImpl: typeof fetch = async (input, init) => { + const req = input instanceof Request ? input : new Request(input, init) + bodies.push(JSON.parse(await req.text())) + return sse( + streams[bodies.length - 1] ?? [ + messageStart, + ...text(0, 'Done.'), + ...end('end_turn'), + ], + ) + } + return { bodies, fetchImpl } +} + +function blocks(body: Body | undefined, role: string) { + const message = body?.messages.filter((m) => m.role === role).at(-1) + return Array.isArray(message?.content) + ? (message.content as Array>) + : [] +} + +/** + * Three scripted runs, answered by a capturing `fetch`: + * + * - `nextTurn`: a redacted block and text, then a second turn through + * `StreamProcessor`, `uiMessagesToWire`, and `chatParamsFromRequestBody`, + * as a browser client and `/api/chat` do. + * - `sameRun`: a redacted block and a tool call. The run's second request + * continues after the tool. + * - `toolError`: the tool throws. The run's second request carries the result. + */ +export const Route = createFileRoute('/api/anthropic-redacted-thinking-wire')({ + server: { + handlers: { + POST: async () => { + try { + const nextTurn = scriptedFetch([ + [ + messageStart, + ...redacted(0, 'opaque-1'), + ...text(1, 'Hello.'), + ...end('end_turn'), + ], + ]) + const adapter = createAnthropicChat('claude-sonnet-4-5', DUMMY_KEY, { + fetch: nextTurn.fetchImpl, + }) + const client = new StreamProcessor({}) + for (const [index, prompt] of ['Hi', 'Say it again'].entries()) { + client.addUserMessage(prompt) + const params = await chatParamsFromRequestBody({ + threadId: 'redacted-thinking', + runId: `redacted-thinking-${index + 1}`, + // `JSON.parse(JSON.stringify(...))` stands in for the HTTP hop. + messages: JSON.parse( + JSON.stringify(uiMessagesToWire(client.getMessages())), + ), + tools: [], + context: [], + }) + for await (const chunk of chat({ + ...createChatOptions({ adapter }), + messages: params.messages, + stream: true, + })) { + client.processChunk(chunk) + } + client.finalizeStream() + } + + const sameRun = scriptedFetch([ + [ + messageStart, + ...redacted(0, 'opaque-2'), + ...toolUse(1, 'check'), + ...end('tool_use'), + ], + ]) + for await (const _ of chat({ + ...createChatOptions({ + adapter: createAnthropicChat('claude-sonnet-4-5', DUMMY_KEY, { + fetch: sameRun.fetchImpl, + }), + }), + messages: [{ role: 'user', content: 'Check it' }], + tools: [ + { + name: 'check', + description: 'Check something.', + execute: () => 'ok', + }, + ], + stream: true, + })) { + // consume + } + + const toolError = scriptedFetch([ + [messageStart, ...toolUse(0, 'fail'), ...end('tool_use')], + ]) + for await (const _ of chat({ + ...createChatOptions({ + adapter: createAnthropicChat('claude-sonnet-4-5', DUMMY_KEY, { + fetch: toolError.fetchImpl, + }), + }), + messages: [{ role: 'user', content: 'Try it' }], + tools: [ + { + name: 'fail', + description: 'Always fails.', + execute: () => { + throw new Error('Tool broke') + }, + }, + ], + stream: true, + })) { + // consume + } + + return Response.json({ + ok: true, + nextTurnBlocks: blocks(nextTurn.bodies[1], 'assistant'), + sameRunBlocks: blocks(sameRun.bodies[1], 'assistant').map( + (block) => block.type, + ), + toolResult: blocks(toolError.bodies[1], 'user').find( + (block) => block.type === 'tool_result', + ), + }) + } catch (error) { + return Response.json({ + ok: false, + error: error instanceof Error ? error.message : String(error), + }) + } + }, + }, + }, +}) diff --git a/testing/e2e/tests/anthropic-redacted-thinking-wire.spec.ts b/testing/e2e/tests/anthropic-redacted-thinking-wire.spec.ts new file mode 100644 index 0000000000..d203d57115 --- /dev/null +++ b/testing/e2e/tests/anthropic-redacted-thinking-wire.spec.ts @@ -0,0 +1,40 @@ +import { test, expect } from './fixtures' + +/** + * Claude encrypts some thinking as a `redacted_thinking` block with opaque + * `data`. The next request must send that block back unchanged, and a failed + * tool result must say `is_error: true`. + * `/api/anthropic-redacted-thinking-wire` scripts the responses and returns + * what the follow-up requests sent. + */ +test.describe('anthropic — redacted thinking and tool errors on the wire', () => { + test('the next turn, the same run, and a failed tool send them back', async ({ + request, + }) => { + const response = await request.post('/api/anthropic-redacted-thinking-wire') + expect(response.ok()).toBe(true) + const result = (await response.json()) as { + ok: boolean + error?: string + nextTurnBlocks: Array> + sameRunBlocks: Array + toolResult?: Record + } + if (!result.ok) throw new Error(`Route failed: ${result.error}`) + + // A new turn from the client history: the block crossed the wire. + expect(result.nextTurnBlocks).toEqual([ + { type: 'redacted_thinking', data: 'opaque-1' }, + { type: 'text', text: 'Hello.' }, + ]) + + // The same run continues after the tool, with the block first. + expect(result.sameRunBlocks).toEqual(['redacted_thinking', 'tool_use']) + + // The failed tool is marked as an error. + expect(result.toolResult).toMatchObject({ + tool_use_id: 'toolu_1', + is_error: true, + }) + }) +})