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.
This commit is contained in:
can1357
2026-06-09 23:17:26 +02:00
parent 228c1ddd8d
commit f80d45aa6b
16 changed files with 721 additions and 99 deletions
+17
View File
@@ -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
@@ -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();
+314 -91
View File
@@ -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<s
Accept: acceptHeader,
Authorization: `Bearer ${options.apiKey}`,
...sharedHeaders,
"anthropic-beta": betaHeader,
...(betaHeader ? { "anthropic-beta": betaHeader } : {}),
...(options.claudeCodeSessionId ? { "X-Claude-Code-Session-Id": options.claudeCodeSessionId } : {}),
"x-client-request-id": nodeCrypto.randomUUID(),
"User-Agent": userAgent,
@@ -302,7 +303,9 @@ let warnedStopSequencesTrim = false;
*/
function supportsAdaptiveThinkingDisplay(modelId: string): boolean {
if (/claude-(?:fable|mythos)-5\b/.test(modelId)) return true;
const match = /claude-opus-(\d+)-(\d+)/.exec(modelId);
// Bound the minor to non-date digits: bare dated ids like
// `claude-opus-4-20250514` (Opus 4.0) must not parse as minor=20250514.
const match = /claude-opus-(\d+)-(\d{1,2})(?!\d)/.exec(modelId);
if (!match) return false;
const major = Number(match[1]);
const minor = Number(match[2]);
@@ -551,7 +554,7 @@ const CCH_PLACEHOLDER = cchEncoder.encode(CCH_PLACEHOLDER_STR);
const BILLING_SYSTEM_MARKER = cchEncoder.encode(`"system":[{"type":"text","text":"${CLAUDE_BILLING_HEADER_PREFIX}`);
const CCH_BILLING_SEARCH_WINDOW = 150;
function patchCch(body: Uint8Array): boolean {
function patchCch(body: Uint8Array): "patched" | "no-billing-header" | "unanchored" {
// Zero-copy Buffer view over the same memory; its `indexOf` is a native memmem,
// ~7.5x faster than a hand-rolled byte loop here — the marker sits ~99% through
// the body because `messages` serializes before `system`, so a JS scan would
@@ -560,19 +563,19 @@ function patchCch(body: Uint8Array): boolean {
// Find the combined system[0] + billing-header prefix marker.
const markerIdx = view.indexOf(BILLING_SYSTEM_MARKER);
if (markerIdx === -1) return false; // no CC billing header injected
if (markerIdx === -1) return "no-billing-header"; // no CC billing header injected
// Placeholder must sit within CCH_BILLING_SEARCH_WINDOW bytes after the marker.
const searchFrom = markerIdx + BILLING_SYSTEM_MARKER.length;
const idx = view.indexOf(CCH_PLACEHOLDER, searchFrom);
if (idx === -1 || idx - searchFrom > 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 = <R>(fn: () => Promise<R>) => Promise<R>;
/**
* 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<void>();
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<ImageContent> {
try {
const inputBuffer = Buffer.from(block.data, "base64");
@@ -775,12 +810,13 @@ async function resizeAnthropicManyImageBlock(block: ImageContent): Promise<Image
async function resizeAnthropicManyImageContent(
content: (TextContent | ImageContent)[],
state: { resized: number },
limit: ResizeLimiter,
): Promise<(TextContent | ImageContent)[]> {
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<Message> {
async function resizeAnthropicManyImageMessage(
message: Message,
state: { resized: number },
limit: ResizeLimiter,
): Promise<Message> {
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<string, FoundryTlsOptions | undefined>();
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<string, string> | undefined)[]): Record<string, string> {
// 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<string, string> = {};
const keyByLower = new Map<string, string>();
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<string, unknown> = {};
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<object>()): 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<string, unknown>,
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;
+2 -1
View File
@@ -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.
+2 -1
View File
@@ -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;
+1 -1
View File
@@ -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 {
+4
View File
@@ -45,6 +45,10 @@ export async function callWithCopilotModelRetry<T>(
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 });
+139 -1
View File
@@ -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<string, string> | 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,
+19
View File
@@ -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", () => {
+2 -1
View File
@@ -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<string, string> | 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({
@@ -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(
[
@@ -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(() => {
@@ -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;
@@ -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",
@@ -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();
@@ -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",