From bbe77969c024cbff17096abd1f4cc529660728d2 Mon Sep 17 00:00:00 2001 From: Alem Tuzlak Date: Wed, 30 Sep 2026 18:49:47 +0200 Subject: [PATCH 1/3] fix(ai, ai-anthropic): send redacted thinking and tool errors back to Claude The Anthropic adapter dropped a redacted_thinking block when it streamed in, so the next request did not send it back, and Claude can refuse a turn whose thinking blocks changed. A failed tool result also went out without is_error. A redacted block now becomes a thinking part with redacted: true, an empty content, and its data in signature. The flag goes through the stream, the engine, the stream processor, the UI and wire converters, and interrupt snapshots, and the adapter sends it back as { type: 'redacted_thinking', data }. A tool message with error sends tool_result.is_error: true. --- .changeset/anthropic-thinking-replay.md | 10 + docs/chat/thinking-content.md | 3 + docs/config.json | 2 +- packages/ai-anthropic/src/adapters/text.ts | 31 +++ .../tests/thinking-replay.test.ts | 226 ++++++++++++++++ packages/ai/src/activities/chat/index.ts | 11 +- packages/ai/src/activities/chat/messages.ts | 13 +- .../chat/stream/message-updaters.ts | 3 + .../src/activities/chat/stream/processor.ts | 8 + .../ai/src/activities/chat/stream/types.ts | 2 + packages/ai/src/types.ts | 15 +- .../ai/src/utilities/adapter-yield-chunk.ts | 2 + packages/ai/src/utilities/ag-ui-wire.ts | 3 + .../src/utilities/normalize-stream-chunk.ts | 1 + .../utilities/reasoning-encrypted-value.ts | 10 +- packages/ai/tests/redacted-thinking.test.ts | 68 +++++ testing/e2e/src/routeTree.gen.ts | 22 ++ .../api.anthropic-redacted-thinking-wire.ts | 241 ++++++++++++++++++ .../anthropic-redacted-thinking-wire.spec.ts | 40 +++ 19 files changed, 702 insertions(+), 9 deletions(-) create mode 100644 .changeset/anthropic-thinking-replay.md create mode 100644 packages/ai-anthropic/tests/thinking-replay.test.ts create mode 100644 packages/ai/tests/redacted-thinking.test.ts create mode 100644 testing/e2e/src/routes/api.anthropic-redacted-thinking-wire.ts create mode 100644 testing/e2e/tests/anthropic-redacted-thinking-wire.spec.ts diff --git a/.changeset/anthropic-thinking-replay.md b/.changeset/anthropic-thinking-replay.md new file mode 100644 index 0000000000..495784e096 --- /dev/null +++ b/.changeset/anthropic-thinking-replay.md @@ -0,0 +1,10 @@ +--- +'@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 }`. +- `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 aea1c76faa..06fdddb605 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..7beef02269 100644 --- a/packages/ai-anthropic/src/adapters/text.ts +++ b/packages/ai-anthropic/src/adapters/text.ts @@ -824,6 +824,7 @@ export class AnthropicTextAdapter< : typeof toolContent === 'string' ? toolContent : '', + ...(message.error !== undefined && { is_error: true }), }, ], }) @@ -952,6 +953,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 +1222,29 @@ 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 travels as the thinking step's signature. + const redactedStepId = genId() + yield { + type: EventType.STEP_STARTED, + stepName: redactedStepId, + stepId: redactedStepId, + model, + timestamp: Date.now(), + stepType: 'thinking', + } + yield { + type: EventType.STEP_FINISHED, + stepName: redactedStepId, + stepId: redactedStepId, + model, + timestamp: Date.now(), + delta: '', + content: '', + signature: event.content_block.data, + redacted: true, + } } } else if (event.type === 'content_block_delta') { if (event.delta.type === 'text_delta') { 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..5520f567ec --- /dev/null +++ b/packages/ai-anthropic/tests/thinking-replay.test.ts @@ -0,0 +1,226 @@ +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, 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('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 a5bd865082..d585a06112 100644 --- a/packages/ai/src/activities/chat/index.ts +++ b/packages/ai/src/activities/chat/index.ts @@ -853,8 +853,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 @@ -866,6 +865,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 @@ -1524,6 +1524,7 @@ class TextEngine< this.turnParts = [] this.currentThinkingContent = '' this.currentThinkingSignature = '' + this.currentThinkingRedacted = false this.finishedEvent = null this.streamedToolErrorResults.clear() @@ -1992,6 +1993,7 @@ class TextEngine< ...(this.currentThinkingSignature && { signature: this.currentThinkingSignature, }), + ...(this.currentThinkingRedacted && { redacted: true }), }) if (this.turnParts) { const placeholder = [...this.turnParts] @@ -2009,6 +2011,7 @@ class TextEngine< } this.currentThinkingContent = '' this.currentThinkingSignature = '' + this.currentThinkingRedacted = false } } @@ -2036,6 +2039,7 @@ class TextEngine< if (typeof chunk.signature === 'string' && chunk.signature !== '') { this.noteThinkingStepPosition() this.currentThinkingSignature = chunk.signature + this.currentThinkingRedacted = chunk.redacted === true } } @@ -2065,6 +2069,7 @@ class TextEngine< } this.noteThinkingStepPosition() this.currentThinkingSignature = chunk.encryptedValue + this.currentThinkingRedacted = tanstackMetadata(chunk)?.redacted === true } /** @@ -2579,7 +2584,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..5022867615 100644 --- a/packages/ai/src/activities/chat/messages.ts +++ b/packages/ai/src/activities/chat/messages.ts @@ -72,6 +72,11 @@ function encryptedValueFrom(value: object): string | undefined { return nonEmptyString(tanstackMetadata(value)?.signature) } +/** `{ redacted: true }` when a reasoning message carries a redacted block. */ +function redactedFrom(value: object) { + return tanstackMetadata(value)?.redacted === true ? { redacted: true } : {} +} + function toolCallFromWire(toolCall: ToolCall, bag: unknown): ToolCall { const fromBag = bag != null && typeof bag === 'object' && !Array.isArray(bag) @@ -246,7 +251,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 +278,7 @@ function convertOwnMessages( pendingThinking.push({ content: typeof content === 'string' ? content : '', ...(signature !== undefined ? { signature } : {}), + ...redactedFrom(msg), }) } continue @@ -663,7 +669,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 +770,7 @@ function buildAssistantMessages(uiMessage: UIMessage): Array { pendingThinking.push({ content: part.content, ...(part.signature && { signature: part.signature }), + ...(part.redacted && { redacted: true }), }) } break @@ -898,6 +905,7 @@ export function modelMessageToUIMessage( type: 'thinking', content: thinking.content, ...(thinking.signature && { signature: thinking.signature }), + ...(thinking.redacted && { redacted: true }), }) } } @@ -1121,6 +1129,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..1e20eb4066 100644 --- a/packages/ai/src/activities/chat/stream/message-updaters.ts +++ b/packages/ai/src/activities/chat/stream/message-updaters.ts @@ -451,6 +451,7 @@ export function updateThinkingPart( stepId: string, content: string, signature?: string, + redacted?: boolean, ): Array { return messages.map((msg) => { if (msg.id !== messageId) { @@ -483,12 +484,14 @@ 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 + const nextRedacted = redacted === true || adopted?.redacted === true const thinkingPart: ThinkingPart = { type: 'thinking', content, stepId, ...(nextSignature && { signature: nextSignature }), + ...(nextRedacted && { redacted: true }), } if (thinkingPartIndex >= 0) { diff --git a/packages/ai/src/activities/chat/stream/processor.ts b/packages/ai/src/activities/chat/stream/processor.ts index 04e72128ea..b7fbee32e0 100644 --- a/packages/ai/src/activities/chat/stream/processor.ts +++ b/packages/ai/src/activities/chat/stream/processor.ts @@ -786,6 +786,7 @@ export class StreamProcessor { hasSeenReasoningEvents: false, thinkingSteps: new Map(), thinkingStepSignatures: new Map(), + thinkingStepRedacted: new Set(), thinkingStepOrder: [], currentThinkingStepId: null, toolCalls: new Map(), @@ -2356,12 +2357,14 @@ export class StreamProcessor { if (thinking === undefined) return state.thinkingStepSignatures.set(stepId, signature) + if (extra.redacted === true) state.thinkingStepRedacted.add(stepId) this.messages = updateThinkingPart( this.messages, messageId, stepId, thinking, signature, + state.thinkingStepRedacted.has(stepId), ) this.emitMessagesChange() } @@ -2400,6 +2403,7 @@ export class StreamProcessor { stepId, nextThinking, state.thinkingStepSignatures.get(stepId), + state.thinkingStepRedacted.has(stepId), ) this.emitMessagesChange() @@ -2427,6 +2431,9 @@ export class StreamProcessor { ) const stepId = state.currentThinkingStepId ?? chunk.entityId state.thinkingStepSignatures.set(stepId, encryptedValue) + if (tanstackMetadata(chunk)?.redacted === true) { + state.thinkingStepRedacted.add(stepId) + } const content = state.thinkingSteps.get(stepId) ?? '' if (!state.thinkingSteps.has(stepId)) { state.thinkingSteps.set(stepId, content) @@ -2438,6 +2445,7 @@ export class StreamProcessor { stepId, content, encryptedValue, + state.thinkingStepRedacted.has(stepId), ) this.emitMessagesChange() } diff --git a/packages/ai/src/activities/chat/stream/types.ts b/packages/ai/src/activities/chat/stream/types.ts index cbdf333c3b..e4fb7f979e 100644 --- a/packages/ai/src/activities/chat/stream/types.ts +++ b/packages/ai/src/activities/chat/stream/types.ts @@ -64,6 +64,8 @@ export interface MessageStreamState { hasSeenReasoningEvents: boolean thinkingSteps: Map thinkingStepSignatures: Map + /** Thinking steps whose signature is a redacted block's data. */ + thinkingStepRedacted: Set thinkingStepOrder: Array currentThinkingStepId: string | null toolCalls: Map diff --git a/packages/ai/src/types.ts b/packages/ai/src/types.ts index 36bfaecb4c..466152428c 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. + */ + 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`. */ @@ -466,6 +471,12 @@ export interface ThinkingPart { content: string stepId?: string 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. + */ + redacted?: boolean } /** @@ -584,6 +595,8 @@ export interface TanStackMessageMetadata { subagent?: SubagentWireInfo /** Thinking signature for a `role: 'reasoning'` fan-out message. */ signature?: string + /** Set with `signature` when the provider redacted the thinking block. */ + redacted?: boolean /** Per-tool-call provider metadata keyed by tool call id (e.g. Gemini thoughtSignature). */ toolCallMetadata?: Record toolResult?: { diff --git a/packages/ai/src/utilities/adapter-yield-chunk.ts b/packages/ai/src/utilities/adapter-yield-chunk.ts index d733c86984..7b4c752730 100644 --- a/packages/ai/src/utilities/adapter-yield-chunk.ts +++ b/packages/ai/src/utilities/adapter-yield-chunk.ts @@ -24,6 +24,8 @@ type AdapterExtras = { stepType?: string delta?: string | ReadonlyArray signature?: string + /** With `signature`: the provider redacted this thinking block. */ + redacted?: boolean error?: { message: string; code?: string } 'tanstack:interruptErrors'?: ReadonlyArray threadId?: string diff --git a/packages/ai/src/utilities/ag-ui-wire.ts b/packages/ai/src/utilities/ag-ui-wire.ts index da2aac2bcc..24f8a921f5 100644 --- a/packages/ai/src/utilities/ag-ui-wire.ts +++ b/packages/ai/src/utilities/ag-ui-wire.ts @@ -225,6 +225,9 @@ export function uiMessagesToWire( if (part.signature) { reasoning.encryptedValue = part.signature } + if (part.redacted) { + reasoning.metadata = { tanstack: { redacted: true } } + } wire.push(reasoning) } diff --git a/packages/ai/src/utilities/normalize-stream-chunk.ts b/packages/ai/src/utilities/normalize-stream-chunk.ts index a706252269..34b8b58542 100644 --- a/packages/ai/src/utilities/normalize-stream-chunk.ts +++ b/packages/ai/src/utilities/normalize-stream-chunk.ts @@ -41,6 +41,7 @@ function encryptedValueExtras(chunk: AdapterYieldChunk): Array { entityId, encryptedValue: chunk.signature, timestamp, + redacted: chunk.redacted === true, }), ) } diff --git a/packages/ai/src/utilities/reasoning-encrypted-value.ts b/packages/ai/src/utilities/reasoning-encrypted-value.ts index 1236a30234..919e50dad8 100644 --- a/packages/ai/src/utilities/reasoning-encrypted-value.ts +++ b/packages/ai/src/utilities/reasoning-encrypted-value.ts @@ -1,18 +1,24 @@ import { EventType } from '../types' +import { withTanstackMetadata } from './merge-metadata' import type { ReasoningEncryptedValueEvent } from '../types' -/** Spec event that carries a provider thinking / tool-call signature blob. */ +/** + * Spec event that carries a provider thinking / tool-call signature blob. + * `redacted: true` marks a redacted thinking block, in `metadata.tanstack`. + */ export function reasoningEncryptedValue(opts: { subtype: 'message' | 'tool-call' entityId: string encryptedValue: string timestamp?: number + redacted?: boolean }): ReasoningEncryptedValueEvent { - return { + const event: ReasoningEncryptedValueEvent = { type: EventType.REASONING_ENCRYPTED_VALUE, subtype: opts.subtype, entityId: opts.entityId, encryptedValue: opts.encryptedValue, ...(opts.timestamp !== undefined ? { timestamp: opts.timestamp } : {}), } + return opts.redacted ? withTanstackMetadata(event, { redacted: true }) : event } diff --git a/packages/ai/tests/redacted-thinking.test.ts b/packages/ai/tests/redacted-thinking.test.ts new file mode 100644 index 0000000000..99a57ab967 --- /dev/null +++ b/packages/ai/tests/redacted-thinking.test.ts @@ -0,0 +1,68 @@ +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('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, + }) + }) +}) From 529774bd08a218644819a3015b15a134614a14e3 Mon Sep 17 00:00:00 2001 From: Tom Beckenham <34339192+tombeckenham@users.noreply.github.com> Date: Thu, 1 Oct 2026 11:15:35 +1000 Subject: [PATCH 2/3] refactor(ai): name the redacted thinking step set, note the signature rename Rename thinkingStepRedacted to redactedThinkingStepIds: it is a set of step ids, not a flag. Document that ThinkingPart.signature holds any provider's opaque reasoning bytes and should become encryptedValue (#1581). --- packages/ai/src/activities/chat/stream/processor.ts | 12 ++++++------ packages/ai/src/activities/chat/stream/types.ts | 2 +- packages/ai/src/types.ts | 8 +++++++- 3 files changed, 14 insertions(+), 8 deletions(-) diff --git a/packages/ai/src/activities/chat/stream/processor.ts b/packages/ai/src/activities/chat/stream/processor.ts index b7fbee32e0..034293286f 100644 --- a/packages/ai/src/activities/chat/stream/processor.ts +++ b/packages/ai/src/activities/chat/stream/processor.ts @@ -786,7 +786,7 @@ export class StreamProcessor { hasSeenReasoningEvents: false, thinkingSteps: new Map(), thinkingStepSignatures: new Map(), - thinkingStepRedacted: new Set(), + redactedThinkingStepIds: new Set(), thinkingStepOrder: [], currentThinkingStepId: null, toolCalls: new Map(), @@ -2357,14 +2357,14 @@ export class StreamProcessor { if (thinking === undefined) return state.thinkingStepSignatures.set(stepId, signature) - if (extra.redacted === true) state.thinkingStepRedacted.add(stepId) + if (extra.redacted === true) state.redactedThinkingStepIds.add(stepId) this.messages = updateThinkingPart( this.messages, messageId, stepId, thinking, signature, - state.thinkingStepRedacted.has(stepId), + state.redactedThinkingStepIds.has(stepId), ) this.emitMessagesChange() } @@ -2403,7 +2403,7 @@ export class StreamProcessor { stepId, nextThinking, state.thinkingStepSignatures.get(stepId), - state.thinkingStepRedacted.has(stepId), + state.redactedThinkingStepIds.has(stepId), ) this.emitMessagesChange() @@ -2432,7 +2432,7 @@ export class StreamProcessor { const stepId = state.currentThinkingStepId ?? chunk.entityId state.thinkingStepSignatures.set(stepId, encryptedValue) if (tanstackMetadata(chunk)?.redacted === true) { - state.thinkingStepRedacted.add(stepId) + state.redactedThinkingStepIds.add(stepId) } const content = state.thinkingSteps.get(stepId) ?? '' if (!state.thinkingSteps.has(stepId)) { @@ -2445,7 +2445,7 @@ export class StreamProcessor { stepId, content, encryptedValue, - state.thinkingStepRedacted.has(stepId), + state.redactedThinkingStepIds.has(stepId), ) this.emitMessagesChange() } diff --git a/packages/ai/src/activities/chat/stream/types.ts b/packages/ai/src/activities/chat/stream/types.ts index e4fb7f979e..cac4c5a0ab 100644 --- a/packages/ai/src/activities/chat/stream/types.ts +++ b/packages/ai/src/activities/chat/stream/types.ts @@ -65,7 +65,7 @@ export interface MessageStreamState { thinkingSteps: Map thinkingStepSignatures: Map /** Thinking steps whose signature is a redacted block's data. */ - thinkingStepRedacted: Set + redactedThinkingStepIds: Set thinkingStepOrder: Array currentThinkingStepId: string | null toolCalls: Map diff --git a/packages/ai/src/types.ts b/packages/ai/src/types.ts index 466152428c..b4dfd0542b 100644 --- a/packages/ai/src/types.ts +++ b/packages/ai/src/types.ts @@ -384,7 +384,7 @@ export interface ModelMessage< /** * 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. + * 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. */ @@ -470,6 +470,12 @@ 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`). From c883b870cdd6e2b631dc1f0e9995da8ecf290463 Mon Sep 17 00:00:00 2001 From: Alem Tuzlak Date: Thu, 1 Oct 2026 10:34:37 +0200 Subject: [PATCH 3/3] fix(ai, ai-anthropic): mark redacted thinking by its reasoning message id AG-UI clients do not copy event metadata onto messages, so metadata.tanstack.redacted did not survive a client that is not TanStack's. A redacted block is now its own reasoning message. Its id starts with redacted_thinking-, and REASONING_ENCRYPTED_VALUE names that message in entityId. The server and the snapshot loader read the flag from the id. Anthropic thinking signatures now also go out as REASONING_ENCRYPTED_VALUE for the reasoning message, not on STEP_FINISHED with the step id, so an AG-UI client can attach them. --- .changeset/anthropic-thinking-replay.md | 2 + packages/ai-anthropic/src/adapters/text.ts | 70 +++++++++++++++---- .../tests/thinking-replay.test.ts | 48 ++++++++++++- packages/ai/src/activities/chat/index.ts | 5 +- packages/ai/src/activities/chat/messages.ts | 7 +- .../chat/stream/message-updaters.ts | 8 ++- .../src/activities/chat/stream/processor.ts | 8 --- .../ai/src/activities/chat/stream/types.ts | 2 - packages/ai/src/adapter-internals.ts | 1 + packages/ai/src/types.ts | 5 +- .../ai/src/utilities/adapter-yield-chunk.ts | 2 - packages/ai/src/utilities/ag-ui-wire.ts | 9 +-- .../src/utilities/normalize-stream-chunk.ts | 1 - .../utilities/reasoning-encrypted-value.ts | 22 +++--- packages/ai/tests/redacted-thinking.test.ts | 29 ++++++++ 15 files changed, 171 insertions(+), 48 deletions(-) diff --git a/.changeset/anthropic-thinking-replay.md b/.changeset/anthropic-thinking-replay.md index 495784e096..e5c18def33 100644 --- a/.changeset/anthropic-thinking-replay.md +++ b/.changeset/anthropic-thinking-replay.md @@ -7,4 +7,6 @@ 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/packages/ai-anthropic/src/adapters/text.ts b/packages/ai-anthropic/src/adapters/text.ts index 7beef02269..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' @@ -1224,26 +1227,60 @@ export class AnthropicTextAdapter< } } else if (event.content_block.type === 'redacted_thinking') { // Encrypted thinking: no text, and its data must go back to - // Anthropic unchanged. It travels as the thinking step's signature. - const redactedStepId = genId() + // 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: redactedStepId, - stepId: redactedStepId, + stepName: redactedId, + stepId: redactedId, model, timestamp: Date.now(), stepType: 'thinking', } yield { type: EventType.STEP_FINISHED, - stepName: redactedStepId, - stepId: redactedStepId, + stepName: redactedId, + stepId: redactedId, model, timestamp: Date.now(), delta: '', content: '', - signature: event.content_block.data, - redacted: true, + } + 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') { @@ -1363,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, @@ -1373,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 index 5520f567ec..282e5b5a6e 100644 --- a/packages/ai-anthropic/tests/thinking-replay.test.ts +++ b/packages/ai-anthropic/tests/thinking-replay.test.ts @@ -2,7 +2,7 @@ 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, Tool } from '@tanstack/ai' +import type { ModelMessage, StreamChunk, Tool } from '@tanstack/ai' const mocks = vi.hoisted(() => { const betaMessagesCreate = vi.fn() @@ -188,6 +188,52 @@ describe('Anthropic replay', () => { }) }) + 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( diff --git a/packages/ai/src/activities/chat/index.ts b/packages/ai/src/activities/chat/index.ts index d585a06112..0a7ba03e29 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 { @@ -2039,7 +2040,7 @@ class TextEngine< if (typeof chunk.signature === 'string' && chunk.signature !== '') { this.noteThinkingStepPosition() this.currentThinkingSignature = chunk.signature - this.currentThinkingRedacted = chunk.redacted === true + this.currentThinkingRedacted = isRedactedThinkingId(chunk.stepId) } } @@ -2069,7 +2070,7 @@ class TextEngine< } this.noteThinkingStepPosition() this.currentThinkingSignature = chunk.encryptedValue - this.currentThinkingRedacted = tanstackMetadata(chunk)?.redacted === true + this.currentThinkingRedacted = isRedactedThinkingId(chunk.entityId) } /** diff --git a/packages/ai/src/activities/chat/messages.ts b/packages/ai/src/activities/chat/messages.ts index 5022867615..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,9 +73,11 @@ function encryptedValueFrom(value: object): string | undefined { return nonEmptyString(tanstackMetadata(value)?.signature) } -/** `{ redacted: true }` when a reasoning message carries a redacted block. */ +/** `{ redacted: true }` when a reasoning message's id marks a redacted block. */ function redactedFrom(value: object) { - return tanstackMetadata(value)?.redacted === true ? { redacted: true } : {} + return 'id' in value && isRedactedThinkingId(value.id) + ? { redacted: true } + : {} } function toolCallFromWire(toolCall: ToolCall, bag: unknown): ToolCall { diff --git a/packages/ai/src/activities/chat/stream/message-updaters.ts b/packages/ai/src/activities/chat/stream/message-updaters.ts index 1e20eb4066..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, @@ -451,7 +452,6 @@ export function updateThinkingPart( stepId: string, content: string, signature?: string, - redacted?: boolean, ): Array { return messages.map((msg) => { if (msg.id !== messageId) { @@ -484,14 +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 - const nextRedacted = redacted === true || adopted?.redacted === true + // 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 }), - ...(nextRedacted && { redacted: true }), + ...(redacted && { redacted: true }), } if (thinkingPartIndex >= 0) { diff --git a/packages/ai/src/activities/chat/stream/processor.ts b/packages/ai/src/activities/chat/stream/processor.ts index 034293286f..04e72128ea 100644 --- a/packages/ai/src/activities/chat/stream/processor.ts +++ b/packages/ai/src/activities/chat/stream/processor.ts @@ -786,7 +786,6 @@ export class StreamProcessor { hasSeenReasoningEvents: false, thinkingSteps: new Map(), thinkingStepSignatures: new Map(), - redactedThinkingStepIds: new Set(), thinkingStepOrder: [], currentThinkingStepId: null, toolCalls: new Map(), @@ -2357,14 +2356,12 @@ export class StreamProcessor { if (thinking === undefined) return state.thinkingStepSignatures.set(stepId, signature) - if (extra.redacted === true) state.redactedThinkingStepIds.add(stepId) this.messages = updateThinkingPart( this.messages, messageId, stepId, thinking, signature, - state.redactedThinkingStepIds.has(stepId), ) this.emitMessagesChange() } @@ -2403,7 +2400,6 @@ export class StreamProcessor { stepId, nextThinking, state.thinkingStepSignatures.get(stepId), - state.redactedThinkingStepIds.has(stepId), ) this.emitMessagesChange() @@ -2431,9 +2427,6 @@ export class StreamProcessor { ) const stepId = state.currentThinkingStepId ?? chunk.entityId state.thinkingStepSignatures.set(stepId, encryptedValue) - if (tanstackMetadata(chunk)?.redacted === true) { - state.redactedThinkingStepIds.add(stepId) - } const content = state.thinkingSteps.get(stepId) ?? '' if (!state.thinkingSteps.has(stepId)) { state.thinkingSteps.set(stepId, content) @@ -2445,7 +2438,6 @@ export class StreamProcessor { stepId, content, encryptedValue, - state.redactedThinkingStepIds.has(stepId), ) this.emitMessagesChange() } diff --git a/packages/ai/src/activities/chat/stream/types.ts b/packages/ai/src/activities/chat/stream/types.ts index cac4c5a0ab..cbdf333c3b 100644 --- a/packages/ai/src/activities/chat/stream/types.ts +++ b/packages/ai/src/activities/chat/stream/types.ts @@ -64,8 +64,6 @@ export interface MessageStreamState { hasSeenReasoningEvents: boolean thinkingSteps: Map thinkingStepSignatures: Map - /** Thinking steps whose signature is a redacted block's data. */ - redactedThinkingStepIds: Set thinkingStepOrder: Array currentThinkingStepId: string | null toolCalls: Map 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 b4dfd0542b..fc13aa5c1c 100644 --- a/packages/ai/src/types.ts +++ b/packages/ai/src/types.ts @@ -480,7 +480,8 @@ export interface ThinkingPart { /** * 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. + * to the provider unchanged. On the AG-UI wire, the reasoning message id + * starts with `redacted_thinking-` instead. */ redacted?: boolean } @@ -601,8 +602,6 @@ export interface TanStackMessageMetadata { subagent?: SubagentWireInfo /** Thinking signature for a `role: 'reasoning'` fan-out message. */ signature?: string - /** Set with `signature` when the provider redacted the thinking block. */ - redacted?: boolean /** Per-tool-call provider metadata keyed by tool call id (e.g. Gemini thoughtSignature). */ toolCallMetadata?: Record toolResult?: { diff --git a/packages/ai/src/utilities/adapter-yield-chunk.ts b/packages/ai/src/utilities/adapter-yield-chunk.ts index 7b4c752730..d733c86984 100644 --- a/packages/ai/src/utilities/adapter-yield-chunk.ts +++ b/packages/ai/src/utilities/adapter-yield-chunk.ts @@ -24,8 +24,6 @@ type AdapterExtras = { stepType?: string delta?: string | ReadonlyArray signature?: string - /** With `signature`: the provider redacted this thinking block. */ - redacted?: boolean error?: { message: string; code?: string } 'tanstack:interruptErrors'?: ReadonlyArray threadId?: string diff --git a/packages/ai/src/utilities/ag-ui-wire.ts b/packages/ai/src/utilities/ag-ui-wire.ts index 24f8a921f5..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' @@ -225,9 +226,6 @@ export function uiMessagesToWire( if (part.signature) { reasoning.encryptedValue = part.signature } - if (part.redacted) { - reasoning.metadata = { tanstack: { redacted: true } } - } wire.push(reasoning) } @@ -617,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/normalize-stream-chunk.ts b/packages/ai/src/utilities/normalize-stream-chunk.ts index 34b8b58542..a706252269 100644 --- a/packages/ai/src/utilities/normalize-stream-chunk.ts +++ b/packages/ai/src/utilities/normalize-stream-chunk.ts @@ -41,7 +41,6 @@ function encryptedValueExtras(chunk: AdapterYieldChunk): Array { entityId, encryptedValue: chunk.signature, timestamp, - redacted: chunk.redacted === true, }), ) } diff --git a/packages/ai/src/utilities/reasoning-encrypted-value.ts b/packages/ai/src/utilities/reasoning-encrypted-value.ts index 919e50dad8..4b2de205b8 100644 --- a/packages/ai/src/utilities/reasoning-encrypted-value.ts +++ b/packages/ai/src/utilities/reasoning-encrypted-value.ts @@ -1,24 +1,30 @@ import { EventType } from '../types' -import { withTanstackMetadata } from './merge-metadata' import type { ReasoningEncryptedValueEvent } from '../types' -/** - * Spec event that carries a provider thinking / tool-call signature blob. - * `redacted: true` marks a redacted thinking block, in `metadata.tanstack`. - */ +/** Spec event that carries a provider thinking / tool-call signature blob. */ export function reasoningEncryptedValue(opts: { subtype: 'message' | 'tool-call' entityId: string encryptedValue: string timestamp?: number - redacted?: boolean }): ReasoningEncryptedValueEvent { - const event: ReasoningEncryptedValueEvent = { + return { type: EventType.REASONING_ENCRYPTED_VALUE, subtype: opts.subtype, entityId: opts.entityId, encryptedValue: opts.encryptedValue, ...(opts.timestamp !== undefined ? { timestamp: opts.timestamp } : {}), } - return opts.redacted ? withTanstackMetadata(event, { redacted: true }) : event +} + +/** + * 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 index 99a57ab967..2794b9410d 100644 --- a/packages/ai/tests/redacted-thinking.test.ts +++ b/packages/ai/tests/redacted-thinking.test.ts @@ -53,6 +53,35 @@ describe('redacted thinking', () => { ) }) + 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))