1
import Anthropic, { type ClientOptions } from "@anthropic-ai/sdk";2
import 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";13
import {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";20
import { calculateCost } from "../models.ts";21
import 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";40
import { appendAssistantMessageDiagnostic } from "../utils/diagnostics.ts";41
import { AssistantMessageEventStream } from "../utils/event-stream.ts";42
import { headersToRecord } from "../utils/headers.ts";43
import { parseJsonWithRepair, parseStreamingJson } from "../utils/json-parse.ts";44
import { getPiUserAgent } from "../utils/pi-user-agent.ts";45
import { getProviderEnvValue } from "../utils/provider-env.ts";46
import { retryProviderRequest } from "../utils/provider-retry.ts";47
import { sanitizeSurrogates } from "../utils/sanitize-unicode.ts";48
import { getSystemMessageText, renderSystemMessageUpdate } from "../utils/text.ts";49
import {50
getCurrentTools,51
getDeclaredTools,52
getInitialSystemMessage,53
hasToolRedefinitions,54
resolveTranscript,55
type TranscriptContext,56
} from "../utils/transcript.ts";58
import {59
getJsonSchemaToolParameters,60
resolveJsonSchemaStrictSampling,61
type UnsupportedStrictSchemaKeywordCheck,62
} from "./constrained-sampling.ts";63
import { buildCopilotDynamicHeaders, hasCopilotVisionInput } from "./github-copilot-headers.ts";64
import { adjustMaxTokensForThinking, buildBaseOptions, clampMaxTokensToContext } from "./simple-options.ts";65
import { transformMessages } from "./transform-messages.ts";67
/**68
* Resolve cache retention preference.69
* Defaults to "short" and uses PI_CACHE_RETENTION for backward compatibility.70
*/71
function 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
}81
function 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
}97
// Stealth mode: Mimic Claude Code's tool naming exactly98
const claudeCodeVersion = "2.1.280";100
// Claude Code 2.x tool names (canonical casing)101
// Source: https://cchistory.mariozechner.at/data/prompts-2.1.11.md102
// To update: https://github.com/badlogic/cchistory103
const 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
];123
const ccToolLookup = new Map(claudeCodeTools.map((t) => [t.toLowerCase(), t]));125
// Convert tool name to CC canonical casing if it matches (case-insensitive)126
const toClaudeCodeName = (name: string) => ccToolLookup.get(name.toLowerCase()) ?? name;127
const 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
};136
/**137
* Convert content blocks to Anthropic API format138
*/139
function convertContentBlocks(content: (TextContent | ImageContent)[]):140
| string141
| 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 simplicity153
const hasImages = content.some((c) => c.type === "image");154
if (!hasImages) {155
return sanitizeSurrogates(content.map((c) => (c as TextContent).text).join("\n"));156
}158
// If we have images, convert to content block array159
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
});176
// If only images (no text), add placeholder text block177
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
}185
return blocks;186
}188
export type AnthropicEffort = "low" | "medium" | "high" | "xhigh" | "max";190
export type AnthropicThinkingDisplay = "summarized" | "omitted";192
const FINE_GRAINED_TOOL_STREAMING_BETA = "fine-grained-tool-streaming-2025-05-14";193
const INTERLEAVED_THINKING_BETA = "interleaved-thinking-2025-05-14";194
const SERVER_SIDE_FALLBACK_BETA = "server-side-fallback-2026-07-01";195
const MID_CONVERSATION_OUTPUT_CONFIG_BETA = "mid-conversation-output-config-2026-07-01";196
const THINKING_BINDING_CONTROLS_BETA = "thinking-binding-controls-2026-08-01";197
const MID_CONVERSATION_TOOL_CHANGES_BETA = "mid-conversation-tool-changes-2026-07-01";199
/**200
* Stable deferred tool declared whenever native tool changes are in use. Anthropic adds201
* hidden prompt scaffolding as soon as any tool has `defer_loading`; declaring this202
* placeholder from the first request keeps that scaffolding in the cached prefix, so the203
* 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
*/206
const 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
};213
function shouldUseServerSideFallbackBeta(model: Model<"anthropic-messages">): boolean {214
return (model.compat?.allowedFallbackModels?.length ?? 0) > 0;215
}217
function 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
}233
export 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()` maps239
* 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 reasoning254
* - "medium": Moderate thinking, may skip for simple queries255
* - "low": Minimal thinking, skips for simple tasks256
* Ignored for older models.257
* Default: omitted unless `streamSimple()` maps a simple reasoning258
* 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 encrypted265
* signature still travels back for multi-turn continuity. Use for faster266
* 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 Preview269
* is "omitted". We default to "summarized" here to keep behavior consistent270
* 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-adaptive276
* 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-in283
* 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 client289
* construction entirely. Use this to inject alternative SDK clients such as290
* `AnthropicVertex` that shares the same messaging API.291
*/292
client?: Anthropic;293
}295
function 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
}305
function mergeClientHeaders(...headerSources: (ProviderHeaders | undefined)[]): ProviderHeaders {306
return mergeHeaders({ "User-Agent": getPiUserAgent() }, ...headerSources);307
}309
function 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
}318
function 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
}327
function 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
}331
/**332
* Anthropic SDK client that never runs the SDK's own credential chain333
* (ANTHROPIC_PROFILE config files, federation env vars). Without this, every334
* client built with `apiKey: null, authToken: null` for header-owned auth would335
* also resolve and exchange SDK credentials behind pi's auth resolver.336
*/337
class PiAnthropic extends Anthropic {338
protected override _shouldResolveDefaultCredentials(): boolean {339
return false;340
}341
}343
type AnthropicFederationConfig = NonNullable<ClientOptions["config"]>;345
/**346
* Workload identity federation config from the ANTHROPIC_* variables the347
* Anthropic SDK documents; the SDK performs the token exchange and refresh.348
* Only for the anthropic provider, since the exchange is an Anthropic API349
* endpoint, and only when no key or auth header was resolved.350
*/351
function 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
}374
/**375
* The SDK caches the federated access token per client, but pi creates a client376
* per request. Keep one client for the current federation config and fetch, and377
* clone it per request with `withOptions()`, which shares the token cache.378
*/379
let federationClient: { key: string; fetch: typeof globalThis.fetch | undefined; client: Anthropic } | undefined;381
interface ServerSentEvent {382
event: string | null;383
data: string;384
raw: string[];385
}387
interface SseDecoderState {388
event: string | null;389
data: string[];390
raw: string[];391
}393
const 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
]);402
function flushSseEvent(state: SseDecoderState): ServerSentEvent | null {403
if (!state.event && state.data.length === 0) {404
return null;405
}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
}418
function decodeSseLine(line: string, state: SseDecoderState): ServerSentEvent | null {419
if (line === "") {420
return flushSseEvent(state);421
}423
state.raw.push(line);424
if (line.startsWith(":")) {425
return null;426
}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
}435
if (fieldName === "event") {436
state.event = value;437
} else if (fieldName === "data") {438
state.data.push(value);439
}441
return null;442
}444
function 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
}456
function consumeLine(text: string): { line: string; rest: string } | null {457
const lineBreakIndex = nextLineBreakIndex(text);458
if (lineBreakIndex === -1) {459
return null;460
}462
let nextIndex = lineBreakIndex + 1;463
if (text[lineBreakIndex] === "\r" && text[nextIndex] === "\n") {464
nextIndex += 1;465
}467
return {468
line: text.slice(0, lineBreakIndex),469
rest: text.slice(nextIndex),470
};471
}473
async 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 = "";482
try {483
while (true) {484
if (signal?.aborted) {485
throw new Error("Request was aborted");486
}488
const { value, done } = await reader.read();489
if (done) {490
break;491
}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
}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
}516
if (buffer.length > 0) {517
const event = decodeSseLine(buffer, state);518
if (event) {519
yield event;520
}521
}523
const trailingEvent = flushSseEvent(state);524
if (trailingEvent) {525
yield trailingEvent;526
}527
} finally {528
reader.releaseLock();529
}530
}532
async 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
}540
let sawMessageStart = false;541
let sawMessageEnd = false;543
for await (const sse of iterateSseMessages(response.body, signal)) {544
if (sse.event === "error") {545
throw new Error(sse.data);546
}548
if (!ANTHROPIC_MESSAGE_EVENTS.has(sse.event ?? "")) {549
continue;550
}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
}568
if (sawMessageStart && !sawMessageEnd) {569
throw new Error("Anthropic stream ended before message_stop");570
}571
}573
export 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);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
};603
try {604
let client: Anthropic;605
let isOAuth: boolean;606
let usageModel = model;607
let inputTransformations: BetaThinkingDroppedInputTransformation[] | undefined;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);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
}626
const cacheRetention = resolveCacheRetention(options?.cacheRetention, options?.env);627
const cacheSessionId = cacheRetention === "none" ? undefined : options?.sessionId;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 });662
type Block = (ThinkingContent | TextContent | (ToolCall & { partialJson: string })) & { index: number };663
const blocks = output.content as Block[];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.id675
? undefined676
: 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 event681
// This ensures we have input token counts even if the stream is aborted early682
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 components688
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: isOAuth730
? 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 only807
// 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 components857
output.usage.totalTokens =858
output.usage.input + output.usage.output + output.usage.cacheRead + output.usage.cacheWrite;859
calculateCost(usageModel, output.usage);860
}861
}863
if (options?.signal?.aborted) {864
throw new Error("Request was aborted");865
}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
}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
})();902
return stream;903
};905
/**906
* Map ThinkingLevel to Anthropic effort levels for adaptive thinking.907
* Note: effort "max" is available on all adaptive-thinking Claude models, while native908
* "xhigh" is only available on Opus 4.7/4.8, Sonnet 5, and Fable 5.909
*/910
function 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;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
}930
export 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
}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
}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
}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
);970
const maxTokens = clampMaxTokensToContext(model, context, adjusted.maxTokens);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
};980
function isOAuthToken(apiKey: string): boolean {981
return apiKey.includes("sk-ant-oat");982
}984
function 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
});1012
return { client, isOAuthToken: false };1013
}1015
// OAuth: Bearer auth, Claude Code identity headers1016
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
});1035
return { client, isOAuthToken: true };1036
}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
}1070
const client = new PiAnthropic({1071
apiKey: apiKey ?? null,1072
authToken: null,1073
baseURL: model.baseUrl,1074
dangerouslyAllowBrowser: true,1075
fetch,1076
defaultHeaders,1077
});1079
return { client, isOAuthToken: false };1080
}1082
function 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
configuredFeatures1100
.split(",")1101
.map((feature) => feature.trim())1102
.filter((feature) => feature.length > 0),1103
),1104
];1105
}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 !== true1115
) {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
}1126
function 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 an1140
// 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 === true1161
? insertThinkingLevelMessages(converted, activeEffort)1162
: converted.messages,1163
max_tokens: options?.maxTokens ?? model.maxTokens,1164
stream: true,1165
...(betaFeatures.length > 0 ? { betas: betaFeatures } : {}),1166
};1168
// For OAuth tokens, we MUST include Claude Code identity1169
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 tokens1186
params.system = [1187
{1188
type: "text",1189
text: sanitizeSurrogates(initialSystemText),1190
...(cacheControl ? { cache_control: cacheControl } : {}),1191
},1192
];1193
}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.supportsTemperature1201
) {1202
params.temperature = options.temperature;1203
}1205
const toolCacheControl = compat.supportsCacheControlOnTools ? cacheControl : undefined;1206
if (nativeToolChanges) {1207
// Initial tools stay active with the cache breakpoint on the last one. Every later1208
// declaration is deferred and only surfaced by its `tool_addition` block; removed1209
// tools stay declared and are withdrawn by `tool_removal`. The request-level list1210
// 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
}1242
// Managed effort models always use adaptive thinking so prefix mismatches can1243
// 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 like1254
// 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 models1264
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
}1275
if (options?.metadata) {1276
const userId = options.metadata.user_id;1277
if (typeof userId === "string") {1278
params.metadata = { user_id: userId };1279
}1280
}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
}1290
const allowedFallbackModels = model.compat?.allowedFallbackModels;1291
if (allowedFallbackModels && allowedFallbackModels.length > 0) {1292
params.fallbacks = allowedFallbackModels.map((fallback) => ({ model: fallback.model }));1293
}1295
return params;1296
}1298
// Normalize tool call IDs to match Anthropic's required pattern and length1299
function normalizeToolCallId(id: string): string {1300
return id.replace(/[^a-zA-Z0-9_-]/g, "_").slice(0, 64);1301
}1303
function 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
}1312
interface ConvertedAnthropicMessages {1313
messages: MessageParam[];1314
assistantLevels: Map<number, AnthropicEffort>;1315
}1317
function 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 assistant1328
// message (or at the end of the transcript). Anthropic requires `tool_result` blocks to1329
// immediately follow their `tool_use`, so a system message between them is rejected; this1330
// also mirrors where the managed-effort system messages are inserted. As a result an1331
// 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
};1338
for (let i = 0; i < transformedMessages.length; i++) {1339
const msg = transformedMessages[i];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[] = [];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_thinking1413
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 emit1425
// and accept empty signatures, so let marked models preserve the block.1426
if (!hasThinkingSignature) {1427
blocks.push(1428
allowEmptySignature1429
? {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
}1478
// Skip the messages we've already processed.1479
i = j - 1;1481
params.push({1482
role: "user",1483
content: toolResults,1484
});1485
}1486
}1488
flushPendingSystemMessages();1490
// Add cache_control to the last user or system message to cache conversation history1491
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
}1518
return { messages: params, assistantLevels };1519
}1521
function isAnthropicEffort(value: unknown): value is AnthropicEffort {1522
return value === "low" || value === "medium" || value === "high" || value === "xhigh" || value === "max";1523
}1525
function 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
}1541
function shouldUseFineGrainedToolStreamingBeta(1542
model: Model<"anthropic-messages">,1543
context: TranscriptContext,1544
): boolean {1545
return getCurrentTools(context.messages).length > 0 && !getAnthropicCompat(model).supportsEagerToolInputStreaming;1546
}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-limitations1550
const 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
]);1563
const 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
]);1576
const 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
};1583
function convertTools(1584
tools: Tool[],1585
isOAuthToken: boolean,1586
supportsEagerToolInputStreaming: boolean,1587
supportsStrictTools: boolean,1588
cacheControl?: CacheControlEphemeral,1589
): BetaTool[] {1590
if (!tools) return [];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 === true1603
? {1604
...(parameters as Record<string, unknown>),1605
...legacyInputSchema,1606
}1607
: legacyInputSchema;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
}1620
function 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 -> resubmit1637
return { stopReason: "stop" };1638
case "stop_sequence":1639
return { stopReason: "stop" }; // We don't supply stop sequences, so this should never happen1640
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
}