返回源码地图

packages/ai/src/api/anthropic-messages.ts

v1.0.0 · a13d35a742c6 · 04/20:中途工具声明适配与版本差异;非全部provider分支审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1import Anthropic, { type ClientOptions } from "@anthropic-ai/sdk";
2import type {
3 BetaStopReason,
4 BetaThinkingDroppedInputTransformation,
5 BetaTool,
6 BetaCacheControlEphemeral as CacheControlEphemeral,
7 BetaContentBlockParam as ContentBlockParam,
8 MessageCreateParamsStreaming,
9 BetaMessageParam as MessageParam,
10 BetaRawMessageStreamEvent as RawMessageStreamEvent,
11 BetaRefusalStopDetails as RefusalStopDetails,
12} from "@anthropic-ai/sdk/resources/beta/messages/messages.js";
13import {
14 ANTHROPIC_FEDERATION_RULE_ID_ENV,
15 ANTHROPIC_IDENTITY_TOKEN_FILE_ENV,
16 ANTHROPIC_ORGANIZATION_ID_ENV,
17 ANTHROPIC_SERVICE_ACCOUNT_ID_ENV,
18 ANTHROPIC_WORKSPACE_ID_ENV,
19} from "../env-api-keys.ts";
20import { calculateCost } from "../models.ts";
21import type {
22 Api,
23 AssistantMessage,
24 CacheRetention,
25 ImageContent,
26 Message,
27 Model,
28 ProviderEnv,
29 ProviderHeaders,
30 SimpleStreamOptions,
31 StopReason,
32 StreamFunction,
33 StreamOptions,
34 TextContent,
35 ThinkingContent,
36 Tool,
37 ToolCall,
38 ToolResultMessage,
39} from "../types.ts";
40import { appendAssistantMessageDiagnostic } from "../utils/diagnostics.ts";
41import { AssistantMessageEventStream } from "../utils/event-stream.ts";
42import { headersToRecord } from "../utils/headers.ts";
43import { parseJsonWithRepair, parseStreamingJson } from "../utils/json-parse.ts";
44import { getPiUserAgent } from "../utils/pi-user-agent.ts";
45import { getProviderEnvValue } from "../utils/provider-env.ts";
46import { retryProviderRequest } from "../utils/provider-retry.ts";
47import { sanitizeSurrogates } from "../utils/sanitize-unicode.ts";
48import { getSystemMessageText, renderSystemMessageUpdate } from "../utils/text.ts";
49import {
50 getCurrentTools,
51 getDeclaredTools,
52 getInitialSystemMessage,
53 hasToolRedefinitions,
54 resolveTranscript,
55 type TranscriptContext,
56} from "../utils/transcript.ts";
57
58import {
59 getJsonSchemaToolParameters,
60 resolveJsonSchemaStrictSampling,
61 type UnsupportedStrictSchemaKeywordCheck,
62} from "./constrained-sampling.ts";
63import { buildCopilotDynamicHeaders, hasCopilotVisionInput } from "./github-copilot-headers.ts";
64import { adjustMaxTokensForThinking, buildBaseOptions, clampMaxTokensToContext } from "./simple-options.ts";
65import { transformMessages } from "./transform-messages.ts";
66
67/**
68 * Resolve cache retention preference.
69 * Defaults to "short" and uses PI_CACHE_RETENTION for backward compatibility.
70 */
71function resolveCacheRetention(cacheRetention?: CacheRetention, env?: ProviderEnv): CacheRetention {
72 if (cacheRetention) {
73 return cacheRetention;
74 }
75 if (getProviderEnvValue("PI_CACHE_RETENTION", env) === "long") {
76 return "long";
77 }
78 return "short";
79}
80
81function getCacheControl(
82 model: Model<"anthropic-messages">,
83 cacheRetention?: CacheRetention,
84 env?: ProviderEnv,
85): { retention: CacheRetention; cacheControl?: CacheControlEphemeral } {
86 const retention = resolveCacheRetention(cacheRetention, env);
87 if (retention === "none") {
88 return { retention };
89 }
90 const ttl = retention === "long" && getAnthropicCompat(model).supportsLongCacheRetention ? "1h" : undefined;
91 return {
92 retention,
93 cacheControl: { type: "ephemeral", ...(ttl && { ttl }) },
94 };
95}
96
97// Stealth mode: Mimic Claude Code's tool naming exactly
98const claudeCodeVersion = "2.1.280";
99
100// Claude Code 2.x tool names (canonical casing)
101// Source: https://cchistory.mariozechner.at/data/prompts-2.1.11.md
102// To update: https://github.com/badlogic/cchistory
103const claudeCodeTools = [
104 "Read",
105 "Write",
106 "Edit",
107 "Bash",
108 "Grep",
109 "Glob",
110 "AskUserQuestion",
111 "EnterPlanMode",
112 "ExitPlanMode",
113 "KillShell",
114 "NotebookEdit",
115 "Skill",
116 "Task",
117 "TaskOutput",
118 "TodoWrite",
119 "WebFetch",
120 "WebSearch",
121];
122
123const ccToolLookup = new Map(claudeCodeTools.map((t) => [t.toLowerCase(), t]));
124
125// Convert tool name to CC canonical casing if it matches (case-insensitive)
126const toClaudeCodeName = (name: string) => ccToolLookup.get(name.toLowerCase()) ?? name;
127const fromClaudeCodeName = (name: string, tools?: Tool[]) => {
128 if (tools && tools.length > 0) {
129 const lowerName = name.toLowerCase();
130 const matchedTool = tools.find((tool) => tool.name.toLowerCase() === lowerName);
131 if (matchedTool) return matchedTool.name;
132 }
133 return name;
134};
135
136/**
137 * Convert content blocks to Anthropic API format
138 */
139function convertContentBlocks(content: (TextContent | ImageContent)[]):
140 | string
141 | Array<
142 | { type: "text"; text: string }
143 | {
144 type: "image";
145 source: {
146 type: "base64";
147 media_type: "image/jpeg" | "image/png" | "image/gif" | "image/webp";
148 data: string;
149 };
150 }
151 > {
152 // If only text blocks, return as concatenated string for simplicity
153 const hasImages = content.some((c) => c.type === "image");
154 if (!hasImages) {
155 return sanitizeSurrogates(content.map((c) => (c as TextContent).text).join("\n"));
156 }
157
158 // If we have images, convert to content block array
159 const blocks = content.map((block) => {
160 if (block.type === "text") {
161 return {
162 type: "text" as const,
163 text: sanitizeSurrogates(block.text),
164 };
165 }
166 return {
167 type: "image" as const,
168 source: {
169 type: "base64" as const,
170 media_type: block.mimeType as "image/jpeg" | "image/png" | "image/gif" | "image/webp",
171 data: block.data,
172 },
173 };
174 });
175
176 // If only images (no text), add placeholder text block
177 const hasText = blocks.some((b) => b.type === "text");
178 if (!hasText) {
179 blocks.unshift({
180 type: "text" as const,
181 text: "(see attached image)",
182 });
183 }
184
185 return blocks;
186}
187
188export type AnthropicEffort = "low" | "medium" | "high" | "xhigh" | "max";
189
190export type AnthropicThinkingDisplay = "summarized" | "omitted";
191
192const FINE_GRAINED_TOOL_STREAMING_BETA = "fine-grained-tool-streaming-2025-05-14";
193const INTERLEAVED_THINKING_BETA = "interleaved-thinking-2025-05-14";
194const SERVER_SIDE_FALLBACK_BETA = "server-side-fallback-2026-07-01";
195const MID_CONVERSATION_OUTPUT_CONFIG_BETA = "mid-conversation-output-config-2026-07-01";
196const THINKING_BINDING_CONTROLS_BETA = "thinking-binding-controls-2026-08-01";
197const MID_CONVERSATION_TOOL_CHANGES_BETA = "mid-conversation-tool-changes-2026-07-01";
198
199/**
200 * Stable deferred tool declared whenever native tool changes are in use. Anthropic adds
201 * hidden prompt scaffolding as soon as any tool has `defer_loading`; declaring this
202 * placeholder from the first request keeps that scaffolding in the cached prefix, so the
203 * first real late tool does not invalidate the cache (measured: full miss without it).
204 * It is never activated and the model cannot see it.
205 */
206const DEFERRED_TOOL_PLACEHOLDER: BetaTool = {
207 name: "__pi_deferred_placeholder__",
208 description: "Reserved placeholder. Never available. Never call this.",
209 input_schema: { type: "object", properties: {}, required: [] },
210 defer_loading: true,
211};
212
213function shouldUseServerSideFallbackBeta(model: Model<"anthropic-messages">): boolean {
214 return (model.compat?.allowedFallbackModels?.length ?? 0) > 0;
215}
216
217function getAnthropicCompat(model: Model<"anthropic-messages">) {
218 const isOpenRouter = model.provider === "openrouter" || model.baseUrl.includes("openrouter.ai");
219 return {
220 supportsEagerToolInputStreaming: model.compat?.supportsEagerToolInputStreaming ?? true,
221 supportsLongCacheRetention: model.compat?.supportsLongCacheRetention ?? true,
222 sendSessionAffinityHeaders: model.compat?.sendSessionAffinityHeaders ?? isOpenRouter,
223 sessionAffinityFormat: model.compat?.sessionAffinityFormat ?? (isOpenRouter ? "openrouter" : undefined),
224 supportsCacheControlOnTools: model.compat?.supportsCacheControlOnTools ?? true,
225 supportsTemperature: model.compat?.supportsTemperature ?? true,
226 allowEmptySignature: model.compat?.allowEmptySignature ?? false,
227 supportsStrictTools: model.compat?.supportsStrictTools ?? false,
228 supportsMidConvoSystemMessages: model.compat?.supportsMidConvoSystemMessages ?? false,
229 supportsMidConvoToolChanges: model.compat?.supportsMidConvoToolChanges ?? false,
230 };
231}
232
233export interface AnthropicOptions extends StreamOptions {
234 /**
235 * Enable extended thinking.
236 * For adaptive thinking models: the model decides when/how much to think.
237 * For older models: uses budget-based thinking with thinkingBudgetTokens.
238 * Default: undefined (thinking is omitted unless `streamSimple()` maps
239 * a simple reasoning level to this option, or callers set it explicitly).
240 */
241 thinkingEnabled?: boolean;
242 /**
243 * Token budget for extended thinking (older models only).
244 * Ignored for adaptive thinking models.
245 * Default: 1024 when `thinkingEnabled` is true and no budget is provided.
246 */
247 thinkingBudgetTokens?: number;
248 /**
249 * Effort level for adaptive thinking models.
250 * Controls how much thinking Claude allocates:
251 * - "max": Always thinks with no constraints (Opus 4.6 only)
252 * - "xhigh": Highest reasoning level (Opus 4.7+, Fable 5)
253 * - "high": Always thinks, deep reasoning
254 * - "medium": Moderate thinking, may skip for simple queries
255 * - "low": Minimal thinking, skips for simple tasks
256 * Ignored for older models.
257 * Default: omitted unless `streamSimple()` maps a simple reasoning
258 * level to this option.
259 */
260 effort?: AnthropicEffort;
261 /**
262 * Controls how thinking content is returned in API responses.
263 * - "summarized": Thinking blocks contain summarized thinking text.
264 * - "omitted": Thinking blocks return an empty thinking field; the encrypted
265 * signature still travels back for multi-turn continuity. Use for faster
266 * time-to-first-text-token when your UI does not surface thinking.
267 *
268 * Note: Anthropic's API default for Claude Opus 4.7 and Claude Mythos Preview
269 * is "omitted". We default to "summarized" here to keep behavior consistent
270 * with older Claude 4 models. Set this explicitly to "omitted" to opt in.
271 * Default: "summarized" when thinking is enabled.
272 */
273 thinkingDisplay?: AnthropicThinkingDisplay;
274 /**
275 * Whether to request the interleaved thinking beta header for non-adaptive
276 * thinking models. Adaptive thinking models have interleaved thinking built in,
277 * so the header is skipped for them regardless of this setting.
278 * Default: true.
279 */
280 interleavedThinking?: boolean;
281 /**
282 * Anthropic tool choice behavior. String values map to Anthropic's built-in
283 * choices; `{ type: "tool", name }` forces a specific tool.
284 * Default: omitted (Anthropic default behavior, currently equivalent to auto).
285 */
286 toolChoice?: "auto" | "any" | "none" | { type: "tool"; name: string };
287 /**
288 * Pre-built Anthropic client instance. When provided, skips internal client
289 * construction entirely. Use this to inject alternative SDK clients such as
290 * `AnthropicVertex` that shares the same messaging API.
291 */
292 client?: Anthropic;
293}
294
295function mergeHeaders(...headerSources: (ProviderHeaders | undefined)[]): ProviderHeaders {
296 const merged: ProviderHeaders = {};
297 for (const headers of headerSources) {
298 if (headers) {
299 Object.assign(merged, headers);
300 }
301 }
302 return merged;
303}
304
305function mergeClientHeaders(...headerSources: (ProviderHeaders | undefined)[]): ProviderHeaders {
306 return mergeHeaders({ "User-Agent": getPiUserAgent() }, ...headerSources);
307}
308
309function hasHeader(headers: ProviderHeaders | undefined, name: string): boolean {
310 if (!headers) return false;
311 const expected = name.toLowerCase();
312 for (const [key, value] of Object.entries(headers)) {
313 if (key.toLowerCase() === expected && value !== null && value.trim().length > 0) return true;
314 }
315 return false;
316}
317
318function hasRequestAuth(apiKey: string | undefined, headers: ProviderHeaders | undefined): boolean {
319 return (
320 !!apiKey ||
321 hasHeader(headers, "authorization") ||
322 hasHeader(headers, "x-api-key") ||
323 hasHeader(headers, "cf-aig-authorization")
324 );
325}
326
327function assertRequestAuth(provider: string, apiKey: string | undefined, headers: ProviderHeaders | undefined): void {
328 if (!hasRequestAuth(apiKey, headers)) throw new Error(`No API key for provider: ${provider}`);
329}
330
331/**
332 * Anthropic SDK client that never runs the SDK's own credential chain
333 * (ANTHROPIC_PROFILE config files, federation env vars). Without this, every
334 * client built with `apiKey: null, authToken: null` for header-owned auth would
335 * also resolve and exchange SDK credentials behind pi's auth resolver.
336 */
337class PiAnthropic extends Anthropic {
338 protected override _shouldResolveDefaultCredentials(): boolean {
339 return false;
340 }
341}
342
343type AnthropicFederationConfig = NonNullable<ClientOptions["config"]>;
344
345/**
346 * Workload identity federation config from the ANTHROPIC_* variables the
347 * Anthropic SDK documents; the SDK performs the token exchange and refresh.
348 * Only for the anthropic provider, since the exchange is an Anthropic API
349 * endpoint, and only when no key or auth header was resolved.
350 */
351function getAnthropicFederation(
352 model: Model<"anthropic-messages">,
353 apiKey: string | undefined,
354 headers: ProviderHeaders | undefined,
355 env: ProviderEnv | undefined,
356): AnthropicFederationConfig | undefined {
357 if (model.provider !== "anthropic" || hasRequestAuth(apiKey, headers)) return undefined;
358 const federationRuleId = getProviderEnvValue(ANTHROPIC_FEDERATION_RULE_ID_ENV, env);
359 const organizationId = getProviderEnvValue(ANTHROPIC_ORGANIZATION_ID_ENV, env);
360 const identityTokenFile = getProviderEnvValue(ANTHROPIC_IDENTITY_TOKEN_FILE_ENV, env);
361 if (!federationRuleId || !organizationId || !identityTokenFile) return undefined;
362 return {
363 organization_id: organizationId,
364 workspace_id: getProviderEnvValue(ANTHROPIC_WORKSPACE_ID_ENV, env),
365 authentication: {
366 type: "oidc_federation",
367 federation_rule_id: federationRuleId,
368 service_account_id: getProviderEnvValue(ANTHROPIC_SERVICE_ACCOUNT_ID_ENV, env),
369 identity_token: { source: "file", path: identityTokenFile },
370 },
371 };
372}
373
374/**
375 * The SDK caches the federated access token per client, but pi creates a client
376 * per request. Keep one client for the current federation config and fetch, and
377 * clone it per request with `withOptions()`, which shares the token cache.
378 */
379let federationClient: { key: string; fetch: typeof globalThis.fetch | undefined; client: Anthropic } | undefined;
380
381interface ServerSentEvent {
382 event: string | null;
383 data: string;
384 raw: string[];
385}
386
387interface SseDecoderState {
388 event: string | null;
389 data: string[];
390 raw: string[];
391}
392
393const ANTHROPIC_MESSAGE_EVENTS: ReadonlySet<string> = new Set([
394 "message_start",
395 "message_delta",
396 "message_stop",
397 "content_block_start",
398 "content_block_delta",
399 "content_block_stop",
400]);
401
402function flushSseEvent(state: SseDecoderState): ServerSentEvent | null {
403 if (!state.event && state.data.length === 0) {
404 return null;
405 }
406
407 const event: ServerSentEvent = {
408 event: state.event,
409 data: state.data.join("\n"),
410 raw: [...state.raw],
411 };
412 state.event = null;
413 state.data = [];
414 state.raw = [];
415 return event;
416}
417
418function decodeSseLine(line: string, state: SseDecoderState): ServerSentEvent | null {
419 if (line === "") {
420 return flushSseEvent(state);
421 }
422
423 state.raw.push(line);
424 if (line.startsWith(":")) {
425 return null;
426 }
427
428 const delimiterIndex = line.indexOf(":");
429 const fieldName = delimiterIndex === -1 ? line : line.slice(0, delimiterIndex);
430 let value = delimiterIndex === -1 ? "" : line.slice(delimiterIndex + 1);
431 if (value.startsWith(" ")) {
432 value = value.slice(1);
433 }
434
435 if (fieldName === "event") {
436 state.event = value;
437 } else if (fieldName === "data") {
438 state.data.push(value);
439 }
440
441 return null;
442}
443
444function nextLineBreakIndex(text: string): number {
445 const carriageReturnIndex = text.indexOf("\r");
446 const newlineIndex = text.indexOf("\n");
447 if (carriageReturnIndex === -1) {
448 return newlineIndex;
449 }
450 if (newlineIndex === -1) {
451 return carriageReturnIndex;
452 }
453 return Math.min(carriageReturnIndex, newlineIndex);
454}
455
456function consumeLine(text: string): { line: string; rest: string } | null {
457 const lineBreakIndex = nextLineBreakIndex(text);
458 if (lineBreakIndex === -1) {
459 return null;
460 }
461
462 let nextIndex = lineBreakIndex + 1;
463 if (text[lineBreakIndex] === "\r" && text[nextIndex] === "\n") {
464 nextIndex += 1;
465 }
466
467 return {
468 line: text.slice(0, lineBreakIndex),
469 rest: text.slice(nextIndex),
470 };
471}
472
473async function* iterateSseMessages(
474 body: ReadableStream<Uint8Array>,
475 signal?: AbortSignal,
476): AsyncGenerator<ServerSentEvent> {
477 const reader = body.getReader();
478 const decoder = new TextDecoder();
479 const state: SseDecoderState = { event: null, data: [], raw: [] };
480 let buffer = "";
481
482 try {
483 while (true) {
484 if (signal?.aborted) {
485 throw new Error("Request was aborted");
486 }
487
488 const { value, done } = await reader.read();
489 if (done) {
490 break;
491 }
492
493 buffer += decoder.decode(value, { stream: true });
494 let consumed = consumeLine(buffer);
495 while (consumed) {
496 buffer = consumed.rest;
497 const event = decodeSseLine(consumed.line, state);
498 if (event) {
499 yield event;
500 }
501 consumed = consumeLine(buffer);
502 }
503 }
504
505 buffer += decoder.decode();
506 let consumed = consumeLine(buffer);
507 while (consumed) {
508 buffer = consumed.rest;
509 const event = decodeSseLine(consumed.line, state);
510 if (event) {
511 yield event;
512 }
513 consumed = consumeLine(buffer);
514 }
515
516 if (buffer.length > 0) {
517 const event = decodeSseLine(buffer, state);
518 if (event) {
519 yield event;
520 }
521 }
522
523 const trailingEvent = flushSseEvent(state);
524 if (trailingEvent) {
525 yield trailingEvent;
526 }
527 } finally {
528 reader.releaseLock();
529 }
530}
531
532async function* iterateAnthropicEvents(
533 response: Response,
534 signal?: AbortSignal,
535): AsyncGenerator<RawMessageStreamEvent> {
536 if (!response.body) {
537 throw new Error("Attempted to iterate over an Anthropic response with no body");
538 }
539
540 let sawMessageStart = false;
541 let sawMessageEnd = false;
542
543 for await (const sse of iterateSseMessages(response.body, signal)) {
544 if (sse.event === "error") {
545 throw new Error(sse.data);
546 }
547
548 if (!ANTHROPIC_MESSAGE_EVENTS.has(sse.event ?? "")) {
549 continue;
550 }
551
552 try {
553 const event = parseJsonWithRepair<RawMessageStreamEvent>(sse.data);
554 if (event.type === "message_start") {
555 sawMessageStart = true;
556 } else if (event.type === "message_stop") {
557 sawMessageEnd = true;
558 }
559 yield event;
560 } catch (error) {
561 const message = error instanceof Error ? error.message : String(error);
562 throw new Error(
563 `Could not parse Anthropic SSE event ${sse.event}: ${message}; data=${sse.data}; raw=${sse.raw.join("\\n")}`,
564 );
565 }
566 }
567
568 if (sawMessageStart && !sawMessageEnd) {
569 throw new Error("Anthropic stream ended before message_stop");
570 }
571}
572
573export const stream: StreamFunction<"anthropic-messages", AnthropicOptions> = (
574 model: Model<"anthropic-messages">,
575 context: TranscriptContext,
576 options?: AnthropicOptions,
577): AssistantMessageEventStream => {
578 const stream = new AssistantMessageEventStream();
579 const normalizedContext = resolveTranscript(context, getAnthropicCompat(model).supportsMidConvoSystemMessages);
580 const currentTools = getCurrentTools(normalizedContext.messages);
581
582 (async () => {
583 const providerThinkingLevel = model.compat?.supportsMidConvoEffort ? (options?.effort ?? "high") : undefined;
584 const output: AssistantMessage = {
585 role: "assistant",
586 content: [],
587 api: model.api as Api,
588 provider: model.provider,
589 model: model.id,
590 ...(providerThinkingLevel === undefined ? {} : { providerThinkingLevel }),
591 usage: {
592 input: 0,
593 output: 0,
594 cacheRead: 0,
595 cacheWrite: 0,
596 totalTokens: 0,
597 cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
598 },
599 stopReason: "pending",
600 timestamp: Date.now(),
601 };
602
603 try {
604 let client: Anthropic;
605 let isOAuth: boolean;
606 let usageModel = model;
607 let inputTransformations: BetaThinkingDroppedInputTransformation[] | undefined;
608
609 if (options?.client) {
610 client = options.client;
611 isOAuth = false;
612 } else {
613 const apiKey = options?.apiKey;
614 const federation = getAnthropicFederation(model, apiKey, options?.headers, options?.env);
615 if (!federation) assertRequestAuth(model.provider, apiKey, options?.headers);
616
617 let copilotDynamicHeaders: Record<string, string> | undefined;
618 if (model.provider === "github-copilot") {
619 const hasImages = hasCopilotVisionInput(normalizedContext.messages);
620 copilotDynamicHeaders = buildCopilotDynamicHeaders({
621 messages: normalizedContext.messages,
622 hasImages,
623 });
624 }
625
626 const cacheRetention = resolveCacheRetention(options?.cacheRetention, options?.env);
627 const cacheSessionId = cacheRetention === "none" ? undefined : options?.sessionId;
628
629 const created = createClient(
630 model,
631 apiKey,
632 options?.headers,
633 options?.fetch,
634 copilotDynamicHeaders,
635 cacheSessionId,
636 federation,
637 );
638 client = created.client;
639 isOAuth = created.isOAuthToken;
640 }
641 let params = buildParams(model, normalizedContext, isOAuth, options);
642 const nextParams = await options?.onPayload?.(params, model);
643 if (nextParams !== undefined) {
644 params = { ...(nextParams as MessageCreateParamsStreaming), stream: true };
645 }
646 const requestOptions = {
647 ...(options?.signal ? { signal: options.signal } : {}),
648 ...(options?.timeoutMs !== undefined ? { timeout: options.timeoutMs } : {}),
649 maxRetries: 0,
650 };
651 const response = await retryProviderRequest(
652 () => client.beta.messages.create(params, requestOptions).asResponse(),
653 {
654 maxRetries: options?.maxRetries,
655 maxRetryDelayMs: options?.maxRetryDelayMs,
656 signal: options?.signal,
657 },
658 );
659 await options?.onResponse?.({ status: response.status, headers: headersToRecord(response.headers) }, model);
660 stream.push({ type: "start", partial: output });
661
662 type Block = (ThinkingContent | TextContent | (ToolCall & { partialJson: string })) & { index: number };
663 const blocks = output.content as Block[];
664
665 for await (const event of iterateAnthropicEvents(response, options?.signal)) {
666 await options?.onProviderStreamEvent?.(event, model);
667 if (event.type === "message_start") {
668 output.responseId = event.message.id;
669 const transformations = event.message.input_transformations;
670 if (Array.isArray(transformations)) inputTransformations = transformations;
671 const responseModel = event.message.model;
672 if (responseModel !== model.id) output.responseModel = responseModel;
673 const fallbackCost =
674 responseModel === model.id
675 ? undefined
676 : model.compat?.allowedFallbackModels?.find(
677 (fallback) => fallback.provider === model.provider && fallback.model === responseModel,
678 )?.cost;
679 usageModel = fallbackCost ? { ...model, id: responseModel, cost: fallbackCost } : model;
680 // Capture initial token usage from message_start event
681 // This ensures we have input token counts even if the stream is aborted early
682 output.usage.input = event.message.usage.input_tokens || 0;
683 output.usage.output = event.message.usage.output_tokens || 0;
684 output.usage.cacheRead = event.message.usage.cache_read_input_tokens || 0;
685 output.usage.cacheWrite = event.message.usage.cache_creation_input_tokens || 0;
686 output.usage.cacheWrite1h = event.message.usage.cache_creation?.ephemeral_1h_input_tokens || 0;
687 // Anthropic doesn't provide total_tokens, compute from components
688 output.usage.totalTokens =
689 output.usage.input + output.usage.output + output.usage.cacheRead + output.usage.cacheWrite;
690 calculateCost(usageModel, output.usage);
691 } else if (event.type === "content_block_start") {
692 if (event.content_block.type === "fallback") {
693 if (output.content.length > 0) {
694 throw new Error("Anthropic performed an unsupported mid-output model fallback");
695 }
696 continue;
697 }
698 if (event.content_block.type === "text") {
699 const block: Block = {
700 type: "text",
701 text: event.content_block.text ?? "",
702 index: event.index,
703 };
704 output.content.push(block);
705 stream.push({ type: "text_start", contentIndex: output.content.length - 1, partial: output });
706 } else if (event.content_block.type === "thinking") {
707 const block: Block = {
708 type: "thinking",
709 thinking: event.content_block.thinking ?? "",
710 thinkingSignature: event.content_block.signature ?? "",
711 index: event.index,
712 };
713 output.content.push(block);
714 stream.push({ type: "thinking_start", contentIndex: output.content.length - 1, partial: output });
715 } else if (event.content_block.type === "redacted_thinking") {
716 const block: Block = {
717 type: "thinking",
718 thinking: "[Reasoning redacted]",
719 thinkingSignature: event.content_block.data,
720 redacted: true,
721 index: event.index,
722 };
723 output.content.push(block);
724 stream.push({ type: "thinking_start", contentIndex: output.content.length - 1, partial: output });
725 } else if (event.content_block.type === "tool_use") {
726 const block: Block = {
727 type: "toolCall",
728 id: event.content_block.id,
729 name: isOAuth
730 ? fromClaudeCodeName(event.content_block.name, currentTools)
731 : event.content_block.name,
732 arguments: (event.content_block.input as Record<string, any>) ?? {},
733 partialJson: "",
734 index: event.index,
735 };
736 output.content.push(block);
737 stream.push({ type: "toolcall_start", contentIndex: output.content.length - 1, partial: output });
738 }
739 } else if (event.type === "content_block_delta") {
740 if (event.delta.type === "text_delta") {
741 const index = blocks.findIndex((b) => b.index === event.index);
742 const block = blocks[index];
743 if (block && block.type === "text") {
744 block.text += event.delta.text;
745 stream.push({
746 type: "text_delta",
747 contentIndex: index,
748 delta: event.delta.text,
749 partial: output,
750 });
751 }
752 } else if (event.delta.type === "thinking_delta") {
753 const index = blocks.findIndex((b) => b.index === event.index);
754 const block = blocks[index];
755 if (block && block.type === "thinking") {
756 block.thinking += event.delta.thinking;
757 stream.push({
758 type: "thinking_delta",
759 contentIndex: index,
760 delta: event.delta.thinking,
761 partial: output,
762 });
763 }
764 } else if (event.delta.type === "input_json_delta") {
765 const index = blocks.findIndex((b) => b.index === event.index);
766 const block = blocks[index];
767 if (block && block.type === "toolCall") {
768 block.partialJson += event.delta.partial_json;
769 block.arguments = parseStreamingJson(block.partialJson);
770 stream.push({
771 type: "toolcall_delta",
772 contentIndex: index,
773 delta: event.delta.partial_json,
774 partial: output,
775 });
776 }
777 } else if (event.delta.type === "signature_delta") {
778 const index = blocks.findIndex((b) => b.index === event.index);
779 const block = blocks[index];
780 if (block && block.type === "thinking") {
781 block.thinkingSignature = block.thinkingSignature || "";
782 block.thinkingSignature += event.delta.signature;
783 }
784 }
785 } else if (event.type === "content_block_stop") {
786 const index = blocks.findIndex((b) => b.index === event.index);
787 const block = blocks[index];
788 if (block) {
789 delete (block as any).index;
790 if (block.type === "text") {
791 stream.push({
792 type: "text_end",
793 contentIndex: index,
794 content: block.text,
795 partial: output,
796 });
797 } else if (block.type === "thinking") {
798 stream.push({
799 type: "thinking_end",
800 contentIndex: index,
801 content: block.thinking,
802 partial: output,
803 });
804 } else if (block.type === "toolCall") {
805 block.arguments = parseStreamingJson(block.partialJson);
806 // Finalize in-place and strip the scratch buffer so replay only
807 // carries parsed arguments.
808 delete (block as { partialJson?: string }).partialJson;
809 stream.push({
810 type: "toolcall_end",
811 contentIndex: index,
812 toolCall: block,
813 partial: output,
814 });
815 }
816 }
817 } else if (event.type === "message_delta") {
818 const transformations = event.input_transformations;
819 if (Array.isArray(transformations)) inputTransformations = transformations;
820 if (event.delta.stop_reason) {
821 output.rawStopReason = event.delta.stop_reason;
822 const stopReasonResult = mapStopReason(event.delta.stop_reason, event.delta.stop_details);
823 output.stopReason = stopReasonResult.stopReason;
824 if (stopReasonResult.errorMessage) {
825 output.errorMessage = stopReasonResult.errorMessage;
826 }
827 }
828 // Only update usage fields if present (not null).
829 // Preserves input_tokens from message_start when proxies omit it in message_delta.
830 if (event.usage) {
831 if (event.usage.input_tokens != null) {
832 output.usage.input = event.usage.input_tokens;
833 }
834 if (event.usage.output_tokens != null) {
835 output.usage.output = event.usage.output_tokens;
836 }
837 if (event.usage.cache_read_input_tokens != null) {
838 output.usage.cacheRead = event.usage.cache_read_input_tokens;
839 }
840 if (event.usage.cache_creation_input_tokens != null) {
841 output.usage.cacheWrite = event.usage.cache_creation_input_tokens;
842 }
843 // Vercel AI Gateway includes the TTL breakdown in deltas, though the SDK only types it on message_start.
844 const cacheCreation = (
845 event.usage as typeof event.usage & { cache_creation?: { ephemeral_1h_input_tokens?: number } }
846 ).cache_creation;
847 if (cacheCreation?.ephemeral_1h_input_tokens != null) {
848 output.usage.cacheWrite1h = cacheCreation.ephemeral_1h_input_tokens;
849 }
850 // Anthropic reports reasoning tokens as a subset of output tokens.
851 const thinkingTokens = event.usage.output_tokens_details?.thinking_tokens;
852 if (thinkingTokens != null) {
853 output.usage.reasoning = thinkingTokens;
854 }
855 }
856 // Anthropic doesn't provide total_tokens, compute from components
857 output.usage.totalTokens =
858 output.usage.input + output.usage.output + output.usage.cacheRead + output.usage.cacheWrite;
859 calculateCost(usageModel, output.usage);
860 }
861 }
862
863 if (options?.signal?.aborted) {
864 throw new Error("Request was aborted");
865 }
866
867 if (output.stopReason === "pending") {
868 throw new Error("Anthropic stream ended without a stop reason");
869 }
870 if (output.stopReason === "aborted" || output.stopReason === "error") {
871 throw new Error(output.errorMessage || "An unknown error occurred");
872 }
873 if (inputTransformations && inputTransformations.length > 0) {
874 appendAssistantMessageDiagnostic(output, {
875 type: "anthropic_input_transformations",
876 timestamp: Date.now(),
877 details: {
878 transformations: inputTransformations.map((transformation) => ({
879 type: transformation.type ?? undefined,
880 path: transformation.path ?? undefined,
881 reason: transformation.reason ?? undefined,
882 })),
883 },
884 });
885 }
886
887 stream.push({ type: "done", reason: output.stopReason, message: output });
888 stream.end();
889 } catch (error) {
890 for (const block of output.content) {
891 delete (block as { index?: number }).index;
892 // partialJson is only a streaming scratch buffer; never persist it.
893 delete (block as { partialJson?: string }).partialJson;
894 }
895 output.stopReason = options?.signal?.aborted ? "aborted" : "error";
896 output.errorMessage = error instanceof Error ? error.message : JSON.stringify(error);
897 stream.push({ type: "error", reason: output.stopReason, error: output });
898 stream.end();
899 }
900 })();
901
902 return stream;
903};
904
905/**
906 * Map ThinkingLevel to Anthropic effort levels for adaptive thinking.
907 * Note: effort "max" is available on all adaptive-thinking Claude models, while native
908 * "xhigh" is only available on Opus 4.7/4.8, Sonnet 5, and Fable 5.
909 */
910function mapThinkingLevelToEffort(
911 model: Model<"anthropic-messages">,
912 level: SimpleStreamOptions["reasoning"],
913): AnthropicEffort {
914 const mapped = level ? model.thinkingLevelMap?.[level] : undefined;
915 if (typeof mapped === "string") return mapped as AnthropicEffort;
916
917 switch (level) {
918 case "minimal":
919 case "low":
920 return "low";
921 case "medium":
922 return "medium";
923 case "high":
924 return "high";
925 default:
926 return "high";
927 }
928}
929
930export const streamSimple: StreamFunction<"anthropic-messages", SimpleStreamOptions> = (
931 model: Model<"anthropic-messages">,
932 context: TranscriptContext,
933 options?: SimpleStreamOptions,
934): AssistantMessageEventStream => {
935 if (!getAnthropicFederation(model, options?.apiKey, options?.headers, options?.env)) {
936 assertRequestAuth(model.provider, options?.apiKey, options?.headers);
937 }
938
939 const base = {
940 ...buildBaseOptions(model, context, options, options?.apiKey),
941 toolChoice: options?.toolChoice,
942 } satisfies AnthropicOptions;
943 if (!options?.reasoning) {
944 return stream(model, context, {
945 ...base,
946 thinkingEnabled: false,
947 } satisfies AnthropicOptions);
948 }
949
950 // For models with adaptive thinking: use an effort level.
951 // For older models: use budget-based thinking.
952 if (model.compat?.forceAdaptiveThinking === true) {
953 const effort = mapThinkingLevelToEffort(model, options.reasoning);
954 return stream(model, context, {
955 ...base,
956 thinkingEnabled: true,
957 effort,
958 } satisfies AnthropicOptions);
959 }
960
961 // Undefined means the caller did not request an output cap; let the helper use the model cap.
962 // Do not coerce to 0 here, or the thinking budget would become the entire max_tokens value.
963 const adjusted = adjustMaxTokensForThinking(
964 base.maxTokens,
965 model.maxTokens,
966 options.reasoning,
967 options.thinkingBudgets,
968 );
969
970 const maxTokens = clampMaxTokensToContext(model, context, adjusted.maxTokens);
971
972 return stream(model, context, {
973 ...base,
974 maxTokens,
975 thinkingEnabled: true,
976 thinkingBudgetTokens: Math.min(adjusted.thinkingBudget, Math.max(0, maxTokens - 1024)),
977 } satisfies AnthropicOptions);
978};
979
980function isOAuthToken(apiKey: string): boolean {
981 return apiKey.includes("sk-ant-oat");
982}
983
984function createClient(
985 model: Model<"anthropic-messages">,
986 apiKey: string | undefined,
987 optionsHeaders?: ProviderHeaders,
988 fetch?: typeof globalThis.fetch,
989 dynamicHeaders?: Record<string, string>,
990 sessionId?: string,
991 federation?: AnthropicFederationConfig,
992): { client: Anthropic; isOAuthToken: boolean } {
993 // Copilot: Bearer auth.
994 if (model.provider === "github-copilot") {
995 const client = new PiAnthropic({
996 apiKey: null,
997 authToken: apiKey ?? null,
998 baseURL: model.baseUrl,
999 dangerouslyAllowBrowser: true,
1000 fetch,
1001 defaultHeaders: mergeClientHeaders(
1002 {
1003 accept: "application/json",
1004 "anthropic-dangerous-direct-browser-access": "true",
1005 },
1006 model.headers,
1007 dynamicHeaders,
1008 optionsHeaders,
1009 ),
1010 });
1011
1012 return { client, isOAuthToken: false };
1013 }
1014
1015 // OAuth: Bearer auth, Claude Code identity headers
1016 if (apiKey && isOAuthToken(apiKey)) {
1017 const client = new PiAnthropic({
1018 apiKey: null,
1019 authToken: apiKey,
1020 baseURL: model.baseUrl,
1021 dangerouslyAllowBrowser: true,
1022 fetch,
1023 defaultHeaders: mergeClientHeaders(
1024 {
1025 accept: "application/json",
1026 "anthropic-dangerous-direct-browser-access": "true",
1027 "user-agent": `claude-cli/${claudeCodeVersion}`,
1028 "x-app": "cli",
1029 },
1030 model.headers,
1031 optionsHeaders,
1032 ),
1033 });
1034
1035 return { client, isOAuthToken: true };
1036 }
1037
1038 // API key, header-owned auth, or workload identity federation.
1039 const compat = getAnthropicCompat(model);
1040 const sessionAffinityHeaders: ProviderHeaders = {};
1041 if (sessionId && compat.sendSessionAffinityHeaders) {
1042 const header = compat.sessionAffinityFormat === "openrouter" ? "x-session-id" : "x-session-affinity";
1043 sessionAffinityHeaders[header] = sessionId;
1044 }
1045 const defaultHeaders = mergeClientHeaders(
1046 {
1047 accept: "application/json",
1048 "anthropic-dangerous-direct-browser-access": "true",
1049 },
1050 sessionAffinityHeaders,
1051 model.headers,
1052 optionsHeaders,
1053 );
1054 if (federation) {
1055 const key = JSON.stringify([model.baseUrl, federation]);
1056 if (federationClient?.key !== key || federationClient.fetch !== fetch) {
1057 const client = new PiAnthropic({
1058 apiKey: null,
1059 authToken: null,
1060 config: federation,
1061 baseURL: model.baseUrl,
1062 dangerouslyAllowBrowser: true,
1063 fetch,
1064 });
1065 federationClient = { key, fetch, client };
1066 }
1067 return { client: federationClient.client.withOptions({ defaultHeaders }), isOAuthToken: false };
1068 }
1069
1070 const client = new PiAnthropic({
1071 apiKey: apiKey ?? null,
1072 authToken: null,
1073 baseURL: model.baseUrl,
1074 dangerouslyAllowBrowser: true,
1075 fetch,
1076 defaultHeaders,
1077 });
1078
1079 return { client, isOAuthToken: false };
1080}
1081
1082function getBetaFeatures(
1083 model: Model<"anthropic-messages">,
1084 context: TranscriptContext,
1085 isOAuthToken: boolean,
1086 nativeToolChanges: boolean,
1087 options?: AnthropicOptions,
1088): NonNullable<MessageCreateParamsStreaming["betas"]> {
1089 let configuredFeatures: string | null | undefined;
1090 for (const headers of [model.headers, options?.headers]) {
1091 for (const [name, value] of Object.entries(headers ?? {})) {
1092 if (name.toLowerCase() === "anthropic-beta") configuredFeatures = value;
1093 }
1094 }
1095 if (configuredFeatures === null) return [];
1096 if (configuredFeatures !== undefined) {
1097 return [
1098 ...new Set(
1099 configuredFeatures
1100 .split(",")
1101 .map((feature) => feature.trim())
1102 .filter((feature) => feature.length > 0),
1103 ),
1104 ];
1105 }
1106
1107 const features: NonNullable<MessageCreateParamsStreaming["betas"]> = [];
1108 if (isOAuthToken) features.push("claude-code-20250219", "oauth-2025-04-20");
1109 if (shouldUseFineGrainedToolStreamingBeta(model, context)) features.push(FINE_GRAINED_TOOL_STREAMING_BETA);
1110 if (
1111 model.reasoning &&
1112 options?.thinkingEnabled === true &&
1113 (options.interleavedThinking ?? true) &&
1114 model.compat?.forceAdaptiveThinking !== true
1115 ) {
1116 features.push(INTERLEAVED_THINKING_BETA);
1117 }
1118 if (shouldUseServerSideFallbackBeta(model)) features.push(SERVER_SIDE_FALLBACK_BETA);
1119 if (model.compat?.supportsMidConvoEffort === true) {
1120 features.push(MID_CONVERSATION_OUTPUT_CONFIG_BETA, THINKING_BINDING_CONTROLS_BETA);
1121 }
1122 if (nativeToolChanges) features.push(MID_CONVERSATION_TOOL_CHANGES_BETA);
1123 return [...new Set(features)];
1124}
1125
1126function buildParams(
1127 model: Model<"anthropic-messages">,
1128 context: TranscriptContext,
1129 isOAuthToken: boolean,
1130 options?: AnthropicOptions,
1131): MessageCreateParamsStreaming {
1132 const { cacheControl } = getCacheControl(model, options?.cacheRetention, options?.env);
1133 const compat = getAnthropicCompat(model);
1134 const initialSystemMessage = getInitialSystemMessage(context.messages);
1135 const initialSystemText = initialSystemMessage ? getSystemMessageText(initialSystemMessage) : "";
1136 const transformedMessages = transformMessages(context.messages, model, normalizeToolCallId);
1137 const conversationMessages = initialSystemMessage ? transformedMessages.slice(1) : transformedMessages;
1138 // Native tool changes reference tools by name, so a redefined name cannot be expressed,
1139 // and Anthropic rejects a tool list where every tool is deferred, so there must be an
1140 // initial active tool to anchor the deferred ones. Otherwise the current tool list is sent.
1141 const initialTools = initialSystemMessage?.toolsAdded ?? [];
1142 const nativeToolChanges =
1143 compat.supportsMidConvoSystemMessages &&
1144 compat.supportsMidConvoToolChanges &&
1145 initialTools.length > 0 &&
1146 !hasToolRedefinitions(context.messages);
1147 const converted = convertMessages(
1148 conversationMessages,
1149 isOAuthToken,
1150 cacheControl,
1151 compat.allowEmptySignature,
1152 model.compat?.supportsMidConvoEffort === true ? model.provider : undefined,
1153 nativeToolChanges,
1154 );
1155 const activeEffort = options?.effort ?? "high";
1156 const betaFeatures = getBetaFeatures(model, context, isOAuthToken, nativeToolChanges, options);
1157 const params: MessageCreateParamsStreaming = {
1158 model: model.id,
1159 messages:
1160 model.compat?.supportsMidConvoEffort === true
1161 ? insertThinkingLevelMessages(converted, activeEffort)
1162 : converted.messages,
1163 max_tokens: options?.maxTokens ?? model.maxTokens,
1164 stream: true,
1165 ...(betaFeatures.length > 0 ? { betas: betaFeatures } : {}),
1166 };
1167
1168 // For OAuth tokens, we MUST include Claude Code identity
1169 if (isOAuthToken) {
1170 params.system = [
1171 {
1172 type: "text",
1173 text: "You are Claude Code, Anthropic's official CLI for Claude.",
1174 ...(cacheControl ? { cache_control: cacheControl } : {}),
1175 },
1176 ];
1177 if (initialSystemText) {
1178 params.system.push({
1179 type: "text",
1180 text: sanitizeSurrogates(initialSystemText),
1181 ...(cacheControl ? { cache_control: cacheControl } : {}),
1182 });
1183 }
1184 } else if (initialSystemText) {
1185 // Add cache control to system prompt for non-OAuth tokens
1186 params.system = [
1187 {
1188 type: "text",
1189 text: sanitizeSurrogates(initialSystemText),
1190 ...(cacheControl ? { cache_control: cacheControl } : {}),
1191 },
1192 ];
1193 }
1194
1195 // Temperature is incompatible with extended thinking and unsupported on Claude Opus 4.7+.
1196 if (
1197 options?.temperature !== undefined &&
1198 !options?.thinkingEnabled &&
1199 model.compat?.supportsMidConvoEffort !== true &&
1200 compat.supportsTemperature
1201 ) {
1202 params.temperature = options.temperature;
1203 }
1204
1205 const toolCacheControl = compat.supportsCacheControlOnTools ? cacheControl : undefined;
1206 if (nativeToolChanges) {
1207 // Initial tools stay active with the cache breakpoint on the last one. Every later
1208 // declaration is deferred and only surfaced by its `tool_addition` block; removed
1209 // tools stay declared and are withdrawn by `tool_removal`. The request-level list
1210 // therefore only grows, keeping the cached prefix intact across tool changes.
1211 const initialNames = new Set(initialTools.map((tool) => tool.name));
1212 const laterTools = getDeclaredTools(context.messages).filter((tool) => !initialNames.has(tool.name));
1213 params.tools = [
1214 ...convertTools(
1215 initialTools,
1216 isOAuthToken,
1217 compat.supportsEagerToolInputStreaming,
1218 compat.supportsStrictTools,
1219 toolCacheControl,
1220 ),
1221 DEFERRED_TOOL_PLACEHOLDER,
1222 ...convertTools(
1223 laterTools,
1224 isOAuthToken,
1225 compat.supportsEagerToolInputStreaming,
1226 compat.supportsStrictTools,
1227 ).map((tool) => ({ ...tool, defer_loading: true })),
1228 ];
1229 } else {
1230 const tools = getCurrentTools(context.messages);
1231 if (tools.length > 0) {
1232 params.tools = convertTools(
1233 tools,
1234 isOAuthToken,
1235 compat.supportsEagerToolInputStreaming,
1236 compat.supportsStrictTools,
1237 toolCacheControl,
1238 );
1239 }
1240 }
1241
1242 // Managed effort models always use adaptive thinking so prefix mismatches can
1243 // be dropped instead of surfacing as persistent 400 responses.
1244 if (model.compat?.supportsMidConvoEffort === true) {
1245 params.thinking = {
1246 type: "adaptive",
1247 display: options?.thinkingDisplay ?? "summarized",
1248 block_binding: { prefix_mismatch_behavior: "drop_block" },
1249 };
1250 params.output_config = { effort: "high" };
1251 } else if (model.reasoning) {
1252 if (options?.thinkingEnabled) {
1253 // Default to "summarized" so Opus 4.7 and Mythos Preview behave like
1254 // older Claude 4 models (whose API default is also "summarized").
1255 const display: AnthropicThinkingDisplay = options.thinkingDisplay ?? "summarized";
1256 if (model.compat?.forceAdaptiveThinking === true) {
1257 // Adaptive thinking: Claude decides when and how much to think.
1258 params.thinking = { type: "adaptive", display };
1259 if (options.effort) {
1260 params.output_config = { effort: options.effort };
1261 }
1262 } else {
1263 // Budget-based thinking for older models
1264 params.thinking = {
1265 type: "enabled",
1266 budget_tokens: options.thinkingBudgetTokens || 1024,
1267 display,
1268 };
1269 }
1270 } else if (options?.thinkingEnabled === false && model.thinkingLevelMap?.off !== null) {
1271 params.thinking = { type: "disabled" };
1272 }
1273 }
1274
1275 if (options?.metadata) {
1276 const userId = options.metadata.user_id;
1277 if (typeof userId === "string") {
1278 params.metadata = { user_id: userId };
1279 }
1280 }
1281
1282 if (options?.toolChoice) {
1283 if (typeof options.toolChoice === "string") {
1284 params.tool_choice = { type: options.toolChoice };
1285 } else {
1286 params.tool_choice = options.toolChoice;
1287 }
1288 }
1289
1290 const allowedFallbackModels = model.compat?.allowedFallbackModels;
1291 if (allowedFallbackModels && allowedFallbackModels.length > 0) {
1292 params.fallbacks = allowedFallbackModels.map((fallback) => ({ model: fallback.model }));
1293 }
1294
1295 return params;
1296}
1297
1298// Normalize tool call IDs to match Anthropic's required pattern and length
1299function normalizeToolCallId(id: string): string {
1300 return id.replace(/[^a-zA-Z0-9_-]/g, "_").slice(0, 64);
1301}
1302
1303function convertToolResult(msg: ToolResultMessage): ContentBlockParam {
1304 return {
1305 type: "tool_result",
1306 tool_use_id: msg.toolCallId,
1307 content: convertContentBlocks(msg.content),
1308 is_error: msg.isError,
1309 };
1310}
1311
1312interface ConvertedAnthropicMessages {
1313 messages: MessageParam[];
1314 assistantLevels: Map<number, AnthropicEffort>;
1315}
1316
1317function convertMessages(
1318 transformedMessages: Message[],
1319 isOAuthToken: boolean,
1320 cacheControl?: CacheControlEphemeral,
1321 allowEmptySignature = false,
1322 managedProvider?: string,
1323 nativeToolChanges = false,
1324): ConvertedAnthropicMessages {
1325 const params: MessageParam[] = [];
1326 const assistantLevels = new Map<number, AnthropicEffort>();
1327 // Later system messages are held back and emitted directly before the next assistant
1328 // message (or at the end of the transcript). Anthropic requires `tool_result` blocks to
1329 // immediately follow their `tool_use`, so a system message between them is rejected; this
1330 // also mirrors where the managed-effort system messages are inserted. As a result an
1331 // update placed before a user message in the transcript lands after it on the wire.
1332 const pendingSystemMessages: MessageParam[] = [];
1333 const flushPendingSystemMessages = (): void => {
1334 params.push(...pendingSystemMessages);
1335 pendingSystemMessages.length = 0;
1336 };
1337
1338 for (let i = 0; i < transformedMessages.length; i++) {
1339 const msg = transformedMessages[i];
1340
1341 if (msg.role === "system") {
1342 // Later system messages only reach this point when the model accepts them natively;
1343 // otherwise the transcript was collapsed into the leading message before conversion.
1344 const text = renderSystemMessageUpdate(msg);
1345 const blocks: ContentBlockParam[] = [];
1346 if (text.length > 0) blocks.push({ type: "text", text: sanitizeSurrogates(text) });
1347 if (nativeToolChanges) {
1348 for (const tool of msg.toolsRemoved ?? []) {
1349 blocks.push({
1350 type: "tool_removal",
1351 tool: { type: "tool_reference", name: isOAuthToken ? toClaudeCodeName(tool.name) : tool.name },
1352 });
1353 }
1354 for (const tool of msg.toolsAdded ?? []) {
1355 blocks.push({
1356 type: "tool_addition",
1357 tool: { type: "tool_reference", name: isOAuthToken ? toClaudeCodeName(tool.name) : tool.name },
1358 });
1359 }
1360 }
1361 if (blocks.length > 0) pendingSystemMessages.push({ role: "system", content: blocks });
1362 } else if (msg.role === "user") {
1363 if (typeof msg.content === "string") {
1364 if (msg.content.trim().length > 0) {
1365 params.push({
1366 role: "user",
1367 content: sanitizeSurrogates(msg.content),
1368 });
1369 }
1370 } else {
1371 const blocks: ContentBlockParam[] = msg.content.map((item) => {
1372 if (item.type === "text") {
1373 return {
1374 type: "text",
1375 text: sanitizeSurrogates(item.text),
1376 };
1377 } else {
1378 return {
1379 type: "image",
1380 source: {
1381 type: "base64",
1382 media_type: item.mimeType as "image/jpeg" | "image/png" | "image/gif" | "image/webp",
1383 data: item.data,
1384 },
1385 };
1386 }
1387 });
1388 const filteredBlocks = blocks.filter((b) => {
1389 if (b.type === "text") {
1390 return b.text.trim().length > 0;
1391 }
1392 return true;
1393 });
1394 if (filteredBlocks.length === 0) continue;
1395 params.push({
1396 role: "user",
1397 content: filteredBlocks,
1398 });
1399 }
1400 } else if (msg.role === "assistant") {
1401 flushPendingSystemMessages();
1402 const blocks: ContentBlockParam[] = [];
1403
1404 for (const block of msg.content) {
1405 if (block.type === "text") {
1406 if (block.text.trim().length === 0) continue;
1407 blocks.push({
1408 type: "text",
1409 text: sanitizeSurrogates(block.text),
1410 });
1411 } else if (block.type === "thinking") {
1412 // Redacted thinking: pass the opaque payload back as redacted_thinking
1413 if (block.redacted) {
1414 blocks.push({
1415 type: "redacted_thinking",
1416 data: block.thinkingSignature!,
1417 });
1418 continue;
1419 }
1420 const thinkingSignature = block.thinkingSignature;
1421 const hasThinkingSignature = !!thinkingSignature && thinkingSignature.trim().length > 0;
1422 if (block.thinking.trim().length === 0 && !hasThinkingSignature) continue;
1423 // If thinking signature is missing/empty (e.g., from aborted stream),
1424 // convert to plain text for Anthropic. Some compatible providers emit
1425 // and accept empty signatures, so let marked models preserve the block.
1426 if (!hasThinkingSignature) {
1427 blocks.push(
1428 allowEmptySignature
1429 ? {
1430 type: "thinking",
1431 thinking: sanitizeSurrogates(block.thinking),
1432 signature: "",
1433 }
1434 : {
1435 type: "text",
1436 text: sanitizeSurrogates(block.thinking),
1437 },
1438 );
1439 } else {
1440 blocks.push({
1441 type: "thinking",
1442 thinking: sanitizeSurrogates(block.thinking),
1443 signature: thinkingSignature,
1444 });
1445 }
1446 } else if (block.type === "toolCall") {
1447 blocks.push({
1448 type: "tool_use",
1449 id: block.id,
1450 name: isOAuthToken ? toClaudeCodeName(block.name) : block.name,
1451 input: block.arguments ?? {},
1452 });
1453 }
1454 }
1455 if (blocks.length === 0) continue;
1456 const messageIndex = params.length;
1457 params.push({
1458 role: "assistant",
1459 content: blocks,
1460 });
1461 if (
1462 managedProvider !== undefined &&
1463 msg.api === "anthropic-messages" &&
1464 msg.provider === managedProvider &&
1465 isAnthropicEffort(msg.providerThinkingLevel)
1466 ) {
1467 assistantLevels.set(messageIndex, msg.providerThinkingLevel);
1468 }
1469 } else if (msg.role === "toolResult") {
1470 // Collect all consecutive toolResult messages, needed for z.ai Anthropic endpoint.
1471 const toolResults: ContentBlockParam[] = [];
1472 let j = i;
1473 while (j < transformedMessages.length && transformedMessages[j].role === "toolResult") {
1474 toolResults.push(convertToolResult(transformedMessages[j] as ToolResultMessage));
1475 j++;
1476 }
1477
1478 // Skip the messages we've already processed.
1479 i = j - 1;
1480
1481 params.push({
1482 role: "user",
1483 content: toolResults,
1484 });
1485 }
1486 }
1487
1488 flushPendingSystemMessages();
1489
1490 // Add cache_control to the last user or system message to cache conversation history
1491 if (cacheControl && params.length > 0) {
1492 const lastMessage = params[params.length - 1];
1493 if (lastMessage.role === "user" || lastMessage.role === "system") {
1494 if (Array.isArray(lastMessage.content)) {
1495 const lastBlock = lastMessage.content[lastMessage.content.length - 1];
1496 if (
1497 lastBlock &&
1498 (lastBlock.type === "text" ||
1499 lastBlock.type === "image" ||
1500 lastBlock.type === "tool_result" ||
1501 lastBlock.type === "tool_addition" ||
1502 lastBlock.type === "tool_removal")
1503 ) {
1504 (lastBlock as any).cache_control = cacheControl;
1505 }
1506 } else if (typeof lastMessage.content === "string") {
1507 lastMessage.content = [
1508 {
1509 type: "text",
1510 text: lastMessage.content,
1511 cache_control: cacheControl,
1512 },
1513 ] as any;
1514 }
1515 }
1516 }
1517
1518 return { messages: params, assistantLevels };
1519}
1520
1521function isAnthropicEffort(value: unknown): value is AnthropicEffort {
1522 return value === "low" || value === "medium" || value === "high" || value === "xhigh" || value === "max";
1523}
1524
1525function insertThinkingLevelMessages(
1526 converted: ConvertedAnthropicMessages,
1527 activeEffort: AnthropicEffort,
1528): MessageParam[] {
1529 const messages: MessageParam[] = [];
1530 for (let index = 0; index < converted.messages.length; index++) {
1531 const historicalEffort = converted.assistantLevels.get(index);
1532 if (historicalEffort !== undefined) {
1533 messages.push({ role: "system", content: [], output_config: { effort: historicalEffort } });
1534 }
1535 messages.push(converted.messages[index]);
1536 }
1537 messages.push({ role: "system", content: [], output_config: { effort: activeEffort } });
1538 return messages;
1539}
1540
1541function shouldUseFineGrainedToolStreamingBeta(
1542 model: Model<"anthropic-messages">,
1543 context: TranscriptContext,
1544): boolean {
1545 return getCurrentTools(context.messages).length > 0 && !getAnthropicCompat(model).supportsEagerToolInputStreaming;
1546}
1547
1548// Keywords Anthropic strict tool use rejects with a 400 for the whole request.
1549// https://platform.claude.com/docs/en/build-with-claude/structured-outputs#json-schema-limitations
1550const ANTHROPIC_STRICT_UNSUPPORTED_KEYWORDS = new Set([
1551 "minimum",
1552 "maximum",
1553 "exclusiveMinimum",
1554 "exclusiveMaximum",
1555 "multipleOf",
1556 "maxItems",
1557 "uniqueItems",
1558 "minContains",
1559 "maxContains",
1560 "minProperties",
1561 "maxProperties",
1562]);
1563const ANTHROPIC_STRICT_STRING_FORMATS = new Set([
1564 "date-time",
1565 "time",
1566 "date",
1567 "duration",
1568 "email",
1569 "hostname",
1570 "uri",
1571 "ipv4",
1572 "ipv6",
1573 "uuid",
1574]);
1575
1576const isAnthropicStrictUnsupportedKeyword: UnsupportedStrictSchemaKeywordCheck = (key, value) => {
1577 if (ANTHROPIC_STRICT_UNSUPPORTED_KEYWORDS.has(key)) return true;
1578 if (key === "minItems") return value !== 0 && value !== 1;
1579 if (key === "format") return typeof value !== "string" || !ANTHROPIC_STRICT_STRING_FORMATS.has(value);
1580 return false;
1581};
1582
1583function convertTools(
1584 tools: Tool[],
1585 isOAuthToken: boolean,
1586 supportsEagerToolInputStreaming: boolean,
1587 supportsStrictTools: boolean,
1588 cacheControl?: CacheControlEphemeral,
1589): BetaTool[] {
1590 if (!tools) return [];
1591
1592 return tools.map((tool, index) => {
1593 const strict = resolveJsonSchemaStrictSampling(tool, supportsStrictTools, isAnthropicStrictUnsupportedKeyword);
1594 const parameters = getJsonSchemaToolParameters(tool, strict);
1595 const schema = parameters as { properties?: unknown; required?: string[] };
1596 const legacyInputSchema = {
1597 type: "object" as const,
1598 properties: schema.properties ?? {},
1599 required: schema.required ?? [],
1600 };
1601 const inputSchema =
1602 strict === true
1603 ? {
1604 ...(parameters as Record<string, unknown>),
1605 ...legacyInputSchema,
1606 }
1607 : legacyInputSchema;
1608
1609 return {
1610 name: isOAuthToken ? toClaudeCodeName(tool.name) : tool.name,
1611 description: tool.description,
1612 ...(supportsEagerToolInputStreaming ? { eager_input_streaming: true } : {}),
1613 ...(strict === true ? { strict: true } : {}),
1614 input_schema: inputSchema,
1615 ...(cacheControl && index === tools.length - 1 ? { cache_control: cacheControl } : {}),
1616 };
1617 });
1618}
1619
1620function mapStopReason(
1621 reason: BetaStopReason | string,
1622 stopDetails?: RefusalStopDetails | null,
1623): { stopReason: StopReason; errorMessage?: string } {
1624 switch (reason) {
1625 case "end_turn":
1626 return { stopReason: "stop" };
1627 case "max_tokens":
1628 return { stopReason: "length" };
1629 case "tool_use":
1630 return { stopReason: "toolUse" };
1631 case "refusal":
1632 return {
1633 stopReason: "error",
1634 errorMessage: stopDetails?.explanation || `The model refused to complete the request`,
1635 };
1636 case "pause_turn": // Stop is good enough -> resubmit
1637 return { stopReason: "stop" };
1638 case "stop_sequence":
1639 return { stopReason: "stop" }; // We don't supply stop sequences, so this should never happen
1640 case "sensitive": // Content flagged by safety filters (not yet in SDK types)
1641 return { stopReason: "error", errorMessage: "Provider stopped with: sensitive" };
1642 default:
1643 // Handle unknown stop reasons gracefully (API may add new values)
1644 throw new Error(`Unhandled stop reason: ${reason}`);
1645 }
1646}