diff --git a/src/browser/features/Tools/WebSearchToolCall.tsx b/src/browser/features/Tools/WebSearchToolCall.tsx index 2a470674920..67d20974206 100644 --- a/src/browser/features/Tools/WebSearchToolCall.tsx +++ b/src/browser/features/Tools/WebSearchToolCall.tsx @@ -17,6 +17,7 @@ import { type ToolStatus, } from "./Shared/toolUtils"; import { JsonHighlight } from "./Shared/HighlightedCode"; +import { stripEncryptedContent } from "@/common/utils/messages/stripEncryptedContent"; interface WebSearchToolCallProps { args: { query?: string }; // Anthropic puts query in args @@ -102,7 +103,9 @@ export const WebSearchToolCall: React.FC = ({ Results
- + {/* Anthropic results keep their ciphertext for native replay (#5887): it is + opaque and large, so the transcript shows the rest. */} +
)} diff --git a/src/browser/utils/messages/modelMessageTransform.test.ts b/src/browser/utils/messages/modelMessageTransform.test.ts index 9981468edd2..d6224fd40b3 100644 --- a/src/browser/utils/messages/modelMessageTransform.test.ts +++ b/src/browser/utils/messages/modelMessageTransform.test.ts @@ -427,6 +427,33 @@ describe("modelMessageTransform", () => { expect(result.valid).toBe(true); }); + it("accepts a provider-executed call whose result rides in the same message (#5887)", () => { + const messages: ModelMessage[] = [ + { role: "user", content: [{ type: "text", text: "search" }] }, + { + role: "assistant", + content: [ + { + type: "tool-call", + toolCallId: "srvtoolu_1", + toolName: "web_search", + input: { query: "xum" }, + providerExecuted: true, + }, + { + type: "tool-result", + toolCallId: "srvtoolu_1", + toolName: "web_search", + output: { type: "json", value: [] }, + }, + { type: "text", text: "found it" }, + ], + }, + ]; + + expect(validateAnthropicCompliance(messages)).toEqual({ valid: true }); + }); + it("should detect tool calls without results", () => { const assistantMsg1: AssistantModelMessage = { role: "assistant", diff --git a/src/browser/utils/messages/modelMessageTransform.ts b/src/browser/utils/messages/modelMessageTransform.ts index f81f9a80497..3e2265e60ae 100644 --- a/src/browser/utils/messages/modelMessageTransform.ts +++ b/src/browser/utils/messages/modelMessageTransform.ts @@ -333,7 +333,11 @@ function splitMixedContentMessages(messages: ModelMessage[]): ModelMessage[] { continue; } - const toolCallParts = assistantMsg.content.filter((c) => c.type === "tool-call"); + // A provider-executed call carries its result inline in this message (Anthropic + // server_tool_use + web_search_tool_result, #5887): it stays with the content around it. + const isClientToolCall = (part: (typeof assistantMsg.content)[number]) => + part.type === "tool-call" && part.providerExecuted !== true; + const toolCallParts = assistantMsg.content.filter(isClientToolCall); if (toolCallParts.length === 0) { result.push(msg); @@ -355,7 +359,7 @@ function splitMixedContentMessages(messages: ModelMessage[]): ModelMessage[] { let currentGroup: { type: "text" | "tool-call"; parts: ContentArray } | null = null; for (const part of assistantMsg.content) { - const partType = part.type === "tool-call" ? "tool-call" : "text"; + const partType = isClientToolCall(part) ? "tool-call" : "text"; // eslint-disable-next-line @typescript-eslint/prefer-optional-chain if (!currentGroup || currentGroup.type !== partType) { @@ -1139,7 +1143,12 @@ function ensureAnthropicThinkingBeforeToolCalls(messages: ModelMessage[]): Model // interleaved thinking, text can sit between two thinking blocks, and moving a block // edits the prefix every later block is bound to (#5887). Preserved thinking: "send it // back unchanged ... in the order received". - if (!hasToolCall || content[0]?.type === "reasoning") { + // The same holds for a message that opens with a natively replayed server tool (#5887): the + // API itself started the response with server_tool_use, and the thinking after the search is + // bound to it. + const opensWithServerTool = + content[0]?.type === "tool-call" && content[0].providerExecuted === true; + if (!hasToolCall || content[0]?.type === "reasoning" || opensWithServerTool) { result.push(msg); continue; } @@ -1321,7 +1330,9 @@ export function validateAnthropicCompliance(messages: ModelMessage[]): { // Track any tool calls in this message for (const content of assistantMsg.content) { - if (content.type === "tool-call") { + // A provider-executed call (Anthropic server tool, #5887) carries its result in the + // same assistant message and needs no tool message after it. + if (content.type === "tool-call" && content.providerExecuted !== true) { pendingToolCalls.set(content.toolCallId, i); } } diff --git a/src/common/orpc/schemas/message.ts b/src/common/orpc/schemas/message.ts index ef3bc2eb991..242d23b0112 100644 --- a/src/common/orpc/schemas/message.ts +++ b/src/common/orpc/schemas/message.ts @@ -50,6 +50,9 @@ const MuxToolPartBase = z.object({ workflowRun: WorkflowRunToolAttachmentSchema.optional(), // Host-authored display data must not enter the model-visible output. mcpServer: MCPToolCallDisplaySchema.optional().catch(undefined), + // The provider ran the tool server-side (Anthropic web_search, ...). Stored only for the + // Anthropic Messages wire, so history can replay it natively (#5887). + providerExecuted: z.boolean().optional(), }); /** diff --git a/src/common/utils/messages/anthropicNativeServerTools.test.ts b/src/common/utils/messages/anthropicNativeServerTools.test.ts new file mode 100644 index 00000000000..2974c744372 --- /dev/null +++ b/src/common/utils/messages/anthropicNativeServerTools.test.ts @@ -0,0 +1,152 @@ +import { describe, expect, test } from "bun:test"; +import type { DynamicToolPart } from "@/common/types/toolParts"; +import { ANTHROPIC_NATIVE_SERVER_TOOL_MAX_ROW_CIPHERTEXT_CHARS } from "@/constants/anthropicServerTools"; +import { createMuxMessage, type MuxMessage } from "@/common/types/message"; +import { + isNativeAnthropicReplayable, + projectAnthropicServerTools, + rowCiphertextChars, + toStoredServerToolPart, +} from "./anthropicNativeServerTools"; + +/** A completed Anthropic web_search part whose results carry these ciphertext lengths. */ +function webSearchPart(id: string, ciphertextLengths: number[]): DynamicToolPart { + return { + type: "dynamic-tool", + toolCallId: id, + toolName: "web_search", + state: "output-available", + input: { query: "xum" }, + providerExecuted: true, + output: ciphertextLengths.map((length, index) => ({ + type: "web_search_result", + url: `https://example.com/${index}`, + title: "Xum", + pageAge: null, + encryptedContent: "e".repeat(length), + })), + }; +} + +/** Ciphertext lengths of at most 12,000 chars (the generic sanitizer bound) that add up to `total`. */ +function chunks(total: number): number[] { + const out: number[] = []; + for (let left = total; left > 0; left -= 12_000) out.push(Math.min(left, 12_000)); + return out; +} + +function hasCiphertext(part: DynamicToolPart): boolean { + return part.state === "output-available" && JSON.stringify(part.output).includes("eee"); +} + +/** Stores the parts in order, the way StreamManager completes them within one row. */ +function storeRow(parts: DynamicToolPart[]): DynamicToolPart[] { + const stored: DynamicToolPart[] = []; + for (const part of parts) { + stored.push( + toStoredServerToolPart(part, { + resultFollowsCall: true, + rowCiphertextChars: rowCiphertextChars(stored), + }) + ); + } + return stored; +} + +describe("toStoredServerToolPart ciphertext bound (#5887)", () => { + const limit = ANTHROPIC_NATIVE_SERVER_TOOL_MAX_ROW_CIPHERTEXT_CHARS; + + test("keeps native replay when the row's summed ciphertext is at the limit", () => { + // Three results over two calls: the bound sums every result of every search in the row. + const row = storeRow([ + webSearchPart("srvtoolu_1", [...chunks(limit - 30), 10]), + webSearchPart("srvtoolu_2", [20]), + ]); + expect(row.map((part) => part.providerExecuted)).toEqual([true, true]); + expect(row.every(hasCiphertext)).toBe(true); + expect(rowCiphertextChars(row)).toBe(limit); + }); + + test("stores the search that crosses the limit as the client pair without ciphertext", () => { + const row = storeRow([ + webSearchPart("srvtoolu_1", [...chunks(limit - 30), 10]), + webSearchPart("srvtoolu_2", [21]), + // A later small search still fits: only the ciphertext kept counts. + webSearchPart("srvtoolu_3", [20]), + ]); + expect(row.map((part) => part.providerExecuted)).toEqual([true, undefined, true]); + expect(hasCiphertext(row[1])).toBe(false); + // The search itself stays in history (URLs and titles), only the ciphertext is dropped. + expect(row[1].state === "output-available" && Array.isArray(row[1].output)).toBe(true); + expect(rowCiphertextChars(row)).toBe(limit); + }); +}); + +describe("isNativeAnthropicReplayable fields (#5887)", () => { + // Native replay sends every field back as the API returned it, so a field the generic + // provider-output sanitizer would rewrite (too long, control text) cannot replay natively. + function withTitle(title: string): DynamicToolPart { + const part = webSearchPart("srvtoolu_1", [10]); + if (part.state !== "output-available" || !Array.isArray(part.output)) throw new Error("setup"); + return { ...part, output: [{ ...(part.output[0] as object), title }] }; + } + + test("a plain title replays natively", () => { + expect(isNativeAnthropicReplayable(withTitle("Xum"))).toBe(true); + }); + + function withResult(fields: Record): DynamicToolPart { + const part = webSearchPart("srvtoolu_1", [10]); + if (part.state !== "output-available" || !Array.isArray(part.output)) throw new Error("setup"); + return { ...part, output: [{ ...(part.output[0] as object), ...fields }] }; + } + + test("ciphertext the sanitizer would rewrite does not replay natively", () => { + // Older builds still run the sanitizer on stored native rows: only unchanged bytes are safe. + expect(isNativeAnthropicReplayable(withResult({ encryptedContent: "e".repeat(12_000) }))).toBe( + true + ); + expect(isNativeAnthropicReplayable(withResult({ encryptedContent: "e".repeat(12_001) }))).toBe( + false + ); + }); + + test("a result without pageAge does not replay natively", () => { + // The SDK's replay schema requires pageAge as a string or null: undefined throws. + const part = withResult({}); + if (part.state !== "output-available" || !Array.isArray(part.output)) throw new Error("setup"); + const { pageAge: _dropped, ...withoutPageAge } = part.output[0] as Record; + expect(isNativeAnthropicReplayable({ ...part, output: [withoutPageAge] })).toBe(false); + }); + + test("a title the sanitizer would rewrite does not", () => { + expect(isNativeAnthropicReplayable(withTitle("t".repeat(12_001)))).toBe(false); + expect(isNativeAnthropicReplayable(withTitle("bad\u0000title"))).toBe(false); + }); +}); + +describe("projectAnthropicServerTools demotion across rows (#5887)", () => { + // Preserved thinking binds each thinking block to everything before it, earlier rows too. + function rows(): MuxMessage[] { + return [ + createMuxMessage("user-1", "user", "search", { historySequence: 0 }), + createMuxMessage("assistant-1", "assistant", "", { historySequence: 1 }, [ + webSearchPart("srvtoolu_1", [10]), + { type: "text", text: "found" }, + ]), + createMuxMessage("user-2", "user", "more", { historySequence: 2 }), + createMuxMessage("assistant-2", "assistant", "", { historySequence: 3 }, [ + { type: "reasoning", text: "read", providerOptions: { anthropic: { signature: "sig" } } }, + { type: "text", text: "done" }, + ]), + ]; + } + + test("thinking in a later row after a demoted search strips thinking", () => { + expect(projectAnthropicServerTools(rows(), false).demotedBeforeThinking).toBe(true); + }); + + test("the same rows replayed natively keep thinking", () => { + expect(projectAnthropicServerTools(rows(), true).demotedBeforeThinking).toBe(false); + }); +}); diff --git a/src/common/utils/messages/anthropicNativeServerTools.ts b/src/common/utils/messages/anthropicNativeServerTools.ts new file mode 100644 index 00000000000..0e203d9c3d9 --- /dev/null +++ b/src/common/utils/messages/anthropicNativeServerTools.ts @@ -0,0 +1,183 @@ +import type { MuxMessage } from "@/common/types/message"; +import type { DynamicToolPart } from "@/common/types/toolParts"; +import { ANTHROPIC_NATIVE_SERVER_TOOL_MAX_ROW_CIPHERTEXT_CHARS } from "@/constants/anthropicServerTools"; +import { sanitizeStringForProviderOutput } from "@/common/utils/providerOutputSanitization"; +import { stripEncryptedContent } from "./stripEncryptedContent"; + +/** + * #5887 option A: Anthropic server tools that history replays natively + * (`server_tool_use` + `web_search_tool_result`, in place, next to the signed thinking around + * them) instead of as a client `tool_use`/`tool_result` pair. Preserved thinking binds each + * thinking block to everything before it, so only the native blocks keep later thinking (and + * the prompt cache) valid. + * + * Scope: successful `web_search` results. Error results and other server tools (web_fetch, + * code_execution, ...) replay as a client pair, and the stream-time receipt + * (`MuxMetadata.anthropicThinkingReplay`) keeps the thinking after them out. + * + * Stream time (the receipt trigger) and request time (the projection below) both use this + * predicate, so they cannot disagree about a stored part. + */ +export function isNativeAnthropicReplayable(part: DynamicToolPart): boolean { + if (part.providerExecuted !== true || part.toolName !== "web_search") return false; + if (part.state !== "output-available") return false; + // The SDK validates every native result against a schema that requires encryptedContent and + // a string-or-null title: one missing field throws, and no request is sent at all. + return Array.isArray(part.output) && part.output.every(isReplayableWebSearchResult); +} + +function isReplayableWebSearchResult(item: unknown): boolean { + if (typeof item !== "object" || item === null) return false; + const result = item as Record; + return ( + result.type === "web_search_result" && + isVerbatimSafe(result.url) && + isVerbatimSafe(result.encryptedContent) && + (result.title === null || isVerbatimSafe(result.title)) && + // Present as a string or null: the SDK's replay schema rejects a missing pageAge. + (result.pageAge === null || isVerbatimSafe(result.pageAge)) + ); +} + +/** + * Native replay must send the result back exactly as the API returned it, but every request + * still runs the generic provider-output sanitizer (applyToolOutputRedaction), in this build and + * in older ones that read the same stored rows. So only fields that sanitizer leaves unchanged + * (at most 12,000 chars, no control text) can replay natively, the ciphertext included. Any other + * result replays as the client pair, and the receipt keeps the thinking after it out. + */ +function isVerbatimSafe(value: unknown): boolean { + return typeof value === "string" && sanitizeStringForProviderOutput(value) === value; +} + +/** + * The form a completed provider-executed part is stored in. It keeps the flag (and its + * ciphertext) only when history can replay it natively. Everything else is stored exactly as + * before #5887: a client pair without ciphertext. Older builds read these rows too: their SDK + * would throw on a flagged result it cannot validate, so the flag never outlives the check. + * + * Two more cases are stored demoted, and the stream-time receipt then keeps the thinking after + * them out: + * - `resultFollowsCall` false: other parts arrived between the call and its result. The API + * does this when Claude calls a client tool in the same parallel group: the response ends + * after both calls, and the server tool's result opens the next step. One stored part cannot + * replay the call and the result at their two positions, and moving them changes the prefix + * the later thinking is bound to. + * - Ciphertext that would take the row's total above + * ANTHROPIC_NATIVE_SERVER_TOOL_MAX_ROW_CIPHERTEXT_CHARS (row size bound). + * `rowCiphertextChars` is what the row's other parts already keep (see rowCiphertextChars). + */ +export function toStoredServerToolPart( + part: DynamicToolPart, + options: { resultFollowsCall: boolean; rowCiphertextChars: number } +): DynamicToolPart { + if (part.providerExecuted !== true || part.state !== "output-available") return part; + const titled = withNullTitles(part); + const native = + options.resultFollowsCall && + isNativeAnthropicReplayable(titled) && + options.rowCiphertextChars + ciphertextChars(part.output) <= + ANTHROPIC_NATIVE_SERVER_TOOL_MAX_ROW_CIPHERTEXT_CHARS; + return native ? titled : demote(part); +} + +/** Ciphertext the stored provider-executed parts of one row keep (only native parts keep any). */ +export function rowCiphertextChars(parts: ReadonlyArray): number { + let total = 0; + for (const part of parts) { + if ( + part.type === "dynamic-tool" && + part.providerExecuted === true && + part.state === "output-available" + ) { + total += ciphertextChars(part.output); + } + } + return total; +} + +/** Total encryptedContent length of a replayable web_search output (an array of results). */ +function ciphertextChars(output: unknown): number { + if (!Array.isArray(output)) return 0; + let total = 0; + for (const item of output as Array<{ encryptedContent?: unknown }>) { + if (typeof item.encryptedContent === "string") total += item.encryptedContent.length; + } + return total; +} + +export interface AnthropicServerToolProjection { + messages: MuxMessage[]; + /** + * An assistant row has a demoted server tool with reasoning after it. That thinking is + * bound to native blocks this request does not send, so the caller must send no Anthropic + * thinking (the same strip as the replay receipt). It is a pure function of the rows, so + * every request in the context segment decides the same way. + */ + demotedBeforeThinking: boolean; +} + +/** + * Request-only projection (history on disk keeps the native identity and ciphertext). + * `native` true: replayable parts stay provider-executed; every other provider-executed part + * is demoted. `native` false (another provider's wire, or a request that must not send + * Anthropic-native blocks): every provider-executed part is demoted. A demoted part replays as + * the client pair Xum sent before #5887, without the ciphertext. + */ +export function projectAnthropicServerTools( + messages: MuxMessage[], + native: boolean +): AnthropicServerToolProjection { + let demotedBeforeThinking = false; + // Preserved thinking binds a thinking block to everything before it, earlier rows included, + // so one demotion strips thinking that follows it anywhere later in the request. + let demoted = false; + const projected = messages.map((message) => { + if (message.role !== "assistant") return message; + // Most rows hold no server tool: return them as is, with no new parts array (perf). + if (!message.parts.some(isProviderExecutedTool)) { + if (demoted && message.parts.some((part) => part.type === "reasoning")) { + demotedBeforeThinking = true; + } + return message; + } + let changed = false; + const parts = message.parts.map((part) => { + if (part.type === "reasoning" && demoted) demotedBeforeThinking = true; + if (part.type !== "dynamic-tool" || part.providerExecuted !== true) return part; + if (native && isNativeAnthropicReplayable(part)) return part; + changed = true; + demoted = true; + return demote(part); + }); + return changed ? { ...message, parts } : message; + }); + return { messages: projected, demotedBeforeThinking }; +} + +function isProviderExecutedTool(part: MuxMessage["parts"][number]): boolean { + return part.type === "dynamic-tool" && part.providerExecuted === true; +} + +function demote(part: DynamicToolPart): DynamicToolPart { + const { providerExecuted: _demoted, ...clientPart } = part; + return clientPart.state === "output-available" + ? { ...clientPart, output: stripEncryptedContent(clientPart.output) } + : clientPart; +} + +/** + * The SDK stream omits a null `title`, but its replay schema requires `title` as a string or + * null. The API returned null, so null is the faithful value. + */ +function withNullTitles(part: DynamicToolPart): DynamicToolPart { + if (part.state !== "output-available" || !Array.isArray(part.output)) return part; + const output: unknown[] = part.output; + const untitled = (item: unknown) => + typeof item === "object" && item !== null && !("title" in item); + if (!output.some(untitled)) return part; + return { + ...part, + output: output.map((item) => (untitled(item) ? { ...(item as object), title: null } : item)), + }; +} diff --git a/src/node/utils/messages/stripEncryptedContent.test.ts b/src/common/utils/messages/stripEncryptedContent.test.ts similarity index 100% rename from src/node/utils/messages/stripEncryptedContent.test.ts rename to src/common/utils/messages/stripEncryptedContent.test.ts diff --git a/src/node/utils/messages/stripEncryptedContent.ts b/src/common/utils/messages/stripEncryptedContent.ts similarity index 100% rename from src/node/utils/messages/stripEncryptedContent.ts rename to src/common/utils/messages/stripEncryptedContent.ts diff --git a/src/constants/anthropicServerTools.ts b/src/constants/anthropicServerTools.ts new file mode 100644 index 00000000000..878af700d70 --- /dev/null +++ b/src/constants/anthropicServerTools.ts @@ -0,0 +1,10 @@ +/** + * Largest total `encryptedContent` (UTF-16 chars) that one assistant row keeps for native replay + * of Anthropic web searches (#5887), summed over every search in the row. A search that would + * exceed it is stored as the pre-#5887 client pair without ciphertext, and the thinking-replay + * receipt keeps the thinking after it out. Why a bound: the partial holding the row is rewritten + * on every throttled partial write for the rest of the turn, the row crosses IPC, and chat.jsonl + * readers skip rows over 1 MiB. One turn can run many searches, so the bound is per row, not per + * call. + */ +export const ANTHROPIC_NATIVE_SERVER_TOOL_MAX_ROW_CIPHERTEXT_CHARS = 256 * 1024; diff --git a/src/node/services/agentSession.continuousCompaction.test.ts b/src/node/services/agentSession.continuousCompaction.test.ts index c923b004e47..9547a631343 100644 --- a/src/node/services/agentSession.continuousCompaction.test.ts +++ b/src/node/services/agentSession.continuousCompaction.test.ts @@ -1800,6 +1800,83 @@ describe("AgentSession continuous compaction wiring", () => { expect(JSON.stringify(requests[0].prompt)).toContain("earlier thinking"); }); + test("the summary sends a native search as a client pair and the thinking after it not at all (#5887)", async () => { + // The summary request declares no tools; native server-tool blocks there are unproven. + const { h, args } = await summarySetup(); + const opus = "anthropic:claude-opus-5-5"; + const requests: LanguageModelV3CallOptions[] = []; + const sdkModel = new MockLanguageModelV3({ + doStream: (request) => { + requests.push(request); + return Promise.resolve({ stream: simulateReadableStream({ chunks: modelChunks() }) }); + }, + }); + const anthropicConfig = { apiKeySet: true, isEnabled: true, isConfigured: true }; + spyOn(h.aiService, "getProvidersConfig").mockReturnValue({ anthropic: anthropicConfig }); + spyOn(h.aiService, "createModelWithPinnedOptions").mockResolvedValue( + Ok({ + ...pinnedSummaryModel(sdkModel, opus), + wireProviderName: "anthropic", + optionsRouteProvider: "anthropic" as const, + optionsProvidersConfig: { anthropic: anthropicConfig }, + }) + ); + const nativeSearch: MuxMessage = { + id: "assistant-search", + role: "assistant", + metadata: { thinkingLevel: "high" }, + parts: [ + { + type: "reasoning", + text: "plan", + providerOptions: { anthropic: { signature: "sig-plan" } }, + }, + { + type: "dynamic-tool", + toolCallId: "srvtoolu_1", + toolName: "web_search", + state: "output-available", + input: { query: "xum" }, + providerExecuted: true, + output: [ + { + type: "web_search_result", + url: "https://example.com/xum", + title: "Xum", + pageAge: null, + encryptedContent: "enc-1", + }, + ], + }, + { + type: "reasoning", + text: "bound thinking", + providerOptions: { anthropic: { signature: "sig-read" } }, + }, + { type: "text", text: "earlier answer" }, + ], + }; + const head = [...args.head, nativeSearch, createMuxMessage("next", "user", "next step")]; + await summarizeContinuousCompaction({ + ...args, + head, + receiptRows: head, + compactOptions: { ...args.compactOptions, model: opus, thinkingLevel: "high" }, + }); + + const parts = requests[0].prompt.flatMap((message) => + Array.isArray(message.content) + ? (message.content as unknown as Array>) + : [] + ); + expect(parts.filter((part) => part.providerExecuted === true)).toEqual([]); + expect(parts.filter((part) => part.type === "tool-call")).toHaveLength(1); + const prompt = JSON.stringify(requests[0].prompt); + expect(prompt).not.toContain("bound thinking"); + expect(prompt).not.toContain("enc-1"); + expect(prompt).toContain("earlier answer"); + }); + test("a receipt in the retained tail keeps the head's Anthropic thinking out of the summary (#5996)", async () => { const { h, args } = await summarySetup(); const opus = "anthropic:claude-opus-5-5"; diff --git a/src/node/services/continuousCompactionSummary.ts b/src/node/services/continuousCompactionSummary.ts index ed2753583dd..eeeebe2631b 100644 --- a/src/node/services/continuousCompactionSummary.ts +++ b/src/node/services/continuousCompactionSummary.ts @@ -160,6 +160,10 @@ export async function summarizeContinuousCompaction(args: { const messages = await prepareMessagesForProvider({ messagesWithSentinel: addInterruptedSentinel(prepared.providerRequestMessages), replayReceiptMessages: args.receiptRows, + // The summary request declares no tools, and native server-tool blocks without the + // declared tool are not proven to be accepted (#5887). Headless and one-shot: sending + // the client pair (and so no thinking after it) costs no cache the turn would reuse. + nativeServerToolReplay: false, effectiveAgentId: "compact", toolNamesForSentinel: [], postCompactionAttachments: null, diff --git a/src/node/services/messagePipeline.ts b/src/node/services/messagePipeline.ts index ca304d0215f..147b7b2da52 100644 --- a/src/node/services/messagePipeline.ts +++ b/src/node/services/messagePipeline.ts @@ -38,6 +38,7 @@ import { } from "@/common/utils/tools/toolCatalog"; import { applyCacheControl, type AnthropicCacheTtl } from "@/common/utils/ai/cacheStrategy"; import { findLatestContextBoundaryIndex } from "@/common/utils/messages/compactionBoundary"; +import { projectAnthropicServerTools } from "@/common/utils/messages/anthropicNativeServerTools"; import { log } from "./log"; /** Options for the full message preparation pipeline. */ @@ -83,6 +84,12 @@ export interface PrepareMessagesOptions { * filter pass the rows from before it (#5886). */ replayReceiptMessages?: MuxMessage[]; + /** + * False: send Anthropic server tools as the client pair even on the Anthropic wire. For + * requests that declare no server tool (the headless compaction summary), where native + * blocks are not proven to be accepted (#5887). + */ + nativeServerToolReplay?: boolean; } /** @@ -129,6 +136,7 @@ export async function prepareMessagesForProvider( workspaceId, deferLoadingToolNames, replayReceiptMessages, + nativeServerToolReplay, } = opts; // --- XumMessage-level transforms --- @@ -191,11 +199,19 @@ export async function prepareMessagesForProvider( // providerMetadata, the only field convertToModelMessages forwards to the request. const messagesWithReasoningReplay = attachReasoningReplayMetadata(messagesWithSdkSafeFileParts); + // #5887: Anthropic server tools replay natively only on the Anthropic wire; other wires + // (whose converters drop Anthropic provider-executed calls) get the client pair. Request-only: + // history keeps the native identity and ciphertext. + const serverTools = projectAnthropicServerTools( + messagesWithReasoningReplay, + providerForMessages === "anthropic" && nativeServerToolReplay !== false + ); + // --- Convert to ModelMessage format --- // Type assertion needed because XumMessage has custom tool parts for interrupted tools // eslint-disable-next-line @typescript-eslint/no-explicit-any, @typescript-eslint/no-unsafe-argument - const rawModelMessages = await convertToModelMessages(messagesWithReasoningReplay as any, { + const rawModelMessages = await convertToModelMessages(serverTools.messages as any, { // Drop unfinished tool calls (input-streaming/input-available) so downstream // transforms only see tool calls that actually produced outputs. ignoreIncompleteToolCalls: true, @@ -234,9 +250,12 @@ export async function prepareMessagesForProvider( // the next turn would 400 and pay the repair again. Keep them out for the rest of the // context segment. Same strip as the one-request repair (stripReasoningReplay), which also // covers adaptive thinking, where transformModelMessages ignores anthropicStripReasoning. + // A demoted server tool with thinking after it (#5887) strips the same way: that thinking is + // bound to native blocks this request does not send. const segmentMessages = providerForMessages === "anthropic" && - hasAnthropicReplayReceipt(replayReceiptMessages ?? messagesWithSentinel) + (serverTools.demotedBeforeThinking || + hasAnthropicReplayReceipt(replayReceiptMessages ?? messagesWithSentinel)) ? stripReasoningReplay(transformedMessages, "anthropic") : transformedMessages; diff --git a/src/node/services/streamManager.continuousCompaction.test.ts b/src/node/services/streamManager.continuousCompaction.test.ts index 540daec5570..aeed67ed22a 100644 --- a/src/node/services/streamManager.continuousCompaction.test.ts +++ b/src/node/services/streamManager.continuousCompaction.test.ts @@ -977,6 +977,64 @@ describe("continuous prefix prepareStep and journal", () => { expect(serialized).not.toContain("removed thinking"); }); + it("prefix replay keeps a native search in place, as the main request does (#5887)", async () => { + const journal = journalFixture(); + journal.preparation.modelString = "anthropic:claude-opus-5-5"; + journal.preparation.effectiveThinkingLevel = "high"; + const nativeSearch = createMuxMessage("native-search", "assistant", "", undefined, [ + { + type: "reasoning", + text: "plan", + providerOptions: { anthropic: { signature: "sig-plan" } }, + }, + { + type: "dynamic-tool", + toolCallId: "srvtoolu_1", + toolName: "web_search", + state: "output-available", + input: { query: "xum" }, + providerExecuted: true, + output: [ + { + type: "web_search_result", + url: "https://example.com/xum", + title: "Xum", + pageAge: null, + encryptedContent: "enc-1", + }, + ], + }, + { + type: "reasoning", + text: "read", + providerOptions: { anthropic: { signature: "sig-read" } }, + }, + { type: "text", text: "kept answer" }, + ]); + journal.prefixSourceRows = [ + journal.boundary, + createMuxMessage("user-1", "user", "search"), + nativeSearch, + createMuxMessage("user-2", "user", "next"), + ]; + const { deferLoadingToolNames: _deferred, ...preparation } = journal.preparation; + const expected = await assemblePromptPayload({ + ...preparation, + workspaceId, + history: journal.prefixSourceRows, + systemMessage: "", + postCompactionAttachments: journal.postCompactionAttachments, + }); + const actual = (await rebuildContinuousPrefix(journal, workspaceId)).filter( + (message) => message.role !== "system" + ); + expect(actual).toEqual(expected.messages.filter((message) => message.role !== "system")); + const serialized = JSON.stringify(actual); + expect(serialized).toContain('"providerExecuted":true'); + expect(serialized).toContain("enc-1"); + expect(serialized).toContain("sig-read"); + }); + it("prefix replay keeps the thinking after a server tool out, tool pairs intact (#5887)", async () => { const journal = journalFixture(); journal.preparation.modelString = "anthropic:claude-opus-5-5"; diff --git a/src/node/services/streamManager.preservedThinking.test.ts b/src/node/services/streamManager.preservedThinking.test.ts index 2d05a3f4c64..0235f3aecfa 100644 --- a/src/node/services/streamManager.preservedThinking.test.ts +++ b/src/node/services/streamManager.preservedThinking.test.ts @@ -8,6 +8,8 @@ import { createMuxMessage, type MuxMessage } from "@/common/types/message"; import { Ok } from "@/common/types/result"; import { prepareMessagesForProvider } from "./messagePipeline"; import { assemblePromptPayload } from "./turnContextAssembler"; +import type { TurnEngineEvent } from "./streamManager"; +import { ANTHROPIC_NATIVE_SERVER_TOOL_MAX_ROW_CIPHERTEXT_CHARS } from "@/constants/anthropicServerTools"; import { createStreamManagerForTests, engineInternals, @@ -147,17 +149,28 @@ const expectedBlocks = [ { type: "tool_use", id: "toolu_1", name: "bash", input: { script: "pwd" } }, ]; -/** - * One step with a server tool between thinking blocks (#5887): thinking "plan", - * server_tool_use web_search, web_search_tool_result, thinking "read", client tool_use. - */ -const serverToolResponse = () => - sse([ - messageStart, - ...thinkingBlock(0, "plan", "sig-plan"), +/** A successful web search result block's content, as the API returns it. */ +const searchResults = [ + { + type: "web_search_result", + url: "https://example.com/xum", + title: "Xum", + encrypted_content: "enc-1", + page_age: null, + }, +]; +/** A failed web search: history does not replay it natively (#5887 scope). */ +const searchError = { type: "web_search_tool_result_error", error_code: "unavailable" }; + +/** The SDK's stream output for a failed search: a server tool history does not replay natively. */ +const nonNativeSearchOutput = { type: "web_search_tool_result_error", errorCode: "unavailable" }; + +/** server_tool_use web_search at `index`, then its web_search_tool_result at `index + 1`. */ +function webSearchBlocks(index: number, content: unknown): unknown[] { + return [ { type: "content_block_start", - index: 1, + index, content_block: { type: "server_tool_use", id: "srvtoolu_1", @@ -165,40 +178,42 @@ const serverToolResponse = () => input: { query: "xum" }, }, }, - { type: "content_block_stop", index: 1 }, - { - type: "content_block_start", - index: 2, - content_block: { - type: "web_search_tool_result", - tool_use_id: "srvtoolu_1", - content: [ - { - type: "web_search_result", - url: "https://example.com/xum", - title: "Xum", - encrypted_content: "enc-1", - page_age: null, - }, - ], - }, - }, - { type: "content_block_stop", index: 2 }, - ...thinkingBlock(3, "read", "sig-read"), + { type: "content_block_stop", index }, { type: "content_block_start", - index: 4, - content_block: { type: "tool_use", id: "toolu_2", name: "bash", input: {} }, + index: index + 1, + content_block: { type: "web_search_tool_result", tool_use_id: "srvtoolu_1", content }, }, - { - type: "content_block_delta", - index: 4, - delta: { type: "input_json_delta", partial_json: '{"script":"pwd"}' }, - }, - { type: "content_block_stop", index: 4 }, - { type: "message_delta", delta: { stop_reason: "tool_use" }, usage: { output_tokens: 5 } }, - { type: "message_stop" }, - ]); + { type: "content_block_stop", index: index + 1 }, + ]; +} + +/** + * One step with a server tool between thinking blocks (#5887): thinking "plan", + * server_tool_use web_search, web_search_tool_result, thinking "read", client tool_use. + */ +const serverToolResponse = + (content: unknown = searchResults) => + () => + sse([ + messageStart, + ...thinkingBlock(0, "plan", "sig-plan"), + ...webSearchBlocks(1, content), + ...thinkingBlock(3, "read", "sig-read"), + { + type: "content_block_start", + index: 4, + content_block: { type: "tool_use", id: "toolu_2", name: "bash", input: {} }, + }, + { + type: "content_block_delta", + index: 4, + delta: { type: "input_json_delta", partial_json: '{"script":"pwd"}' }, + }, + { type: "content_block_stop", index: 4 }, + { type: "message_delta", delta: { stop_reason: "tool_use" }, usage: { output_tokens: 5 } }, + { type: "message_stop" }, + ]); describe("StreamManager - Anthropic preserved thinking replay", () => { test("replays every thinking block from history as the API returned it", async () => { @@ -259,22 +274,27 @@ describe("StreamManager - Anthropic preserved thinking replay", () => { expect(firstAssistantBlocks(bodies[2])).toEqual(expectedBlocks); }); - describe("server tool between thinking blocks (#5887, containment)", () => { - // Xum stores an Anthropic server tool (web_search) as a client tool call without its - // encrypted results, so history replays it as a tool_use/tool_result pair instead of - // the API's server_tool_use + web_search_tool_result blocks. Thinking after the server - // tool is bound to a prefix Xum never sends again. Until native replay exists, such a - // turn writes the thinking-replay receipt, and later requests in the context segment - // send no thinking (removing all thinking blocks is valid per the preserved-thinking - // docs). A scripted fixture proves the drift and the strip, not upstream acceptance. + describe("server tool between thinking blocks (#5887)", () => { + // History replays a successful Anthropic web_search natively (server_tool_use + + // web_search_tool_result, ciphertext included), so the thinking after it stays valid. + // Any other server tool (a failed search here) replays as a client tool_use/tool_result + // pair: thinking after it is bound to a prefix Xum never sends again, so such a turn + // writes the thinking-replay receipt and later requests in the context segment send no + // thinking (removing all thinking blocks is valid per the preserved-thinking docs). + // A scripted fixture proves the request shape, not upstream acceptance. const bashTool = () => tool({ inputSchema: z.object({ script: z.string() }), execute: () => Promise.resolve("/tmp"), }); - async function runServerToolTurn(workspaceId: string, first: () => Response) { - const scripted = scriptedAnthropicModel([first, textResponse, textResponse]); + async function runServerToolTurn( + workspaceId: string, + first: () => Response, + // The turn's later steps; the last textResponse answers the next turn's request. + rest: Array<() => Response> = [textResponse] + ) { + const scripted = scriptedAnthropicModel([first, ...rest, textResponse]); const tools = { // Same cast as production (src/common/utils/tools/tools.ts). web_search: scripted.provider.tools.webSearch_20250305({ maxUses: 5 }) as Tool, @@ -285,8 +305,12 @@ describe("StreamManager - Anthropic preserved thinking replay", () => { createMuxMessage("user-1", "user", "search", { historySequence: 0 }) ); if (!seeded.success) throw new Error(seeded.error); + const events: TurnEngineEvent[] = []; const streamManager = createStreamManagerForTests(historyService, { streamText: fakeStreamText((options) => aiSdk.streamText(options)), + eventSink: (event) => { + events.push(event); + }, }); const { messageId } = await runTurnForTests(streamManager, { workspaceId, @@ -295,7 +319,7 @@ describe("StreamManager - Anthropic preserved thinking replay", () => { messages: [{ role: "user", content: "search" }], tools, }); - return { ...scripted, tools, messageId }; + return { ...scripted, tools, messageId, events }; } /** The next turn's request body, rebuilt from committed history the production way. */ @@ -364,9 +388,9 @@ describe("StreamManager - Anthropic preserved thinking replay", () => { expect(uses).toBeGreaterThan(0); } - test("a server tool before thinking turns thinking replay off for the segment", async () => { + test("a non-native server tool before thinking turns thinking replay off for the segment", async () => { const workspaceId = "preserved-thinking-server-tool"; - const scripted = await runServerToolTurn(workspaceId, serverToolResponse); + const scripted = await runServerToolTurn(workspaceId, serverToolResponse(searchError)); // The SDK's own in-turn replay keeps the API's block order, native result included: // that is the prefix the "read" signature is bound to. @@ -383,6 +407,12 @@ describe("StreamManager - Anthropic preserved thinking replay", () => { const row = history.data.find((message) => message.id === scripted.messageId); expect(row?.metadata?.partial).not.toBe(true); expect(row?.metadata?.anthropicThinkingReplay).toBe("off"); + // Stored as the client pair it was before native replay: an older build's SDK would + // throw on a flagged result it cannot validate. + const search = row?.parts.find( + (part) => part.type === "dynamic-tool" && part.toolCallId === "srvtoolu_1" + ); + expect(search?.type === "dynamic-tool" && search.providerExecuted).toBeFalsy(); const next = await nextTurnBody(workspaceId, scripted); expect(thinkingTypes(next)).toEqual([]); @@ -390,6 +420,132 @@ describe("StreamManager - Anthropic preserved thinking replay", () => { expect(JSON.stringify(next)).toContain("done"); }); + test("a response that opens with a native search replays in the API's order", async () => { + // The shape a live Opus 5.5 run returned (#5887): server_tool_use, its result, then + // signed thinking, then text, with no thinking before the search. The thinking is bound to + // the search before it, so the replay must not move it to the front of the message. + const workspaceId = "preserved-thinking-search-first"; + const searchFirst = () => + sse([ + messageStart, + ...webSearchBlocks(0, searchResults), + ...thinkingBlock(2, "read", "sig-read"), + { type: "content_block_start", index: 3, content_block: { type: "text", text: "" } }, + { type: "content_block_delta", index: 3, delta: { type: "text_delta", text: "found" } }, + { type: "content_block_stop", index: 3 }, + { + type: "message_delta", + delta: { stop_reason: "end_turn" }, + usage: { output_tokens: 1 }, + }, + { type: "message_stop" }, + ]); + const scripted = await runServerToolTurn(workspaceId, searchFirst, []); + const history = await historyService.getHistoryFromLatestBoundary(workspaceId); + if (!history.success) throw new Error(history.error); + const row = history.data.find((message) => message.id === scripted.messageId); + expect(row?.metadata?.anthropicThinkingReplay).toBeUndefined(); + + const next = await nextTurnBody(workspaceId, scripted); + expect(firstAssistantBlocks(next).map((block) => block.type)).toEqual([ + "server_tool_use", + "web_search_tool_result", + "thinking", + "text", + ]); + }); + + test("a search whose result opens the next step is stored as the client pair", async () => { + // Claude called web_search and a client tool in one parallel group: the response ends + // after both calls, and the API runs the search at the start of the next step + // (platform.claude.com/docs/en/agents-and-tools/tool-use/server-tools). One stored part + // cannot replay the call and its result at their two positions, so the search is stored + // as the client pair and the thinking after it is kept out (receipt). + const workspaceId = "preserved-thinking-split-server-tool"; + const parallelCalls = () => + sse([ + messageStart, + ...thinkingBlock(0, "plan", "sig-plan"), + { + type: "content_block_start", + index: 1, + content_block: { + type: "server_tool_use", + id: "srvtoolu_1", + name: "web_search", + input: { query: "xum" }, + }, + }, + { type: "content_block_stop", index: 1 }, + { + type: "content_block_start", + index: 2, + content_block: { type: "tool_use", id: "toolu_2", name: "bash", input: {} }, + }, + { + type: "content_block_delta", + index: 2, + delta: { type: "input_json_delta", partial_json: '{"script":"pwd"}' }, + }, + { type: "content_block_stop", index: 2 }, + { + type: "message_delta", + delta: { stop_reason: "tool_use" }, + usage: { output_tokens: 5 }, + }, + { type: "message_stop" }, + ]); + const resultThenThinking = () => + sse([ + messageStart, + { + type: "content_block_start", + index: 0, + content_block: { + type: "web_search_tool_result", + tool_use_id: "srvtoolu_1", + content: searchResults, + }, + }, + { type: "content_block_stop", index: 0 }, + ...thinkingBlock(1, "read", "sig-read"), + { type: "content_block_start", index: 2, content_block: { type: "text", text: "" } }, + { type: "content_block_delta", index: 2, delta: { type: "text_delta", text: "done" } }, + { type: "content_block_stop", index: 2 }, + { + type: "message_delta", + delta: { stop_reason: "end_turn" }, + usage: { output_tokens: 1 }, + }, + { type: "message_stop" }, + ]); + const scripted = await runServerToolTurn(workspaceId, parallelCalls, [resultThenThinking]); + expect(scripted.bodies).toHaveLength(2); + + const history = await historyService.getHistoryFromLatestBoundary(workspaceId); + if (!history.success) throw new Error(history.error); + const row = history.data.find((message) => message.id === scripted.messageId); + expect(row?.metadata?.partial).not.toBe(true); + expect(row?.metadata?.anthropicThinkingReplay).toBe("off"); + const search = row?.parts.find( + (part) => part.type === "dynamic-tool" && part.toolCallId === "srvtoolu_1" + ); + expect(search?.type === "dynamic-tool" && search.state).toBe("output-available"); + expect(search?.type === "dynamic-tool" && search.providerExecuted).toBeFalsy(); + expect(JSON.stringify(search)).not.toContain("enc-1"); + // The renderer gets the stored output: the dropped ciphertext does not cross IPC. + const searchEnd = scripted.events.find( + (event) => event.type === "tool-call-end" && event.toolCallId === "srvtoolu_1" + ); + expect(searchEnd).toBeDefined(); + expect(JSON.stringify(searchEnd)).not.toContain("enc-1"); + + const next = await nextTurnBody(workspaceId, scripted); + expect(thinkingTypes(next)).toEqual([]); + expectPairedTools(next); + expect(JSON.stringify(next)).not.toContain("server_tool_use"); + }); + test("the first partial that holds the bound thinking carries the receipt through a crash", async () => { const workspaceId = "preserved-thinking-server-tool-partial"; const written: MuxMessage[] = []; @@ -402,7 +558,7 @@ describe("StreamManager - Anthropic preserved thinking replay", () => { ); let scripted: Awaited>; try { - scripted = await runServerToolTurn(workspaceId, serverToolResponse); + scripted = await runServerToolTurn(workspaceId, serverToolResponse(searchError)); } finally { writePartial.mockRestore(); } @@ -559,7 +715,7 @@ describe("StreamManager - Anthropic preserved thinking replay", () => { type: "tool-result", toolCallId: "srvtoolu_1", toolName: "web_search", - output: [], + output: nonNativeSearchOutput, providerExecuted: true, }, // A retry that does not preserve parts drops the server-tool part, while @@ -619,7 +775,7 @@ describe("StreamManager - Anthropic preserved thinking replay", () => { type: "tool-result", toolCallId: "srvtoolu_1", toolName: "web_search", - output: [], + output: nonNativeSearchOutput, providerExecuted: true, }, REFUSAL_FINISH, @@ -697,7 +853,7 @@ describe("StreamManager - Anthropic preserved thinking replay", () => { type: "tool-result", toolCallId: "srvtoolu_1", toolName: "web_search", - output: [], + output: nonNativeSearchOutput, providerExecuted: true, }, { type: "reasoning-delta", text: "read" }, @@ -722,6 +878,521 @@ describe("StreamManager - Anthropic preserved thinking replay", () => { expect(seen).toEqual([undefined, "off"]); }); + test("replays a successful web search natively, in place, and writes no receipt", async () => { + const workspaceId = "preserved-thinking-native-search"; + const scripted = await runServerToolTurn(workspaceId, serverToolResponse()); + // The SDK's in-turn replay is the reference: the API's block order, ciphertext included. + const reference = firstAssistantBlocks(scripted.bodies[1]); + expect(reference.map((block) => block.type)).toEqual([ + "thinking", + "server_tool_use", + "web_search_tool_result", + "thinking", + "tool_use", + ]); + + const history = await historyService.getHistoryFromLatestBoundary(workspaceId); + if (!history.success) throw new Error(history.error); + const row = history.data.find((message) => message.id === scripted.messageId); + expect(row?.metadata?.anthropicThinkingReplay).toBeUndefined(); + const search = row?.parts.find( + (part) => part.type === "dynamic-tool" && part.toolCallId === "srvtoolu_1" + ); + expect(search?.type === "dynamic-tool" && search.providerExecuted).toBe(true); + expect(JSON.stringify(search)).toContain("enc-1"); + + // Each thinking signature binds to everything before it: the next turn sends the same + // blocks, so the "read" thinking stays valid and nothing is stripped. + const next = await nextTurnBody(workspaceId, scripted); + expect(firstAssistantBlocks(next)).toEqual(reference); + expect(thinkingTypes(next)).toEqual(["thinking", "thinking"]); + }); + + test("a search result with a null title replays natively without failing the request", async () => { + // The SDK stream omits a null title, but its replay schema requires the key. + const workspaceId = "preserved-thinking-null-title"; + const untitled = [{ ...searchResults[0], title: null }]; + const scripted = await runServerToolTurn(workspaceId, () => + sse([ + messageStart, + ...thinkingBlock(0, "plan", "sig-plan"), + ...webSearchBlocks(1, untitled), + ...thinkingBlock(3, "read", "sig-read"), + { type: "content_block_start", index: 4, content_block: { type: "text", text: "" } }, + { type: "content_block_delta", index: 4, delta: { type: "text_delta", text: "found" } }, + { type: "content_block_stop", index: 4 }, + { + type: "message_delta", + delta: { stop_reason: "end_turn" }, + usage: { output_tokens: 3 }, + }, + { type: "message_stop" }, + ]) + ); + // Stored with the null the API returned, so any reader's SDK can validate it. + const history = await historyService.getHistoryFromLatestBoundary(workspaceId); + if (!history.success) throw new Error(history.error); + const search = history.data + .find((message) => message.id === scripted.messageId) + ?.parts.find((part) => part.type === "dynamic-tool" && part.toolCallId === "srvtoolu_1"); + expect( + search?.type === "dynamic-tool" && search.state === "output-available" && search.output + ).toEqual([ + { + type: "web_search_result", + url: "https://example.com/xum", + title: null, + pageAge: null, + encryptedContent: "enc-1", + }, + ]); + const next = await nextTurnBody(workspaceId, scripted); + const result = firstAssistantBlocks(next).find( + (block) => block.type === "web_search_tool_result" + ); + expect(result?.content).toEqual([ + { + type: "web_search_result", + url: "https://example.com/xum", + title: null, + encrypted_content: "enc-1", + page_age: null, + }, + ]); + }); + + test("three turns resend each request's blocks unchanged, so the cached prefix holds", async () => { + // Prompt caching reuses the longest unchanged prefix (system, tools, then messages). + // Each request must equal the previous one plus the previous reply, as the API returned + // it, plus the new user turn. cache_control is checked on its own: Xum moves the message + // breakpoint to the newest block on every request, which changes the JSON, not content. + const workspaceId = "preserved-thinking-cache-prefix"; + const thinkingText = (n: number) => () => + sse([ + messageStart, + ...thinkingBlock(0, `t${n}`, `sig-t${n}`), + { type: "content_block_start", index: 1, content_block: { type: "text", text: "" } }, + { + type: "content_block_delta", + index: 1, + delta: { type: "text_delta", text: `answer ${n}` }, + }, + { type: "content_block_stop", index: 1 }, + { + type: "message_delta", + delta: { stop_reason: "end_turn" }, + usage: { output_tokens: 1 }, + }, + { type: "message_stop" }, + ]); + const searchThenText = () => + sse([ + messageStart, + ...thinkingBlock(0, "plan", "sig-plan"), + ...webSearchBlocks(1, searchResults), + ...thinkingBlock(3, "read", "sig-read"), + { type: "content_block_start", index: 4, content_block: { type: "text", text: "" } }, + { type: "content_block_delta", index: 4, delta: { type: "text_delta", text: "found" } }, + { type: "content_block_stop", index: 4 }, + { + type: "message_delta", + delta: { stop_reason: "end_turn" }, + usage: { output_tokens: 3 }, + }, + { type: "message_stop" }, + ]); + // What the API returned on each turn, in request-block form. + const replies = [ + [ + { type: "thinking", thinking: "plan", signature: "sig-plan" }, + { + type: "server_tool_use", + id: "srvtoolu_1", + name: "web_search", + input: { query: "xum" }, + }, + { + type: "web_search_tool_result", + tool_use_id: "srvtoolu_1", + content: [ + { + type: "web_search_result", + url: "https://example.com/xum", + title: "Xum", + encrypted_content: "enc-1", + page_age: null, + }, + ], + }, + { type: "thinking", thinking: "read", signature: "sig-read" }, + { type: "text", text: "found" }, + ], + [ + { type: "thinking", thinking: "t2", signature: "sig-t2" }, + { type: "text", text: "answer 2" }, + ], + ]; + const scripted = scriptedAnthropicModel([searchThenText, thinkingText(2), thinkingText(3)]); + const tools = { + web_search: scripted.provider.tools.webSearch_20250305({ maxUses: 5 }) as Tool, + }; + const streamManager = createStreamManagerForTests(historyService, { + streamText: fakeStreamText((options) => aiSdk.streamText(options)), + }); + for (const [turn, text] of ["search", "next 2", "next 3"].entries()) { + const appended = await historyService.appendToHistory( + workspaceId, + createMuxMessage(`user-${turn}`, "user", text, { historySequence: turn * 2 }) + ); + if (!appended.success) throw new Error(appended.error); + const history = await historyService.getHistoryFromLatestBoundary(workspaceId); + if (!history.success) throw new Error(history.error); + const payload = await assemblePromptPayload({ + history: history.data, + systemMessage: "system", + modelString: "anthropic:claude-opus-5-5", + providerForMessages: "anthropic", + effectiveThinkingLevel: "high", + effectiveAgentId: "exec", + toolNamesForSentinel: [], + workspaceId, + }); + await runTurnForTests(streamManager, { + workspaceId, + messageId: `assistant-${turn}`, + historySequence: turn * 2 + 1, + model: scripted.model, + modelString: "anthropic:claude-opus-5-5", + messages: payload.messages.filter((message) => message.role !== "system"), + tools, + }); + } + + const bodies = scripted.bodies as unknown as Array>; + expect(bodies).toHaveLength(3); + const withoutCacheControl = (value: unknown): unknown => + JSON.parse( + JSON.stringify(value, (key, item: unknown) => + key === "cache_control" ? undefined : item + ) + ); + const cacheControlPaths = (value: unknown, path = ""): string[] => { + if (Array.isArray(value)) + return value.flatMap((item, i) => cacheControlPaths(item, `${path}[${i}]`)); + if (typeof value !== "object" || value === null) return []; + return Object.entries(value).flatMap(([key, item]) => + key === "cache_control" ? [path] : cacheControlPaths(item, `${path}.${key}`) + ); + }; + for (let turn = 1; turn < bodies.length; turn++) { + const previous = bodies[turn - 1]; + const current = bodies[turn]; + const previousMessages = previous.messages as unknown[]; + const currentMessages = current.messages as unknown[]; + expect(withoutCacheControl(current.system)).toEqual(withoutCacheControl(previous.system)); + expect(withoutCacheControl(current.tools)).toEqual(withoutCacheControl(previous.tools)); + expect(withoutCacheControl(currentMessages.slice(0, previousMessages.length))).toEqual( + withoutCacheControl(previousMessages) + ); + expect(withoutCacheControl(currentMessages[previousMessages.length])).toEqual({ + role: "assistant", + content: replies[turn - 1], + }); + expect(currentMessages).toHaveLength(previousMessages.length + 2); + } + // The breakpoints: system and tools keep theirs; the message breakpoint sits on the + // newest block only. + const breakpoints = bodies.map((body) => cacheControlPaths(body)); + const lastBlock = (body: Record) => { + const messages = body.messages as Array<{ content: unknown[] }>; + return `.messages[${messages.length - 1}].content[${messages.at(-1)!.content.length - 1}]`; + }; + for (const [index, paths] of breakpoints.entries()) { + const messagePaths = paths.filter((path) => path.startsWith(".messages")); + expect(messagePaths).toEqual([lastBlock(bodies[index])]); + expect(paths.filter((path) => !path.startsWith(".messages"))).toEqual( + breakpoints[0].filter((path) => !path.startsWith(".messages")) + ); + } + }); + + test("another provider's request gets the search as a client pair, without the ciphertext", async () => { + const workspaceId = "preserved-thinking-native-to-openai"; + const scripted = await runServerToolTurn(workspaceId, serverToolResponse()); + const history = await historyService.getHistoryFromLatestBoundary(workspaceId); + if (!history.success) throw new Error(history.error); + const messages = await prepareMessagesForProvider({ + messagesWithSentinel: [ + ...history.data, + createMuxMessage("user-2", "user", "next", { historySequence: 2 }), + ], + effectiveAgentId: "exec", + toolNamesForSentinel: [], + providerForMessages: "openai", + effectiveThinkingLevel: "high", + modelString: "openai:gpt-5.2", + workspaceId, + }); + expect(scripted.messageId).toBeDefined(); + const parts = messages.flatMap((message) => + Array.isArray(message.content) ? (message.content as Array>) : [] + ); + // The OpenAI converter drops provider-executed calls it did not run: the pair keeps the + // search in the transcript. + expect(parts.filter((part) => part.providerExecuted === true)).toEqual([]); + expect( + parts.filter((part) => part.type === "tool-call" && part.toolName === "web_search") + ).toHaveLength(1); + expect( + parts.filter((part) => part.type === "tool-result" && part.toolName === "web_search") + ).toHaveLength(1); + expect(JSON.stringify(messages)).not.toContain("enc-1"); + }); + + test("a native row that lost its ciphertext replays as a client pair, with no thinking", async () => { + // Without the ciphertext the SDK refuses to build the native block (no request at all), + // and the thinking after it is bound to native blocks that are not sent. + const workspaceId = "preserved-thinking-incomplete-native"; + const scripted = scriptedAnthropicModel([textResponse]); + const row = createMuxMessage("assistant-1", "assistant", "", { historySequence: 1 }, [ + { + type: "reasoning", + text: "plan", + providerOptions: { anthropic: { signature: "sig-plan" } }, + }, + { + type: "dynamic-tool", + toolCallId: "srvtoolu_1", + toolName: "web_search", + state: "output-available", + input: { query: "xum" }, + providerExecuted: true, + output: [ + { + type: "web_search_result", + url: "https://example.com/xum", + title: "Xum", + pageAge: null, + }, + ], + }, + { + type: "reasoning", + text: "read", + providerOptions: { anthropic: { signature: "sig-read" } }, + }, + { type: "text", text: "found" }, + ]); + const messages = await prepareMessagesForProvider({ + messagesWithSentinel: [ + createMuxMessage("user-1", "user", "search", { historySequence: 0 }), + row, + createMuxMessage("user-2", "user", "next", { historySequence: 2 }), + ], + effectiveAgentId: "exec", + toolNamesForSentinel: [], + providerForMessages: "anthropic", + effectiveThinkingLevel: "high", + modelString: "anthropic:claude-opus-5-5", + workspaceId, + }); + const replay = aiSdk.streamText({ + model: scripted.model, + messages, + tools: { web_search: scripted.provider.tools.webSearch_20250305({ maxUses: 5 }) as Tool }, + maxRetries: 0, + }); + await replay.consumeStream(); + expect(scripted.bodies).toHaveLength(1); + const body = scripted.bodies[0]; + expect(thinkingTypes(body)).toEqual([]); + expectPairedTools(body); + expect(JSON.stringify(body)).toContain("found"); + }); + + test("searches past the row's ciphertext budget are stored as the client pair", async () => { + const workspaceId = "preserved-thinking-row-ciphertext-budget"; + const search = (index: number, id: string, ciphertext: string) => [ + { + type: "content_block_start", + index, + content_block: { type: "server_tool_use", id, name: "web_search", input: { query: id } }, + }, + { type: "content_block_stop", index }, + { + type: "content_block_start", + index: index + 1, + content_block: { + type: "web_search_tool_result", + tool_use_id: id, + // One result holds at most 12,000 chars (the generic sanitizer bound). + content: Array.from({ length: Math.ceil(ciphertext.length / 12_000) }, (_, i) => ({ + ...searchResults[0], + encrypted_content: ciphertext.slice(i * 12_000, (i + 1) * 12_000), + })), + }, + }, + { type: "content_block_stop", index: index + 1 }, + ]; + // Two searches that fit the budget only one at a time, then thinking after them. + const half = "h".repeat(ANTHROPIC_NATIVE_SERVER_TOOL_MAX_ROW_CIPHERTEXT_CHARS / 2 + 1); + const twoSearches = () => + sse([ + messageStart, + ...thinkingBlock(0, "plan", "sig-plan"), + ...search(1, "srvtoolu_1", half), + ...search(3, "srvtoolu_2", half), + ...thinkingBlock(5, "read", "sig-read"), + { type: "content_block_start", index: 6, content_block: { type: "text", text: "" } }, + { type: "content_block_delta", index: 6, delta: { type: "text_delta", text: "done" } }, + { type: "content_block_stop", index: 6 }, + { + type: "message_delta", + delta: { stop_reason: "end_turn" }, + usage: { output_tokens: 1 }, + }, + { type: "message_stop" }, + ]); + const scripted = await runServerToolTurn(workspaceId, twoSearches, []); + + const history = await historyService.getHistoryFromLatestBoundary(workspaceId); + if (!history.success) throw new Error(history.error); + const row = history.data.find((message) => message.id === scripted.messageId); + const stored = (row?.parts ?? []).flatMap((part) => + part.type === "dynamic-tool" ? [[part.toolCallId, part.providerExecuted === true]] : [] + ); + expect(stored).toEqual([ + ["srvtoolu_1", true], + ["srvtoolu_2", false], + ]); + expect(JSON.stringify(row).length).toBeLessThan( + ANTHROPIC_NATIVE_SERVER_TOOL_MAX_ROW_CIPHERTEXT_CHARS + ); + // The second search is demoted with thinking after it: the receipt keeps it out. + expect(row?.metadata?.anthropicThinkingReplay).toBe("off"); + }); + + test("a native search replays its ciphertext byte for byte", async () => { + // The signed thinking after the search is bound to the exact ciphertext. 12,000 chars is + // the longest value the generic provider-output sanitizer leaves unchanged. + const workspaceId = "preserved-thinking-long-ciphertext"; + const ciphertext = "c".repeat(12_000); + const scripted = scriptedAnthropicModel([textResponse]); + const row = createMuxMessage("assistant-1", "assistant", "", { historySequence: 1 }, [ + { + type: "reasoning", + text: "plan", + providerOptions: { anthropic: { signature: "sig-plan" } }, + }, + { + type: "dynamic-tool", + toolCallId: "srvtoolu_1", + toolName: "web_search", + state: "output-available", + input: { query: "xum" }, + providerExecuted: true, + output: [ + { + type: "web_search_result", + url: "https://example.com/xum", + title: "Xum", + pageAge: null, + encryptedContent: ciphertext, + }, + ], + }, + { + type: "reasoning", + text: "read", + providerOptions: { anthropic: { signature: "sig-read" } }, + }, + { type: "text", text: "found" }, + ]); + const messages = await prepareMessagesForProvider({ + messagesWithSentinel: [ + createMuxMessage("user-1", "user", "search", { historySequence: 0 }), + row, + createMuxMessage("user-2", "user", "next", { historySequence: 2 }), + ], + effectiveAgentId: "exec", + toolNamesForSentinel: [], + providerForMessages: "anthropic", + effectiveThinkingLevel: "high", + modelString: "anthropic:claude-opus-5-5", + workspaceId, + }); + const replay = aiSdk.streamText({ + model: scripted.model, + messages, + tools: { web_search: scripted.provider.tools.webSearch_20250305({ maxUses: 5 }) as Tool }, + maxRetries: 0, + }); + await replay.consumeStream(); + expect(scripted.bodies).toHaveLength(1); + const blocks = firstAssistantBlocks(scripted.bodies[0]); + expect(blocks.map((block) => block.type)).toEqual([ + "thinking", + "server_tool_use", + "web_search_tool_result", + "thinking", + "text", + ]); + const content = blocks[2].content as Array<{ encrypted_content: string }>; + expect(content[0].encrypted_content).toBe(ciphertext); + }); + + test("a client-pair row from before native replay keeps today's request shape", async () => { + const workspaceId = "preserved-thinking-legacy-row"; + const scripted = scriptedAnthropicModel([textResponse]); + const legacy = createMuxMessage("assistant-1", "assistant", "", { historySequence: 1 }, [ + { + type: "reasoning", + text: "plan", + providerOptions: { anthropic: { signature: "sig-plan" } }, + }, + { + type: "dynamic-tool", + toolCallId: "srvtoolu_1", + toolName: "web_search", + state: "output-available", + input: { query: "xum" }, + output: [ + { + type: "web_search_result", + url: "https://example.com/xum", + title: "Xum", + pageAge: null, + }, + ], + }, + { + type: "reasoning", + text: "read", + providerOptions: { anthropic: { signature: "sig-read" } }, + }, + { type: "text", text: "found" }, + ]); + const messages = await prepareMessagesForProvider({ + messagesWithSentinel: [ + createMuxMessage("user-1", "user", "search", { historySequence: 0 }), + legacy, + createMuxMessage("user-2", "user", "next", { historySequence: 2 }), + ], + effectiveAgentId: "exec", + toolNamesForSentinel: [], + providerForMessages: "anthropic", + effectiveThinkingLevel: "high", + modelString: "anthropic:claude-opus-5-5", + workspaceId, + }); + const replay = aiSdk.streamText({ model: scripted.model, messages, maxRetries: 0 }); + await replay.consumeStream(); + const body = scripted.bodies[0]; + // As on main: a client pair, and the thinking is still sent (no receipt, no native part). + expectPairedTools(body); + expect(thinkingTypes(body)).toEqual(["thinking", "thinking"]); + }); + test("a server tool with no thinking after it keeps thinking replay on", async () => { const workspaceId = "preserved-thinking-server-tool-last"; // thinking "plan", server tool, then the answer: no block is bound to the native prefix. diff --git a/src/node/services/streamManager.ts b/src/node/services/streamManager.ts index 3b7f9a99ac9..9e24d9a0d59 100644 --- a/src/node/services/streamManager.ts +++ b/src/node/services/streamManager.ts @@ -126,7 +126,12 @@ import type { SessionUsageService } from "./sessionUsageService"; import { createDisplayUsage } from "@/common/utils/tokens/displayUsage"; import { extractToolMediaAsUserMessagesFromModelMessages } from "@/node/utils/messages/extractToolMediaAsUserMessagesFromModelMessages"; import { neutralizeAgentEnvelopeLookalikesInModelToolParts } from "@/node/utils/messages/neutralizeAgentEnvelopeLookalikesForProvider"; -import { stripEncryptedContent } from "@/node/utils/messages/stripEncryptedContent"; +import { stripEncryptedContent } from "@/common/utils/messages/stripEncryptedContent"; +import { + isNativeAnthropicReplayable, + rowCiphertextChars, + toStoredServerToolPart, +} from "@/common/utils/messages/anthropicNativeServerTools"; import { countAnthropicInputTransformations } from "@/node/utils/messages/anthropicInputTransformations"; import { stripWorkflowRunRecordsFromModelMessages } from "@/node/utils/messages/stripWorkflowRunRecordsFromModelMessages"; import { stripAnthropicReasoning } from "@/browser/utils/messages/modelMessageTransform"; @@ -922,9 +927,9 @@ interface WorkspaceStreamInfo { // original start timestamp even after they gain output. toolCompletionTimestamps: Map; - // Anthropic server tools (web_search, ...) called in this row (#5887). Xum stores them as - // client tool calls without their encrypted results, so history replays them as a - // tool_use/tool_result pair: thinking after one is bound to a prefix never sent again. + // Anthropic server tools (web_search, ...) called in this row (#5887). History replays + // only successful web searches natively; any other one replays as a tool_use/tool_result + // pair, so thinking after it is bound to a prefix never sent again. anthropicServerToolCallIds?: Set; // Workflow tools can create the durable run before their stream part is stored. Keep the exact @@ -2188,8 +2193,9 @@ export class StreamManager { } /** - * #5887 containment: thinking that follows an Anthropic server tool in this row cannot - * replay as the API returned it (see anthropicServerToolCallIds). Write the same receipt + * #5887 containment: thinking that follows an Anthropic server tool that history does not + * replay natively (anything but a successful web_search, see isNativeAnthropicReplayable) + * cannot replay as the API returned it. Write the same receipt * as the signature repair, so later requests in the context segment send no thinking * (removing every thinking block is valid; a gap or a put-back block is not). It rides * on the partial write that stores this reasoning part (every reasoning append schedules @@ -2204,13 +2210,19 @@ export class StreamManager { // provider keeps the parts (and the IDs) but its reasoning replays nothing Anthropic. if (!isAnthropicMessagesModel(streamInfo.request.model)) return; // Read the parts, not only the ID set: a retry that drops parts drops the server tool too. + // A server tool that history replays natively keeps the thinking after it valid (#5887), + // so only the others count. const followsServerTool = streamInfo.parts.some( - (part) => part.type === "dynamic-tool" && serverToolIds.has(part.toolCallId) + (part) => + part.type === "dynamic-tool" && + serverToolIds.has(part.toolCallId) && + !isNativeAnthropicReplayable(part) ); if (!followsServerTool) { - // Every recorded server tool's part is gone. The tool-call case records an ID and - // stores its part before any later reasoning, so the IDs are stale: drop them, or - // every later reasoning delta would scan all parts again (perf). + // Every recorded server tool's part is gone or replays natively (a stored result does + // not change). The tool-call case records an ID and stores its part before any later + // reasoning, so the IDs are spent: drop them, or every later reasoning delta would scan + // all parts again (perf). streamInfo.anthropicServerToolCallIds = undefined; return; } @@ -3540,16 +3552,32 @@ export class StreamManager { ? this.toolCallDisplayRegistry.take(streamInfo.executionScope, toolCallId) : undefined; + // The output the renderer receives: the stored one, so ciphertext the stored part dropped + // never crosses IPC either. + let emittedOutput = output; if (existingPartIndex !== -1) { const existingPart = streamInfo.parts[existingPartIndex]; if (existingPart.type === "dynamic-tool") { - streamInfo.parts[existingPartIndex] = { - ...existingPart, - ...(pendingAttachment != null ? { workflowRun: pendingAttachment } : {}), - ...(mcpServer ? { mcpServer } : {}), - state: "output-available" as const, - output, - }; + // A provider-executed part keeps its native identity only when history can replay + // it natively (#5887); any other result is stored as the client pair it was before. + const stored = toStoredServerToolPart( + { + ...existingPart, + ...(pendingAttachment != null ? { workflowRun: pendingAttachment } : {}), + ...(mcpServer ? { mcpServer } : {}), + state: "output-available" as const, + output, + }, + { + // Nothing arrived between the call and its result (see toStoredServerToolPart). + resultFollowsCall: existingPartIndex === streamInfo.parts.length - 1, + // The call part itself holds no output yet, so it adds nothing to this sum. + rowCiphertextChars: + existingPart.providerExecuted === true ? rowCiphertextChars(streamInfo.parts) : 0, + } + ); + streamInfo.parts[existingPartIndex] = stored; + if (stored.state === "output-available") emittedOutput = stored.output; } } else { // Fallback: if the matching tool-call part is missing, still persist output so the UI @@ -3592,7 +3620,7 @@ export class StreamManager { messageId: streamInfo.messageId, toolCallId, toolName, - result: output, + result: emittedOutput, ...(mcpServer ? { mcpServer } : {}), ...(providerExecuted === true ? { providerExecuted: true } : {}), timestamp: completionTimestamp, @@ -4760,10 +4788,10 @@ export class StreamManager { toolName: part.toolName, input: part.input, }); - if ( + const anthropicServerTool = part.providerExecuted === true && - isAnthropicMessagesModel(streamInfo.request.model) - ) { + isAnthropicMessagesModel(streamInfo.request.model); + if (anthropicServerTool) { (streamInfo.anthropicServerToolCallIds ??= new Set()).add(part.toolCallId); } @@ -4780,6 +4808,8 @@ export class StreamManager { // eslint-disable-next-line @typescript-eslint/no-unsafe-assignment input: part.input, timestamp: nextPartTimestamp(streamInfo), + // History replays Anthropic server tools natively (#5887). + ...(anthropicServerTool ? { providerExecuted: true } : {}), }; // Emit using shared logic (ensures replay consistency) @@ -4807,9 +4837,21 @@ export class StreamManager { providerExecuted?: boolean; }; - // Strip encrypted content from web search results before storing + // Strip encrypted content from web search results before storing, except for + // Anthropic server tools: native replay must send the ciphertext back as the API + // returned it, or the thinking after the search no longer matches (#5887). + // Keyed on the stored part's flag, so ciphertext is never kept on a part that + // replays as a client pair (an orphan result has no stored call part). + const keepCiphertext = streamInfo.parts.some( + (stored) => + stored.type === "dynamic-tool" && + stored.toolCallId === toolResultPart.toolCallId && + stored.providerExecuted === true + ); const strippedOutput = stripInternalToolResultFields( - stripEncryptedContent(toolResultPart.output) + keepCiphertext + ? toolResultPart.output + : stripEncryptedContent(toolResultPart.output) ); // Tool call completed successfully