From f80d45aa6bfb451b5eb673ee3eb262bb03b0a3e0 Mon Sep 17 00:00:00 2001 From: can1357 Date: Tue, 9 Jun 2026 23:17:26 +0200 Subject: [PATCH] fix(anthropic-provider): resolved Anthropic request and stream handling - Added a 4-slot concurrency limiter for Anthropic multi-image resize work. - Prioritized client-owned request fields and `claudeCodeVersion`-based `User-Agent` headers. - Ignored terminal `message_delta` envelopes and treated missing `usage`/`delta` payloads as non-fatal. - Stopped ping-only streams from consuming first-event watchdog logic and preserved timeout retries. --- packages/ai/CHANGELOG.md | 17 + packages/ai/src/providers/anthropic-client.ts | 6 +- packages/ai/src/providers/anthropic.ts | 405 ++++++++++++++---- packages/ai/src/registry/oauth/anthropic.ts | 3 +- packages/ai/src/usage/claude.ts | 3 +- packages/ai/src/utils.ts | 2 +- packages/ai/src/utils/retry.ts | 4 + packages/ai/test/anthropic-alignment.test.ts | 140 +++++- packages/ai/test/anthropic-client.test.ts | 19 + packages/ai/test/anthropic-oauth.test.ts | 3 +- packages/ai/test/anthropic-prefill.test.ts | 38 ++ .../ai/test/anthropic-stream-envelope.test.ts | 56 +++ .../ai/test/anthropic-stream-timeout.test.ts | 48 +++ ...anthropic-unsigned-thinking-replay.test.ts | 41 ++ packages/ai/test/claude-usage-headers.test.ts | 3 +- .../github-copilot-anthropic-auth.test.ts | 32 ++ 16 files changed, 721 insertions(+), 99 deletions(-) diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index dcbf1877e..ced7e1a47 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -25,6 +25,23 @@ - Stopped sending the `fast-mode-2026-02-01` beta header once a session has learned the endpoint+model rejects fast mode (`fastModeDisabled` provider state), matching the already-dropped `speed` param. - Stopped `buildAnthropicHeaders` defaulting API-key requests onto the full Claude Code OAuth beta list (`oauth-2025-04-20`, `claude-code-20250219`, …). The `claudeCodeBetas` default is now OAuth-gated, matching the streaming path — the web-search header builder was the only caller hitting the default, so API-key search requests now carry just their own betas (e.g. `web-search-2025-03-05`). An empty `anthropic-beta` header is omitted entirely instead of being sent as an empty string. - Fixed image-bearing `developer` messages being upgraded to mid-conversation `system` turns on Opus 4.8+/Fable/Mythos 5. System content is text-only on the wire, so a developer turn carrying image blocks in an upgrade-eligible position produced a 400; it now stays a `user` message. +- Fixed a spliced reconnect's second envelope overwriting the completed Anthropic message: `message_delta` was not gated by the terminal-stop flag (content events and duplicate `message_start` were), so the splice's `stop_reason`/usage replaced the finished turn's — a `tool_use` turn could be relabeled `stop`, and the harness then never executed the streamed tool calls. Post-terminal deltas are now logged as envelope anomalies and skipped. +- Fixed a `ping` arriving before `message_start` consuming the Anthropic first-event watchdog: the stall was then classified as a terminal mid-stream idle timeout instead of a retryable first-event timeout. Pings no longer count as the first item but still refresh the idle deadline once content is flowing. +- Fixed Anthropic-compatible proxies that omit `usage`/`delta` objects from `message_start`/`message_delta`/`content_block_*` envelopes crashing the turn with an unretryable `TypeError`; the missing payloads now degrade to logged envelope anomalies like every other malformed-frame case. +- Fixed `applyPromptCaching` placing `cache_control` on `thinking`/`redacted_thinking` blocks — Anthropic rejects that with a 400. A thinking-only assistant turn inside the trailing cache window (e.g. followed by the synthetic `Continue.` pad) no longer receives a breakpoint. +- Fixed consecutive `assistant` params reaching the wire when an empty user/developer turn between two assistant turns was dropped by the converter (e.g. an empty "nudge" submission after a length-truncated reply); Anthropic 400s on non-alternating assistant turns, and the broken triple replayed on every subsequent request. A `user: "Continue."` separator is now inserted, mirroring the trailing-prefill fallback. +- Fixed `supportsAdaptiveThinkingDisplay` misparsing bare dated Opus ids: `claude-opus-4-20250514` (Opus 4.0) parsed as minor `20250514` ≥ 4.7, which silently dropped the `interleaved-thinking-2025-05-14` beta for API-key Opus 4.0 requests. +- Fixed `output_config.effort` shipping without the `effort-2025-11-24` beta on thinking-off requests against adaptive-only Claude models (the effort:"low" pin), and the mid-conversation `system` role shipping without `mid-conversation-system-2026-04-07` on API-key and OAuth-utility requests; both betas are now added whenever the request can carry the corresponding field. +- Fixed GitHub Copilot anthropic-messages requests going out with no `Content-Type` and no `anthropic-version` header — the copilot branch builds its headers from scratch and Bun's fetch does not default `Content-Type` for string bodies. Both headers are now pinned to match every other branch. +- Fixed Anthropic client/provider retry multiplication: with the first-event watchdog disabled (`PI_STREAM_FIRST_EVENT_TIMEOUT_MS=0`), the client's internal `maxRetries: 5` reactivated and stacked with the provider loop's 3 retries — up to 24 wire attempts with double backoff. The provider now pins per-request `maxRetries: 0` unconditionally. +- Fixed `AnthropicMessagesClient` spreading `fetchOptions` after the core request fields, letting a caller-supplied `signal`/`method`/`body` silently disconnect the timeout controller or corrupt the request. Transport extras (TLS) still pass through; core fields now always win. +- Fixed Foundry mTLS/CA material being cached for the process lifetime when the env vars point at files: the cache key now folds in the file mtime so on-disk certificate rotation takes effect. +- Fixed the Claude Code fingerprint version drifting across surfaces: the usage endpoint (`claude-cli/2.1.160`) and OAuth bootstrap (`claude-code/2.1.160`) pinned a stale version while `/v1/messages` reported 2.1.165; both now derive from `claudeCodeVersion`. +- Fixed a system prompt that merely *mentions* `x-anthropic-billing-header:` mid-text suppressing the entire Claude Code system-block injection (billing header, instruction, and cch attestation); the resumed-session guard now anchors with `startsWith`. +- Fixed lone surrogates in cross-API tool-call arguments reaching Anthropic's strict UTF-8 validation: replayed OpenAI/Google-origin `tool_use.input` string leaves are now deep-sanitized with `toWellFormed()`, while same-API Anthropic arguments stay byte-identical to keep prompt-cache prefixes stable. +- Bounded the many-image resize fan-out to 4 concurrent decodes (it previously decoded every oversized image at once, two encode pipelines each — multi-GB transient memory at the 20+-image threshold that activates the feature). +- Fixed `mergeHeaders` merging case-sensitively on the Copilot/client-options path, where a miscased user-configured header (e.g. `authorization` next to the synthesized `Authorization`) survived as two keys that the `Headers` constructor joins comma-separated on the wire. +- Hardened the Anthropic stream lifecycle: prologue failures (e.g. a malformed Copilot credential in `buildCopilotDynamicHeaders`) and error-finalization failures now surface as an `error` event instead of an unhandled rejection that left `stream.result()` hanging forever; the spurious "cch billing placeholder not patched" warning no longer fires when the placeholder only appears in user content. ## [15.10.9] - 2026-06-09 diff --git a/packages/ai/src/providers/anthropic-client.ts b/packages/ai/src/providers/anthropic-client.ts index 0639f6397..0e752fa8d 100644 --- a/packages/ai/src/providers/anthropic-client.ts +++ b/packages/ai/src/providers/anthropic-client.ts @@ -43,7 +43,9 @@ export interface AnthropicRequestOptions { /** * Extra `RequestInit` fields merged into every fetch call. Bun extends * `RequestInit` with a `tls` option used for the Claude Code TLS profile and - * Foundry mTLS. + * Foundry mTLS. Core request fields (`method`, `headers`, `body`, `signal`) + * are owned by the client and cannot be overridden from here — the timeout + * controller's signal in particular must always win. */ export type AnthropicFetchOptions = RequestInit & { tls?: { @@ -288,11 +290,11 @@ export class AnthropicMessagesClient implements AnthropicMessagesClientLike { callerSignal?.addEventListener("abort", onAbort, { once: true }); try { return await fetchFn(url, { + ...(this.#options.fetchOptions ?? {}), method: "POST", headers, body, signal: controller.signal, - ...(this.#options.fetchOptions ?? {}), }); } catch (error) { if (timedOut && !callerSignal?.aborted) throw new AnthropicConnectionTimeoutError(); diff --git a/packages/ai/src/providers/anthropic.ts b/packages/ai/src/providers/anthropic.ts index bfc1a93fb..e134c3f86 100644 --- a/packages/ai/src/providers/anthropic.ts +++ b/packages/ai/src/providers/anthropic.ts @@ -124,6 +124,7 @@ export function buildBetaHeader(baseBetas: readonly string[], extraBetas: readon return result.join(","); } +const midConversationSystemBeta = "mid-conversation-system-2026-04-07"; const claudeCodeUtilityBetaDefaults = [ "oauth-2025-04-20", "interleaved-thinking-2025-05-14", @@ -137,7 +138,7 @@ const claudeCodeAgentBetaDefaults = [ "interleaved-thinking-2025-05-14", "context-management-2025-06-27", "prompt-caching-scope-2026-01-05", - "mid-conversation-system-2026-04-07", + midConversationSystemBeta, "advanced-tool-use-2025-11-20", ] as const; const claudeCodeAgentPostEffortBetas = ["extended-cache-ttl-2025-04-11"] as const; @@ -243,7 +244,7 @@ export function buildAnthropicHeaders(options: AnthropicHeaderOptions): Record CCH_BILLING_SEARCH_WINDOW) return false; + if (idx === -1 || idx - searchFrom > CCH_BILLING_SEARCH_WINDOW) return "unanchored"; // Hash the body with the placeholder in place (matches CC's in-place behaviour). const h = Bun.hash.xxHash64(body, CCH_SEED); const cch = (h & 0xfffffn).toString(16).padStart(5, "0"); for (let i = 0; i < 5; i++) body[idx + 4 + i] = cch.charCodeAt(i); - return true; + return "patched"; } /** @@ -584,12 +587,13 @@ export function wrapFetchForCch(base: FetchImpl): FetchImpl { return (input, init) => { if (init?.body && typeof init.body === "string" && init.body.includes(CCH_PLACEHOLDER_STR)) { const encoded = cchEncoder.encode(init.body); - if (!patchCch(encoded)) { - // The OAuth billing placeholder is present but we couldn't anchor it to - // system[0] — e.g. an `onPayload` hook reordered the first system block's keys + if (patchCch(encoded) === "unanchored") { + // The OAuth billing placeholder is anchored to system[0] but we couldn't + // patch it — e.g. an `onPayload` hook reordered the first system block's keys // so BILLING_SYSTEM_MARKER no longer matches. Send the body as-is (cch stays // `00000`, the prior behaviour) rather than failing the request, but surface the - // fingerprint regression instead of letting it ship silently. + // fingerprint regression instead of letting it ship silently. A `cch=00000` + // literal in user content alone ("no-billing-header") is not a regression. logger.warn("anthropic: cch billing placeholder present but not patched; sending unattested request"); } return base(input, { ...init, body: encoded }); @@ -740,6 +744,37 @@ function countAnthropicImageBlocks(messages: Message[]): number { return count; } +const ANTHROPIC_IMAGE_RESIZE_CONCURRENCY = 4; + +type ResizeLimiter = (fn: () => Promise) => Promise; + +/** + * Bounded-concurrency gate for image decode/encode work. The many-image path + * fans out over every block of every message; unbounded, 100+ oversized images + * would decode concurrently (two encode pipelines each) and spike memory by + * gigabytes. Slots are handed off directly to the next waiter on release. + */ +function createResizeLimiter(limit: number): ResizeLimiter { + let active = 0; + const queue: (() => void)[] = []; + return async fn => { + if (active >= limit) { + const { promise, resolve } = Promise.withResolvers(); + queue.push(resolve); + await promise; + } else { + active++; + } + try { + return await fn(); + } finally { + const next = queue.shift(); + if (next) next(); + else active--; + } + }; +} + async function resizeAnthropicManyImageBlock(block: ImageContent): Promise { try { const inputBuffer = Buffer.from(block.data, "base64"); @@ -775,12 +810,13 @@ async function resizeAnthropicManyImageBlock(block: ImageContent): Promise { let changed = false; const next = await Promise.all( content.map(async block => { if (block.type !== "image") return block; - const resized = await resizeAnthropicManyImageBlock(block); + const resized = await limit(() => resizeAnthropicManyImageBlock(block)); if (resized !== block) { changed = true; state.resized++; @@ -791,14 +827,18 @@ async function resizeAnthropicManyImageContent( return changed ? next : content; } -async function resizeAnthropicManyImageMessage(message: Message, state: { resized: number }): Promise { +async function resizeAnthropicManyImageMessage( + message: Message, + state: { resized: number }, + limit: ResizeLimiter, +): Promise { if (message.role === "user" || message.role === "developer") { if (!Array.isArray(message.content)) return message; - const content = await resizeAnthropicManyImageContent(message.content, state); + const content = await resizeAnthropicManyImageContent(message.content, state, limit); return content === message.content ? message : { ...message, content }; } if (message.role === "toolResult") { - const content = await resizeAnthropicManyImageContent(message.content, state); + const content = await resizeAnthropicManyImageContent(message.content, state, limit); return content === message.content ? message : { ...message, content }; } return message; @@ -811,9 +851,10 @@ async function prepareAnthropicManyImageContext(context: Context, supportsImages let changed = false; const state = { resized: 0 }; + const limit = createResizeLimiter(ANTHROPIC_IMAGE_RESIZE_CONCURRENCY); const messages = await Promise.all( context.messages.map(async message => { - const next = await resizeAnthropicManyImageMessage(message, state); + const next = await resizeAnthropicManyImageMessage(message, state, limit); if (next !== message) changed = true; return next; }), @@ -1006,11 +1047,27 @@ type FoundryTlsOptions = { const foundryTlsOptionsCache = new Map(); +function foundryTlsCacheKeyComponent(value: string | undefined): string | null { + if (!value) return null; + const trimmed = value.trim(); + // For path-valued vars, fold the file mtime into the key so on-disk cert + // rotation (common for short-lived corporate mTLS certs) invalidates the + // cached TLS options instead of pinning the first read forever. + if (trimmed && !trimmed.includes("-----BEGIN") && looksLikeFilePath(trimmed)) { + try { + return `${trimmed}@${fs.statSync(trimmed).mtimeMs}`; + } catch { + return trimmed; + } + } + return value; +} + function foundryTlsOptionsCacheKey(): string { return JSON.stringify([ - $env.NODE_EXTRA_CA_CERTS ?? null, - $env.CLAUDE_CODE_CLIENT_CERT ?? null, - $env.CLAUDE_CODE_CLIENT_KEY ?? null, + foundryTlsCacheKeyComponent($env.NODE_EXTRA_CA_CERTS), + foundryTlsCacheKeyComponent($env.CLAUDE_CODE_CLIENT_CERT), + foundryTlsCacheKeyComponent($env.CLAUDE_CODE_CLIENT_KEY), ]); } @@ -1150,10 +1207,19 @@ function buildClaudeCodeTlsFetchOptions( }; } function mergeHeaders(...headerSources: (Record | undefined)[]): Record { + // Case-insensitive merge: later sources win and keep their casing. A plain + // Object.assign would let `authorization` and `Authorization` coexist, and + // the Headers constructor then joins both values comma-separated on the wire. const merged: Record = {}; + const keyByLower = new Map(); for (const headers of headerSources) { - if (headers) { - Object.assign(merged, headers); + if (!headers) continue; + for (const [key, value] of Object.entries(headers)) { + const lower = key.toLowerCase(); + const existing = keyByLower.get(lower); + if (existing !== undefined && existing !== key) delete merged[existing]; + keyByLower.set(lower, key); + merged[key] = value; } } return merged; @@ -1333,19 +1399,10 @@ function reportAnthropicEnvelopeAnomaly(detail: string): void { logger.warn(`anthropic: ignoring malformed stream envelope: ${detail}`); } -const ANTHROPIC_PRE_MESSAGE_START_EVENT_TYPES = new Set([ - "content_block_start", - "content_block_delta", - "content_block_stop", - "message_delta", - "message_stop", - "message_start", -]); - function shouldIgnoreAnthropicPreambleEvent(eventType: unknown): boolean { if (typeof eventType !== "string") return false; if (eventType === "ping") return true; - return !ANTHROPIC_PRE_MESSAGE_START_EVENT_TYPES.has(eventType); + return !ANTHROPIC_MESSAGE_EVENTS.has(eventType); } function isTransientStreamEnvelopeError(error: unknown): boolean { @@ -1450,23 +1507,13 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = ( const startTime = Date.now(); let firstTokenTime: number | undefined; - const copilotDynamicHeaders = - model.provider === "github-copilot" - ? buildCopilotDynamicHeaders({ - messages: context.messages, - hasImages: hasCopilotVisionInput(context.messages), - premiumMultiplier: model.premiumMultiplier, - headers: { ...(model.headers ?? {}), ...(options?.headers ?? {}) }, - initiatorOverride: options?.initiatorOverride, - }) - : undefined; const output: AssistantMessage = { role: "assistant", content: [], api: model.api as Api, provider: model.provider, model: model.id, - usage: createEmptyUsage(copilotDynamicHeaders?.premiumRequests), + usage: createEmptyUsage(), stopReason: "stop", timestamp: Date.now(), }; @@ -1477,6 +1524,22 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = ( const rawSseObserver = onSseEvent ? (event: RawSseEvent) => onSseEvent(event, model) : undefined; try { + // Built inside the try so a copilot credential/header failure surfaces as + // an error event instead of an unhandled rejection that leaves the stream + // (and any consumer awaiting `result()`) hanging forever. + const copilotDynamicHeaders = + model.provider === "github-copilot" + ? buildCopilotDynamicHeaders({ + messages: context.messages, + hasImages: hasCopilotVisionInput(context.messages), + premiumMultiplier: model.premiumMultiplier, + headers: { ...(model.headers ?? {}), ...(options?.headers ?? {}) }, + initiatorOverride: options?.initiatorOverride, + }) + : undefined; + if (copilotDynamicHeaders?.premiumRequests !== undefined) { + output.usage.premiumRequests = copilotDynamicHeaders.premiumRequests; + } const apiKey = options?.apiKey ?? getEnvApiKey(model.provider) ?? ""; const baseUrl = resolveAnthropicBaseUrl(model, apiKey) ?? "https://api.anthropic.com"; const providerSessionState = getAnthropicProviderSessionState( @@ -1506,9 +1569,30 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = ( if (options?.taskBudget && !extraBetas.includes(taskBudgetBeta)) { extraBetas.push(taskBudgetBeta); } - if (options?.thinkingEnabled && model.reasoning && !extraBetas.includes(effortBeta)) { + // `output_config.effort` ships on thinking-on requests AND on the + // thinking-off adaptive pin (adaptive-only models get effort:"low" so + // the toggle cannot 400); the beta must accompany the field in both. + const sendsAdaptiveEffortPin = + options?.thinkingEnabled === false && + model.thinking?.mode === "anthropic-adaptive" && + !getAnthropicCompat(model).disableAdaptiveThinking; + if ( + model.reasoning && + (options?.thinkingEnabled || sendsAdaptiveEffortPin) && + !extraBetas.includes(effortBeta) + ) { extraBetas.push(effortBeta); } + if ( + getAnthropicCompat(model).supportsMidConversationSystem && + !extraBetas.includes(midConversationSystemBeta) + ) { + // convertAnthropicMessages may upgrade developer turns to the + // mid-conversation `system` role on these models; API-key requests + // need the beta alongside the role (OAuth agent requests already + // carry it in the Claude Code list). + extraBetas.push(midConversationSystemBeta); + } const created = createClient(model, { model, @@ -1598,7 +1682,11 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = ( while (true) { activeAbortTracker = createAbortSourceTracker(options?.signal); const { requestSignal } = activeAbortTracker; - const requestOptions = createSdkStreamRequestOptions(requestSignal, requestTimeoutMs); + // The provider loop owns retries: pin the client's internal retry loop + // to zero even when no watchdog timeout is configured (the helper only + // pins it alongside a timeout; the client default of 5 would otherwise + // multiply with PROVIDER_MAX_RETRIES into up to 24 wire attempts). + const requestOptions = { ...createSdkStreamRequestOptions(requestSignal, requestTimeoutMs), maxRetries: 0 }; const anthropicRequest: unknown = isOAuthToken && client.beta ? client.beta.messages.create({ ...params, stream: true }, requestOptions) @@ -1642,6 +1730,11 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = ( { contentIndex: number; kind: "text" | "thinking" | "redactedThinking" | "toolCall" | "ignored" } >(); + // Pings keep the idle deadline alive once content is flowing, but a + // ping before message_start must not consume the first-event watchdog: + // it would flip the (retryable) pre-content stall classification into + // a terminal mid-stream idle timeout. + let sawNonPingEvent = false; const timedAnthropicStream = iterateWithIdleTimeout(anthropicStream, { idleTimeoutMs, firstItemTimeoutMs: firstEventTimeoutMs, @@ -1650,6 +1743,11 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = ( onIdle: () => activeAbortTracker.abortLocally(idleTimeoutAbortError), onFirstItemTimeout: () => activeAbortTracker.abortLocally(firstEventTimeoutAbortError), abortSignal: options?.signal, + isProgressItem: item => { + if ((item as AnthropicStreamEvent).type === "ping") return sawNonPingEvent; + sawNonPingEvent = true; + return true; + }, }); const observedAnthropicStream = rawSseObserver && !recordsRawSseEvents @@ -1666,15 +1764,21 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = ( continue; } sawMessageStart = true; - applyAnthropicUsageExtras(output.usage, event.message.usage); - output.responseId = event.message.id; - output.usage.input = event.message.usage.input_tokens || 0; - output.usage.output = event.message.usage.output_tokens || 0; - output.usage.cacheRead = event.message.usage.cache_read_input_tokens || 0; - output.usage.cacheWrite = event.message.usage.cache_creation_input_tokens || 0; - output.usage.totalTokens = - output.usage.input + output.usage.output + output.usage.cacheRead + output.usage.cacheWrite; - calculateCost(model, output.usage); + const startMessage = event.message; + if (startMessage?.id) output.responseId = startMessage.id; + const startUsage = startMessage?.usage; + if (startUsage) { + applyAnthropicUsageExtras(output.usage, startUsage); + output.usage.input = startUsage.input_tokens || 0; + output.usage.output = startUsage.output_tokens || 0; + output.usage.cacheRead = startUsage.cache_read_input_tokens || 0; + output.usage.cacheWrite = startUsage.cache_creation_input_tokens || 0; + output.usage.totalTokens = + output.usage.input + output.usage.output + output.usage.cacheRead + output.usage.cacheWrite; + calculateCost(model, output.usage); + } else { + reportAnthropicEnvelopeAnomaly("message_start missing usage"); + } continue; } @@ -1694,6 +1798,10 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = ( reportAnthropicEnvelopeAnomaly(`duplicate content_block_start index ${event.index}`); continue; } + if (!event.content_block?.type) { + reportAnthropicEnvelopeAnomaly("content_block_start missing content_block payload"); + continue; + } if (!firstTokenTime) firstTokenTime = Date.now(); if (event.content_block.type === "text") { streamedReplayUnsafeContent = true; @@ -1774,6 +1882,10 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = ( continue; } if (openBlock.kind === "ignored") continue; + if (!event.delta?.type) { + reportAnthropicEnvelopeAnomaly("content_block_delta missing delta payload"); + continue; + } const block = blocks[openBlock.contentIndex]; if (event.delta.type === "text_delta") { if (openBlock.kind !== "text" || block?.type !== "text") { @@ -1851,44 +1963,56 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = ( openBlocks.delete(event.index); finalizeStreamBlock(block, openBlock.contentIndex); } else if (event.type === "message_delta") { - const rawStopReason = event.delta.stop_reason; + if (sawTerminalEnvelope) { + // A spliced reconnect's second envelope must not overwrite the + // completed message's stop reason or usage. + reportAnthropicEnvelopeAnomaly("received message_delta after terminal stop signal"); + continue; + } + const delta = event.delta; + const rawStopReason = delta?.stop_reason; if (rawStopReason) { output.stopReason = mapStopReason(rawStopReason); sawTerminalEnvelope = true; } - const stopDetails = event.delta.stop_details; - if (stopDetails && stopDetails.type === "refusal") { - const explanation = stopDetails.explanation?.trim(); - const category = stopDetails.category; - const label = category ? `Refusal (${category})` : "Refusal"; - output.errorMessage = explanation ? `${label}: ${explanation}` : label; - } else if (output.stopReason === "error" && !output.errorMessage) { - // Anthropic flagged an error-class stop (refusal / sensitive) without - // populating stop_details. Surface the raw reason instead of falling - // through to the generic "unknown error" string when we throw below. - output.errorMessage = - rawStopReason === "refusal" - ? "Refusal (no details provided)" - : rawStopReason === "sensitive" - ? "Content flagged by safety filters" - : `Anthropic stream ended with stop_reason: ${rawStopReason ?? "unknown"}`; + if (output.stopReason === "error") { + const stopDetails = delta?.stop_details; + if (stopDetails?.type === "refusal") { + const explanation = stopDetails.explanation?.trim(); + const category = stopDetails.category; + const label = category ? `Refusal (${category})` : "Refusal"; + output.errorMessage = explanation ? `${label}: ${explanation}` : label; + } else if (!output.errorMessage) { + // Anthropic flagged an error-class stop (refusal / sensitive) without + // populating stop_details. Surface the raw reason instead of falling + // through to the generic "unknown error" string when we throw below. + output.errorMessage = + rawStopReason === "refusal" + ? "Refusal (no details provided)" + : rawStopReason === "sensitive" + ? "Content flagged by safety filters" + : `Anthropic stream ended with stop_reason: ${rawStopReason ?? "unknown"}`; + } } - if (event.usage.input_tokens != null) { - output.usage.input = event.usage.input_tokens; + const deltaUsage = event.usage; + if (deltaUsage) { + if (deltaUsage.input_tokens != null) { + output.usage.input = deltaUsage.input_tokens; + } + if (deltaUsage.output_tokens != null) { + output.usage.output = deltaUsage.output_tokens; + } + if (deltaUsage.cache_read_input_tokens != null) { + output.usage.cacheRead = deltaUsage.cache_read_input_tokens; + } + if (deltaUsage.cache_creation_input_tokens != null) { + output.usage.cacheWrite = deltaUsage.cache_creation_input_tokens; + } + applyAnthropicUsageExtras(output.usage, deltaUsage); + output.usage.totalTokens = + output.usage.input + output.usage.output + output.usage.cacheRead + output.usage.cacheWrite; + calculateCost(model, output.usage); } - if (event.usage.output_tokens != null) { - output.usage.output = event.usage.output_tokens; - } - if (event.usage.cache_read_input_tokens != null) { - output.usage.cacheRead = event.usage.cache_read_input_tokens; - } - if (event.usage.cache_creation_input_tokens != null) { - output.usage.cacheWrite = event.usage.cache_creation_input_tokens; - } - applyAnthropicUsageExtras(output.usage, event.usage); - output.usage.totalTokens = - output.usage.input + output.usage.output + output.usage.cacheRead + output.usage.cacheWrite; - calculateCost(model, output.usage); } else if (event.type === "message_stop") { sawTerminalEnvelope = true; sawMessageStop = true; @@ -1946,6 +2070,7 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = ( providerRetryAttempt = 0; output.content.length = 0; output.responseId = undefined; + output.errorMessage = undefined; output.providerPayload = undefined; output.usage = createEmptyUsage(copilotDynamicHeaders?.premiumRequests); output.stopReason = "stop"; @@ -1970,6 +2095,7 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = ( providerRetryAttempt = 0; output.content.length = 0; output.responseId = undefined; + output.errorMessage = undefined; output.providerPayload = undefined; output.usage = createEmptyUsage(copilotDynamicHeaders?.premiumRequests); output.stopReason = "stop"; @@ -2026,8 +2152,15 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = ( const firstEventTimeoutError = activeAbortTracker.getLocalAbortReason(); output.stopReason = activeAbortTracker.wasCallerAbort() ? "aborted" : "error"; output.errorStatus = extractHttpStatusFromError(error); - output.errorMessage = firstEventTimeoutError?.message ?? (await finalizeErrorMessage(error, rawRequestDump)); - output.errorMessage = rewriteCopilotError(output.errorMessage, error, model.provider); + try { + output.errorMessage = + firstEventTimeoutError?.message ?? (await finalizeErrorMessage(error, rawRequestDump)); + output.errorMessage = rewriteCopilotError(output.errorMessage, error, model.provider); + } catch { + // finalizeErrorMessage must never take the stream down with it — a + // throw here would skip stream.end() and hang result() forever. + output.errorMessage = error instanceof Error ? error.message : String(error); + } output.duration = Date.now() - startTime; if (firstTokenTime) output.ttft = firstTokenTime - startTime; stream.push({ type: "error", reason: output.stopReason, error: output }); @@ -2069,7 +2202,7 @@ export function buildAnthropicSystemBlocks( const { includeClaudeCodeInstruction = false, extraInstructions = [], firstUserMessageText, cacheControl } = options; const sanitizedPrompts = normalizeSystemPrompts(systemPrompt); const trimmedInstructions = extraInstructions.map(instruction => instruction.trim()).filter(Boolean); - const hasBillingHeader = sanitizedPrompts.some(prompt => prompt.includes(CLAUDE_BILLING_HEADER_PREFIX)); + const hasBillingHeader = sanitizedPrompts.some(prompt => prompt.startsWith(CLAUDE_BILLING_HEADER_PREFIX)); if (includeClaudeCodeInstruction && !hasBillingHeader) { const blocks: AnthropicSystemBlock[] = [ @@ -2143,6 +2276,8 @@ export function buildAnthropicClientOptions(args: AnthropicClientOptionsArgs): A const defaultHeaders = mergeHeaders( { Accept: stream ? "text/event-stream" : "application/json", + "Content-Type": "application/json", + "anthropic-version": "2023-06-01", "Anthropic-Dangerous-Direct-Browser-Access": "true", Authorization: `Bearer ${copilotApiKey}`, ...(betaFeatures.length > 0 ? { "anthropic-beta": buildBetaHeader([], betaFeatures) } : {}), @@ -2325,7 +2460,16 @@ function applyCacheControlToLastTextBlock( return true; } } - return applyCacheControlToLastBlock(blocks, cacheControl); + // No text block — fall back to the last block that accepts cache_control; + // thinking/redacted_thinking blocks reject the field with a 400. + for (let i = blocks.length - 1; i >= 0; i--) { + const type = blocks[i].type; + if (type === "thinking" || type === "redacted_thinking") continue; + if (blocks[i].cache_control != null) return false; + blocks[i] = { ...blocks[i], cache_control: cloneAnthropicCacheControl(cacheControl) }; + return true; + } + return false; } function applyPromptCaching(params: MessageCreateParamsStreaming, cacheControl?: AnthropicCacheControl): void { @@ -2772,6 +2916,37 @@ function buildToolResultBlock(model: Model<"anthropic-messages">, msg: ToolResul */ export type AnthropicMessageParam = MessageParam; +/** + * Recursively replace lone surrogates in string leaves. Identity-preserving: + * returns the input object/array when nothing changed. + */ +function toWellFormedDeep(value: unknown): unknown { + if (typeof value === "string") { + const wellFormed = value.toWellFormed(); + return wellFormed === value ? value : wellFormed; + } + if (Array.isArray(value)) { + let changed = false; + const next = value.map(entry => { + const sanitized = toWellFormedDeep(entry); + if (sanitized !== entry) changed = true; + return sanitized; + }); + return changed ? next : value; + } + if (isRecord(value)) { + let changed = false; + const next: Record = {}; + for (const [key, entry] of Object.entries(value)) { + const sanitized = toWellFormedDeep(entry); + if (sanitized !== entry) changed = true; + next[key] = sanitized; + } + return changed ? next : value; + } + return value; +} + export function convertAnthropicMessages( messages: Message[], model: Model<"anthropic-messages">, @@ -2871,7 +3046,13 @@ export function convertAnthropicMessages( type: "tool_use", id: block.id, name: isOAuthToken ? applyClaudeToolPrefix(block.name) : block.name, - input: block.arguments ?? {}, + // Anthropic-origin arguments are guaranteed well-formed (they came + // from the API's own JSON); cross-API replays can carry lone + // surrogates that Anthropic's strict UTF-8 validation rejects. + input: + msg.api === "anthropic-messages" + ? (block.arguments ?? {}) + : toWellFormedDeep(block.arguments ?? {}), }); } } @@ -2927,6 +3108,14 @@ export function convertAnthropicMessages( } } } + // Dropped empty user/developer turns can leave two assistant params adjacent; + // the API rejects consecutive assistant messages. Repair with the same neutral + // nudge used for trailing-assistant prefill below. + for (let i = params.length - 1; i > 0; i--) { + if (params[i].role === "assistant" && params[i - 1]?.role === "assistant") { + params.splice(i, 0, { role: "user", content: "Continue." }); + } + } if (params.length > 0 && params[params.length - 1]?.role === "assistant") { params.push({ role: "user", content: "Continue." }); } @@ -3328,6 +3517,38 @@ function normalizeAnthropicStrictSchemaNode( return result; } +const ANTHROPIC_STRICT_INCOMPATIBLE_KEYWORDS = [ + "oneOf", + "allOf", + "$ref", + "patternProperties", + "propertyNames", +] as const; + +/** + * Anthropic's strict grammar subset supports anyOf/type-array unions only. + * oneOf/allOf/$ref compile unpredictably (rejections arrive as 400s the + * grammar-too-large fallback does not recognize, so they would hard-fail the + * turn), and patternProperties/propertyNames describe open key sets that the + * strict pipeline's injected `additionalProperties: false` would contradict. + * Runs against the raw wire schema — the base normalizer spills several of + * these keywords into the description, erasing the evidence. + */ +function hasAnthropicStrictIncompatibleKeyword(schema: unknown, seen = new Set()): boolean { + if (Array.isArray(schema)) { + if (seen.has(schema)) return false; + seen.add(schema); + return schema.some(entry => hasAnthropicStrictIncompatibleKeyword(entry, seen)); + } + if (!isRecord(schema)) return false; + if (seen.has(schema)) return false; + seen.add(schema); + for (const keyword of ANTHROPIC_STRICT_INCOMPATIBLE_KEYWORDS) { + if (schema[keyword] !== undefined) return true; + } + return Object.values(schema).some(value => hasAnthropicStrictIncompatibleKeyword(value, seen)); +} + function normalizeAnthropicStrictSchema( schema: Record, optionalRemaining: number, @@ -3367,7 +3588,9 @@ function buildAnthropicToolSchemaPlans(tools: Tool[], disableStrictTools = false const candidateIndexes = tools.flatMap((tool, index) => { if (!ANTHROPIC_STRICT_TOOL_ALLOWLIST.has(tool.name)) return []; - return tool.strict === false ? [] : [index]; + if (tool.strict === false) return []; + if (hasAnthropicStrictIncompatibleKeyword(toolWireSchema(tool))) return []; + return [index]; }); let strictToolCount = 0; diff --git a/packages/ai/src/registry/oauth/anthropic.ts b/packages/ai/src/registry/oauth/anthropic.ts index 22d530ee8..6c9d0eca4 100644 --- a/packages/ai/src/registry/oauth/anthropic.ts +++ b/packages/ai/src/registry/oauth/anthropic.ts @@ -2,6 +2,7 @@ * Anthropic OAuth flow (Claude Pro/Max) */ +import { claudeCodeVersion } from "../../providers/anthropic"; import type { FetchImpl } from "../../types"; import { OAuthCallbackFlow } from "./callback-server"; import { generatePKCE } from "./pkce"; @@ -13,7 +14,7 @@ const AUTHORIZE_URL = "https://claude.ai/oauth/authorize"; const TOKEN_URL = "https://api.anthropic.com/v1/oauth/token"; const BOOTSTRAP_URL = "https://api.anthropic.com/api/claude_cli/bootstrap"; const CLAUDE_CODE_BOOTSTRAP_MODEL = "claude-opus-4-8"; -const CLAUDE_CODE_BOOTSTRAP_USER_AGENT = "claude-code/2.1.160"; +const CLAUDE_CODE_BOOTSTRAP_USER_AGENT = `claude-code/${claudeCodeVersion}`; const CALLBACK_PORT = 54545; const CALLBACK_PATH = "/callback"; // Scopes required for direct OAuth-token inference (user:inference) plus account/session management. diff --git a/packages/ai/src/usage/claude.ts b/packages/ai/src/usage/claude.ts index 1b4f389c1..a17b2f9ed 100644 --- a/packages/ai/src/usage/claude.ts +++ b/packages/ai/src/usage/claude.ts @@ -1,4 +1,5 @@ import { scheduler } from "node:timers/promises"; +import { claudeCodeVersion } from "../providers/anthropic"; import type { CredentialRankingStrategy, UsageAmount, @@ -24,7 +25,7 @@ const CLAUDE_HEADERS = { "anthropic-beta": "claude-code-20250219,oauth-2025-04-20,interleaved-thinking-2025-05-14,redact-thinking-2026-02-12,context-management-2025-06-27,prompt-caching-scope-2026-01-05,mid-conversation-system-2026-04-07,advanced-tool-use-2025-11-20,effort-2025-11-24,extended-cache-ttl-2025-04-11", "content-type": "application/json", - "user-agent": "claude-cli/2.1.160 (external, cli)", + "user-agent": `claude-cli/${claudeCodeVersion} (external, cli)`, connection: "keep-alive", } as const; diff --git a/packages/ai/src/utils.ts b/packages/ai/src/utils.ts index 226511d91..e6e3bc74f 100644 --- a/packages/ai/src/utils.ts +++ b/packages/ai/src/utils.ts @@ -8,7 +8,7 @@ export { isRecord } from "@oh-my-pi/pi-utils"; export function normalizeSystemPrompts(systemPrompt: readonly string[] | string | undefined | null): string[] { if (systemPrompt === undefined || systemPrompt === null) return []; const prompts = Array.isArray(systemPrompt) ? systemPrompt : typeof systemPrompt === "string" ? [systemPrompt] : []; - return prompts.map(prompt => prompt.toWellFormed()).filter(prompt => prompt.length > 0); + return prompts.map(prompt => prompt.toWellFormed()).filter(prompt => prompt.trim().length > 0); } export function toNumber(value: unknown): number | undefined { diff --git a/packages/ai/src/utils/retry.ts b/packages/ai/src/utils/retry.ts index 732f54914..8a5539186 100644 --- a/packages/ai/src/utils/retry.ts +++ b/packages/ai/src/utils/retry.ts @@ -45,6 +45,10 @@ export async function callWithCopilotModelRetry( return await fn(); } catch (error) { lastError = error; + // A latched abort (caller cancel or local watchdog) makes any retry a + // guaranteed-dead attempt — surface the original error, not the + // scheduler's AbortError. + if (options.signal?.aborted) throw error; if (!isCopilotTransientModelError(error) && !isRetryableError(error)) throw error; if (attempt === COPILOT_MODEL_RETRY_MAX_ATTEMPTS - 1) break; await scheduler.wait(retryBaseDelayMs * (attempt + 1), { signal: options.signal }); diff --git a/packages/ai/test/anthropic-alignment.test.ts b/packages/ai/test/anthropic-alignment.test.ts index cf52aa1e3..526b416f2 100644 --- a/packages/ai/test/anthropic-alignment.test.ts +++ b/packages/ai/test/anthropic-alignment.test.ts @@ -22,7 +22,7 @@ import { stripClaudeToolPrefix, } from "@oh-my-pi/pi-ai/providers/anthropic"; import { getEnvApiKey } from "@oh-my-pi/pi-ai/stream"; -import type { Context, Model, TJsonSchema, TokenTaskBudget, Tool } from "@oh-my-pi/pi-ai/types"; +import type { AssistantMessage, Context, Model, TJsonSchema, TokenTaskBudget, Tool } from "@oh-my-pi/pi-ai/types"; import * as z from "zod/v4"; import { withEnv } from "./helpers"; @@ -313,6 +313,74 @@ describe("Anthropic request fingerprint alignment", () => { expect(payload.max_tokens).toBe(128_000); }); + it("does not place cache_control on thinking blocks in the trailing cache window", async () => { + const thinkingOnlyAssistant: AssistantMessage = { + role: "assistant", + content: [{ type: "thinking", thinking: "long deliberation", thinkingSignature: "sig-1" }], + api: "anthropic-messages", + provider: "anthropic", + model: ANTHROPIC_MODEL.id, + usage: { + input: 0, + output: 0, + cacheRead: 0, + cacheWrite: 0, + totalTokens: 0, + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, + }, + stopReason: "stop", + timestamp: Date.now(), + }; + const payload = (await captureAnthropicPayload( + ANTHROPIC_MODEL, + { + systemPrompt: ["Stay concise."], + messages: [{ role: "user", content: "Think about it", timestamp: Date.now() }, thinkingOnlyAssistant], + }, + { isOAuth: false }, + )) as { messages?: Array<{ role: string; content: string | Array<{ type: string; cache_control?: unknown }> }> }; + + // The thinking-only assistant turn sits inside the trailing two-message + // cache window (the Continue. pad is appended after it) but must not get + // a breakpoint — Anthropic rejects cache_control on thinking blocks. + const assistant = payload.messages?.find(message => message.role === "assistant"); + expect(Array.isArray(assistant?.content)).toBe(true); + for (const block of assistant?.content as Array<{ type: string; cache_control?: unknown }>) { + expect(block.cache_control).toBeUndefined(); + } + const last = payload.messages?.at(-1); + expect((last?.content as Array<{ cache_control?: unknown }>)[0]?.cache_control).toBeDefined(); + }); + + it("adds effort and mid-conversation betas to API-key requests that use those features", async () => { + let capturedBeta: string | undefined; + const fetchMock = (async (_input: string | URL | Request, init?: RequestInit) => { + capturedBeta = (init?.headers as Record | undefined)?.["anthropic-beta"]; + return new Response( + JSON.stringify({ type: "error", error: { type: "invalid_request_error", message: "captured" } }), + { status: 400, headers: { "Content-Type": "application/json" } }, + ); + }) as typeof fetch; + const adaptiveModel: Model<"anthropic-messages"> = { + ...ANTHROPIC_MODEL, + id: "claude-opus-4-8-20260528", + name: "Claude Opus 4.8", + thinking: { mode: "anthropic-adaptive", minLevel: Effort.Minimal, maxLevel: Effort.XHigh }, + }; + + await streamAnthropic( + adaptiveModel, + { systemPrompt: ["Stay concise."], messages: [{ role: "user", content: "Hi", timestamp: Date.now() }] }, + { apiKey: "sk-ant-api-test", thinkingEnabled: false, fetch: fetchMock }, + ).result(); + + // thinking-off on an adaptive-only model still pins output_config.effort, + // and the converter may emit mid-conversation system turns on Opus 4.8 — + // both fields need their betas on API-key requests too. + expect(capturedBeta).toContain("effort-2025-11-24"); + expect(capturedBeta).toContain("mid-conversation-system-2026-04-07"); + }); + it("billing-header fingerprint uses first user message, not leading developer message", async () => { const userText = "Hello from user with enough chars padding here"; @@ -1146,6 +1214,76 @@ describe("Anthropic request fingerprint alignment", () => { expect(strictNames).toEqual(["python"]); }); + it("demotes allowlisted tools with strict-incompatible schema keywords to non-strict", async () => { + const tools: Tool[] = [ + { + name: "edit", + description: "Edit a value", + parameters: { + type: "object", + properties: { q: { oneOf: [{ type: "string" }, { type: "integer" }] } }, + required: ["q"], + } as TJsonSchema, + }, + { + name: "python", + description: "python tool", + parameters: { + type: "object", + properties: { tagged: { type: "object", patternProperties: { "^x-": { type: "string" } } } }, + required: ["tagged"], + } as TJsonSchema, + }, + { + name: "find", + description: "find tool", + parameters: { + type: "object", + properties: { pattern: { type: "string" } }, + required: ["pattern"], + } as TJsonSchema, + }, + ]; + const payload = (await captureAnthropicPayload( + ANTHROPIC_MODEL, + { + systemPrompt: ["Stay concise."], + messages: [{ role: "user", content: "Hi", timestamp: Date.now() }], + tools, + }, + { isOAuth: false }, + )) as { tools?: Array<{ name?: string; strict?: boolean }> }; + + // oneOf/allOf/$ref compile unpredictably under the strict grammar and + // patternProperties contradicts the injected additionalProperties:false; + // such tools must stay non-strict while clean allowlisted tools keep it. + const strictNames = (payload.tools ?? []).filter(tool => tool.strict === true).map(tool => tool.name); + expect(strictNames).toEqual(["find"]); + }); + + it("keeps the interleaved-thinking beta for dated Opus 4.0 ids", () => { + const legacy = buildAnthropicClientOptions({ + model: { ...ANTHROPIC_MODEL, id: "claude-opus-4-20250514", name: "Claude Opus 4" }, + apiKey: "sk-ant-api-test", + extraBetas: [], + stream: true, + interleavedThinking: true, + hasTools: false, + }); + // The date suffix must not parse as minor=20250514 (>= 4.7 display support). + expect(legacy.defaultHeaders["anthropic-beta"]).toContain("interleaved-thinking-2025-05-14"); + + const modern = buildAnthropicClientOptions({ + model: { ...ANTHROPIC_MODEL, id: "claude-opus-4-7", name: "Claude Opus 4.7" }, + apiKey: "sk-ant-api-test", + extraBetas: [], + stream: true, + interleavedThinking: true, + hasTools: false, + }); + expect(modern.defaultHeaders["anthropic-beta"] ?? "").not.toContain("interleaved-thinking-2025-05-14"); + }); + it("adds legacy fine-grained tool-streaming beta only for tool requests on incompatible models", () => { const incompatibleModel: Model<"anthropic-messages"> = { ...ANTHROPIC_MODEL, diff --git a/packages/ai/test/anthropic-client.test.ts b/packages/ai/test/anthropic-client.test.ts index 46abba978..2fa82328a 100644 --- a/packages/ai/test/anthropic-client.test.ts +++ b/packages/ai/test/anthropic-client.test.ts @@ -72,6 +72,25 @@ describe("AnthropicMessagesClient error mapping", () => { expect(error).toBeInstanceOf(AnthropicApiError); expect((error as AnthropicApiError).message).toBe("500 status code (no body)"); }); + + it("does not let fetchOptions override core request fields", async () => { + const { calls, fetch } = createFetchMock([new Response(null, { status: 200 })]); + const preAborted = AbortSignal.abort(); + const client = new AnthropicMessagesClient({ + apiKey: "sk-test", + maxRetries: 0, + fetch, + fetchOptions: { method: "GET", signal: preAborted }, + }); + + const response = await client.messages.create(params).asResponse(); + + // fetchOptions exists for transport extras (tls); a caller-supplied signal + // or method must not disconnect the timeout controller or break the POST. + expect(response.status).toBe(200); + expect(calls[0]?.init.method).toBe("POST"); + expect(calls[0]?.init.signal?.aborted).toBe(false); + }); }); describe("AnthropicMessagesClient retries", () => { diff --git a/packages/ai/test/anthropic-oauth.test.ts b/packages/ai/test/anthropic-oauth.test.ts index 30322edc2..cecae9371 100644 --- a/packages/ai/test/anthropic-oauth.test.ts +++ b/packages/ai/test/anthropic-oauth.test.ts @@ -1,4 +1,5 @@ import { afterEach, describe, expect, it, vi } from "bun:test"; +import { claudeCodeVersion } from "@oh-my-pi/pi-ai/providers/anthropic"; import { AnthropicOAuthFlow, refreshAnthropicToken } from "@oh-my-pi/pi-ai/registry/oauth/anthropic"; import { buildAnthropicAuthConfig, @@ -201,7 +202,7 @@ describe("anthropic oauth alignment", () => { expect(init?.method).toBe("GET"); const headers = init?.headers as Record | undefined; expect(headers?.Authorization).toBe("Bearer access-token"); - expect(headers?.["User-Agent"]).toBe("claude-code/2.1.160"); + expect(headers?.["User-Agent"]).toBe(`claude-code/${claudeCodeVersion}`); expect(headers?.["anthropic-beta"]).toBe("oauth-2025-04-20"); return new Response( JSON.stringify({ diff --git a/packages/ai/test/anthropic-prefill.test.ts b/packages/ai/test/anthropic-prefill.test.ts index f9e8aa470..b64110b4e 100644 --- a/packages/ai/test/anthropic-prefill.test.ts +++ b/packages/ai/test/anthropic-prefill.test.ts @@ -51,6 +51,44 @@ describe("Anthropic assistant-prefill fallback", () => { expect(params.at(-1)?.content).toBe("Continue."); }); + it("repairs consecutive assistant turns left by dropped empty user messages", () => { + const assistant = (text: string): AssistantMessage => ({ + role: "assistant", + content: [{ type: "text", text }], + api: "anthropic-messages", + provider: "anthropic", + model: model.id, + usage: { + input: 0, + output: 0, + cacheRead: 0, + cacheWrite: 0, + totalTokens: 0, + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, + }, + stopReason: "stop", + timestamp: Date.now(), + }); + // An empty nudge submission is dropped by the converter, which would leave + // the two assistant turns adjacent — Anthropic 400s on that shape. + const emptyNudge: UserMessage = { role: "user", content: [{ type: "text", text: "" }], timestamp: Date.now() }; + + const params = convertAnthropicMessages( + [ + { role: "user", content: "answer me", timestamp: Date.now() }, + assistant("partial answer"), + emptyNudge, + assistant("full answer"), + { role: "user", content: "thanks", timestamp: Date.now() }, + ], + model, + false, + ); + + expect(params.map(p => p.role)).toEqual(["user", "assistant", "user", "assistant", "user"]); + expect(params[2]?.content).toBe("Continue."); + }); + it("does not append Continue. when the last turn is already user", () => { const params = convertAnthropicMessages( [ diff --git a/packages/ai/test/anthropic-stream-envelope.test.ts b/packages/ai/test/anthropic-stream-envelope.test.ts index aea22d8e2..dc8c46790 100644 --- a/packages/ai/test/anthropic-stream-envelope.test.ts +++ b/packages/ai/test/anthropic-stream-envelope.test.ts @@ -348,6 +348,62 @@ describe("anthropic stream envelope handling", () => { expect(result.content).toEqual([{ type: "text", text: "hello" }]); }); + it("ignores a spliced second envelope's message_delta after the terminal stop", async () => { + const events: MockAnthropicEvent[] = [ + ...createTextSuccessEvents("hello"), + // Transparent reconnect splices a fresh envelope onto the same stream. + { type: "message_start", message: { id: "msg_second", usage: { input_tokens: 99, output_tokens: 99 } } }, + { type: "message_delta", delta: { stop_reason: "tool_use" }, usage: { input_tokens: 99, output_tokens: 99 } }, + { type: "message_stop" }, + ]; + vi.spyOn(AnthropicMessages.prototype, "create").mockImplementation(() => createMockRequest(events) as never); + + const stream = streamAnthropic(model, context, { apiKey: "sk-ant-test" }); + const collected: AssistantMessageEvent[] = []; + for await (const event of stream) { + collected.push(event); + } + const result = await stream.result(); + + // The completed first envelope owns the stop reason and usage; the splice + // must not relabel a finished turn or overwrite its counters. + expect(countEvents(collected, "error")).toBe(0); + expect(countEvents(collected, "done")).toBe(1); + expect(result.stopReason).toBe("stop"); + expect(result.usage.output).toBe(4); + expect(result.responseId).toBe("msg_text_success"); + expect(result.content).toEqual([{ type: "text", text: "hello" }]); + }); + + it("tolerates envelopes missing usage and delta payloads", async () => { + const events: MockAnthropicEvent[] = [ + { type: "message_start", message: { id: "msg_lenient" } }, + { type: "content_block_start", index: 0, content_block: { type: "text", text: "" } }, + { type: "content_block_delta", index: 0 }, + { type: "content_block_delta", index: 0, delta: { type: "text_delta", text: "hi" } }, + { type: "content_block_stop", index: 0 }, + { type: "message_delta" }, + { type: "message_delta", delta: { stop_reason: "end_turn" } }, + { type: "message_stop" }, + ]; + vi.spyOn(AnthropicMessages.prototype, "create").mockImplementation(() => createMockRequest(events) as never); + + const stream = streamAnthropic(model, context, { apiKey: "sk-ant-test" }); + const collected: AssistantMessageEvent[] = []; + for await (const event of stream) { + collected.push(event); + } + const result = await stream.result(); + + // Proxies that omit usage/delta objects must degrade to anomaly logs, not + // TypeErrors that fail the turn. + expect(countEvents(collected, "error")).toBe(0); + expect(countEvents(collected, "done")).toBe(1); + expect(result.stopReason).toBe("stop"); + expect(result.responseId).toBe("msg_lenient"); + expect(result.content).toEqual([{ type: "text", text: "hi" }]); + }); + it("ignores unknown preamble events before message_start and streams the response once", async () => { let attempt = 0; vi.spyOn(AnthropicMessages.prototype, "create").mockImplementation(() => { diff --git a/packages/ai/test/anthropic-stream-timeout.test.ts b/packages/ai/test/anthropic-stream-timeout.test.ts index 94775eef8..4691c2665 100644 --- a/packages/ai/test/anthropic-stream-timeout.test.ts +++ b/packages/ai/test/anthropic-stream-timeout.test.ts @@ -225,6 +225,54 @@ describe("anthropic first-event timeout retries", () => { expect(result.responseId).toBe("msg_retry_success"); }); + it("keeps the first-event watchdog armed when only pings arrive before message_start", async () => { + vi.useFakeTimers(); + let attempt = 0; + let firstAttemptIteratorStarted = false; + const create = ((_body: unknown, requestOptions?: { signal?: AbortSignal }) => { + attempt += 1; + return createAnthropicMockStream({ + signal: requestOptions?.signal, + events: attempt === 1 ? [{ type: "ping" }] : createSuccessfulAnthropicEvents("retry recovered"), + hangAfterEvents: attempt === 1, + onIteratorStart: + attempt === 1 + ? () => { + firstAttemptIteratorStarted = true; + } + : undefined, + }) as never; + }) as unknown as AnthropicMessagesClientLike["messages"]["create"]; + const client = { messages: { create } } as AnthropicMessagesClientLike; + const providerRetryWait = vi.fn(async () => {}); + + const resultPromise = streamAnthropic(model, context, { + client, + streamFirstEventTimeoutMs: 1, + streamIdleTimeoutMs: 60_000, + providerRetryWait, + }).result(); + + await drainMicrotasksUntil( + () => firstAttemptIteratorStarted, + "Anthropic mock stream did not enter the ping-then-hang first attempt", + ); + await drainMicrotasksUntil(() => vi.getTimerCount() > 0, "Anthropic watchdog timer was not armed"); + + // A keepalive must not consume the first-event watchdog: if it did, the + // stall would be classified as a (non-retryable) 60s idle timeout and + // advancing 1ms would never settle the stream. + vi.advanceTimersByTime(1); + const result = await resolveAfterMicrotasks( + resultPromise, + "Anthropic ping-then-stall did not retry via the first-event watchdog", + ); + + expect(attempt).toBe(2); + expect(result.stopReason).toBe("stop"); + expect(result.content).toEqual([{ type: "text", text: "retry recovered" }]); + }); + it("does not arm the Anthropic first-event watchdog before the stream connects", async () => { let seenRequestTimeout: number | undefined; let seenRequestMaxRetries: number | undefined; diff --git a/packages/ai/test/anthropic-unsigned-thinking-replay.test.ts b/packages/ai/test/anthropic-unsigned-thinking-replay.test.ts index 9b3f8dcbc..d18e77412 100644 --- a/packages/ai/test/anthropic-unsigned-thinking-replay.test.ts +++ b/packages/ai/test/anthropic-unsigned-thinking-replay.test.ts @@ -93,6 +93,47 @@ describe("Anthropic-compatible unsigned thinking replay (#2005)", () => { expect(blocks[1]).toEqual({ type: "text", text: "Sure." }); }); + it("sanitizes lone surrogates in cross-API tool arguments only", () => { + const loneSurrogate = "broken \ud83d end"; + const makeToolCallAssistant = (api: AssistantMessage["api"]): AssistantMessage => ({ + role: "assistant", + content: [ + { + type: "toolCall", + id: "call_1", + name: "write", + arguments: { text: loneSurrogate, nested: { parts: [loneSurrogate] } }, + }, + ], + api, + provider: api === "anthropic-messages" ? "custom-anthropic" : "openai", + model: "reasoning-model", + usage: { + input: 0, + output: 0, + cacheRead: 0, + cacheWrite: 0, + totalTokens: 0, + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, + }, + stopReason: "toolUse", + timestamp: 0, + }); + + // Cross-API replay: Anthropic's strict UTF-8 validation rejects lone + // surrogates, so string leaves are deep-sanitized. + const crossBlocks = assistantWireBlocks([makeUser(), makeToolCallAssistant("openai-responses")], makeModel()); + const crossToolUse = crossBlocks.find(block => block.type === "tool_use") as WireToolUseBlock; + expect(crossToolUse.input.text).toBe("broken \ufffd end"); + expect((crossToolUse.input.nested as { parts: string[] }).parts[0]).toBe("broken \ufffd end"); + + // Same-API replay stays byte-identical (the args came from Anthropic's own + // JSON; rewriting them would destabilize prompt-cache prefixes). + const sameBlocks = assistantWireBlocks([makeUser(), makeToolCallAssistant("anthropic-messages")], makeModel()); + const sameToolUse = sameBlocks.find(block => block.type === "tool_use") as WireToolUseBlock; + expect(sameToolUse.input.text).toBe(loneSurrogate); + }); + it("covers the Xiaomi MiMo Anthropic-compatible reporter configuration without provider allowlists", () => { const model = makeModel({ provider: "user-custom", diff --git a/packages/ai/test/claude-usage-headers.test.ts b/packages/ai/test/claude-usage-headers.test.ts index 3aab042ad..e7c1cf096 100644 --- a/packages/ai/test/claude-usage-headers.test.ts +++ b/packages/ai/test/claude-usage-headers.test.ts @@ -1,4 +1,5 @@ import { describe, expect, it } from "bun:test"; +import { claudeCodeVersion } from "@oh-my-pi/pi-ai/providers/anthropic"; import type { UsageFetchContext } from "@oh-my-pi/pi-ai/usage"; import { claudeUsageProvider } from "@oh-my-pi/pi-ai/usage/claude"; @@ -75,7 +76,7 @@ describe("claude usage request headers", () => { const headers = calls[0]?.init?.headers; expect(getHeaderCaseInsensitive(headers, "authorization")).toBe(`Bearer ${token}`); - expect(getHeaderCaseInsensitive(headers, "user-agent")).toBe("claude-cli/2.1.160 (external, cli)"); + expect(getHeaderCaseInsensitive(headers, "user-agent")).toBe(`claude-cli/${claudeCodeVersion} (external, cli)`); const beta = getHeaderCaseInsensitive(headers, "anthropic-beta"); expect(beta).toBeDefined(); diff --git a/packages/ai/test/github-copilot-anthropic-auth.test.ts b/packages/ai/test/github-copilot-anthropic-auth.test.ts index 7566c6cfb..f22839f4a 100644 --- a/packages/ai/test/github-copilot-anthropic-auth.test.ts +++ b/packages/ai/test/github-copilot-anthropic-auth.test.ts @@ -224,6 +224,38 @@ describe("Anthropic Copilot auth config", () => { expect(result.baseURL).toBe("http://127.0.0.1:8317"); }); + it("sends Content-Type and anthropic-version on Copilot anthropic requests", () => { + const result = buildAnthropicClientOptions({ + model: makeCopilotClaudeModel(), + apiKey: "ghu_test", + extraBetas: [], + stream: true, + dynamicHeaders: {}, + }); + + // The client posts JSON.stringify(params); without these the request goes + // out with no Content-Type at all (Bun does not default it for string + // bodies when a plain headers object is supplied). + expect(result.defaultHeaders["Content-Type"]).toBe("application/json"); + expect(result.defaultHeaders["anthropic-version"]).toBe("2023-06-01"); + }); + + it("merges Copilot headers case-insensitively so auth headers cannot duplicate", () => { + const result = buildAnthropicClientOptions({ + model: { ...makeCopilotClaudeModel(), headers: { ...OPENCODE_HEADERS, authorization: "Bearer override" } }, + apiKey: "ghu_test", + extraBetas: [], + stream: true, + dynamicHeaders: {}, + }); + + // A miscased duplicate would survive Object.assign and the Headers + // constructor then joins both values comma-separated on the wire. + const authKeys = Object.keys(result.defaultHeaders).filter(key => key.toLowerCase() === "authorization"); + expect(authKeys).toHaveLength(1); + expect(result.defaultHeaders[authKeys[0]]).toBe("Bearer override"); + }); + it("builds anthropic auth URLs from the normalized service root", () => { const url = buildAnthropicUrl({ apiKey: "test-key",