feat(claude): merge client anthropic-beta flags and forward rate-limit headers

This commit is contained in:
DavidArthurCole
2026-09-26 11:07:46 +07:00
parent 4a57df8bf9
commit dc198dff1f
9 changed files with 122 additions and 13 deletions

View File

@@ -1,6 +1,6 @@
import { BaseExecutor } from "./base.js";
import { PROVIDERS, PROVIDER_OAUTH } from "../config/providers.js";
import { ANTHROPIC_API_VERSION, OPENAI_COMPAT_BASE, ANTHROPIC_COMPAT_BASE, selectAnthropicBeta } from "../providers/shared.js";
import { ANTHROPIC_API_VERSION, OPENAI_COMPAT_BASE, ANTHROPIC_COMPAT_BASE, selectAnthropicBeta, mergeAnthropicBeta } from "../providers/shared.js";
import { resolveOpenAICompatibleApiType } from "../services/provider.js";
import { OAUTH_ENDPOINTS, buildKimiHeaders } from "../config/appConstants.js";
import { buildClineHeaders } from "../shared/clineAuth.js";
@@ -164,9 +164,12 @@ export class DefaultExecutor extends BaseExecutor {
// a node fronting Kimi or GLM answers on its own ids and never matches, so
// gateways that would choke on unknown beta flags are left untouched.
const isClaudeModel = typeof model === "string" && /^claude-/.test(model);
const clientBeta = credentials?.rawHeaders?.["anthropic-beta"];
if (model && (this.provider === "claude"
|| (this.provider?.startsWith?.("anthropic-compatible-") && isClaudeModel))) {
headers["Anthropic-Beta"] = selectAnthropicBeta(model, body);
headers["Anthropic-Beta"] = mergeAnthropicBeta(selectAnthropicBeta(model, body), clientBeta);
} else if (this.provider === "anthropic" && clientBeta) {
headers["Anthropic-Beta"] = mergeAnthropicBeta(headers["Anthropic-Beta"], clientBeta);
}
// Strip first-party Claude Code identity headers for non-Anthropic anthropic-compatible upstreams

View File

@@ -9,6 +9,7 @@ import { createRequestLogger } from "../utils/requestLogger.js";
import { getModelTargetFormat, getModelSupportedFormats, getModelStrip, getModelUpstreamId, getModelType, PROVIDER_ID_TO_ALIAS } from "../config/providerModels.js";
import { PROVIDERS } from "../config/providers.js";
import { createErrorResult, parseUpstreamError, formatProviderError } from "../utils/error.js";
import { upstreamResponseHeaders } from "../utils/upstreamHeaders.js";
import { HTTP_STATUS, TOKEN_SAVER_HEADER } from "../config/runtimeConfig.js";
import { handleBypassRequest } from "../utils/bypassHandler.js";
import { trackPendingRequest, appendRequestLog, saveRequestDetail } from "@/lib/usageDb.js";
@@ -486,7 +487,7 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
log.errorLine(reqTag, "✗", `ERROR ${statusCode} · ${provider}/${model} · ${Date.now() - requestStartTime}ms${urlStr}\n ${errMsg}`);
}
reqLogger.logError(new Error(message), finalBody || translatedBody);
return createErrorResult(statusCode, errMsg, resetsAtMs);
return createErrorResult(statusCode, errMsg, resetsAtMs, upstreamResponseHeaders(providerResponse.headers));
}
const sharedCtx = { provider, model, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess, pxpipe: pxpipeSummary, reqTag, log };

View File

@@ -4,6 +4,7 @@ import { fromOpenAIFinish } from "../../translator/concerns/finishReason.js";
import { ollamaBodyToOpenAI } from "../../translator/response/ollama-to-openai.js";
import { addBufferToUsage, filterUsageForFormat } from "../../utils/usageTracking.js";
import { createErrorResult } from "../../utils/error.js";
import { upstreamResponseHeaders } from "../../utils/upstreamHeaders.js";
import { HTTP_STATUS } from "../../config/runtimeConfig.js";
import { parseSSEToOpenAIResponse } from "./sseToJsonHandler.js";
import { unwrapClineEnvelope } from "../../shared/clineEnvelope.js";
@@ -399,7 +400,7 @@ export async function handleNonStreamingResponse({ providerResponse, provider, m
return {
success: true,
response: new Response(JSON.stringify(restoreToolNames(translatedResponse, toolNameMap)), {
headers: { "Content-Type": "application/json", "Access-Control-Allow-Origin": "*" }
headers: { "Content-Type": "application/json", "Access-Control-Allow-Origin": "*", ...upstreamResponseHeaders(providerResponse.headers) }
})
};
}

View File

@@ -9,6 +9,7 @@ import { buildStreamErrorBytes } from "../../utils/streamHelpers.js";
import { buildRequestDetail, extractRequestConfig, saveUsageStats, formatDoneLine } from "./requestDetail.js";
import { saveRequestDetail } from "@/lib/usageDb.js";
import { SSE_HEADERS_CORS as SSE_HEADERS } from "../../utils/sseConstants.js";
import { upstreamResponseHeaders } from "../../utils/upstreamHeaders.js";
// Codex returns Responses API SSE → which client format to translate INTO, by request sourceFormat.
// Gemini-family all map to ANTIGRAVITY decoder; unknown sources fall back to OPENAI.
@@ -109,7 +110,7 @@ export async function handleStreamingResponse({ providerResponse, provider, mode
return {
success: true,
response: new Response(transformedBody, { headers: SSE_HEADERS })
response: new Response(transformedBody, { headers: { ...SSE_HEADERS, ...upstreamResponseHeaders(providerResponse.headers) } })
};
}

View File

@@ -77,6 +77,11 @@ export function selectAnthropicBeta(model = "", body = null) {
return flags.join(",");
}
export function mergeAnthropicBeta(...values) {
const flags = values.flatMap((v) => (typeof v === "string" ? v.split(",") : [])).map((f) => f.trim()).filter(Boolean);
return [...new Set(flags)].join(",");
}
// Shared baseUrls
export const KIMI_CODING_BASE_URL = "https://api.kimi.com/coding/v1/messages";

View File

@@ -27,12 +27,13 @@ export function buildErrorBody(statusCode, message) {
* @param {string} message - Error message
* @returns {Response} HTTP Response object
*/
export function errorResponse(statusCode, message) {
export function errorResponse(statusCode, message, extraHeaders = null) {
return new Response(JSON.stringify(buildErrorBody(statusCode, message)), {
status: statusCode,
headers: {
"Content-Type": "application/json",
"Access-Control-Allow-Origin": "*"
"Access-Control-Allow-Origin": "*",
...extraHeaders
}
});
}
@@ -95,13 +96,13 @@ export async function parseUpstreamError(response, executor = null) {
* @param {number} [resetsAtMs] - Optional precise cooldown expiry (ms epoch) for provider-specific quota errors
* @returns {{ success: false, status: number, error: string, response: Response, resetsAtMs?: number }}
*/
export function createErrorResult(statusCode, message, resetsAtMs) {
export function createErrorResult(statusCode, message, resetsAtMs, extraHeaders = null) {
return {
success: false,
status: statusCode,
error: message,
resetsAtMs,
response: errorResponse(statusCode, message)
response: errorResponse(statusCode, message, extraHeaders)
};
}
@@ -113,7 +114,7 @@ export function createErrorResult(statusCode, message, resetsAtMs) {
* @param {string} retryAfterHuman - Human-readable retry info e.g. "reset after 30s"
* @returns {Response}
*/
export function unavailableResponse(statusCode, message, retryAfter, retryAfterHuman) {
export function unavailableResponse(statusCode, message, retryAfter, retryAfterHuman, extraHeaders = null) {
const retryAfterSec = Math.max(Math.ceil((new Date(retryAfter).getTime() - Date.now()) / 1000), 1);
const msg = `${message} (${retryAfterHuman})`;
return new Response(
@@ -121,8 +122,10 @@ export function unavailableResponse(statusCode, message, retryAfter, retryAfterH
{
status: statusCode,
headers: {
...extraHeaders,
"Content-Type": "application/json",
"Retry-After": String(retryAfterSec)
// Intentionally mis-cased to prevent duplicate headers
"retry-after": String(retryAfterSec)
}
}
);

View File

@@ -0,0 +1,12 @@
const FORWARDED = new Set(["retry-after", "x-should-retry"]);
const FORWARDED_PREFIX = "anthropic-ratelimit-";
export function upstreamResponseHeaders(headers) {
const out = {};
if (typeof headers?.forEach !== "function") return out;
headers.forEach((value, name) => {
const key = name.toLowerCase();
if (FORWARDED.has(key) || key.startsWith(FORWARDED_PREFIX)) out[key] = value;
});
return out;
}

View File

@@ -15,6 +15,7 @@ import { DEFAULT_HEADROOM_URL } from "@/lib/headroom/detect";
import { getTransform as getPxpipeTransform } from "@/lib/pxpipe/loader.js";
import { appendPxpipeEvent } from "@/lib/pxpipe/events.js";
import { errorResponse, unavailableResponse } from "open-sse/utils/error.js";
import { upstreamResponseHeaders } from "open-sse/utils/upstreamHeaders.js";
import { handleComboChat, handleFusionChat, detectRequiredCapabilities } from "open-sse/services/combo.js";
import { augmentModelsWithCapacityAdapter, withCapacityAdapterStripping, getActiveAdapterStrategy } from "open-sse/services/capacityAdapter.js";
import { handleBypassRequest } from "open-sse/utils/bypassHandler.js";
@@ -229,6 +230,7 @@ async function handleSingleModelChat(body, modelStr, clientRawRequest = null, re
const excludeConnectionIds = new Set();
let lastError = null;
let lastStatus = null;
let lastHeaders = null;
while (true) {
const credentials = await getProviderCredentials(provider, excludeConnectionIds, model);
@@ -239,14 +241,14 @@ async function handleSingleModelChat(body, modelStr, clientRawRequest = null, re
const errorMsg = lastError || credentials.lastError || "Unavailable";
const status = HTTP_STATUS.SERVICE_UNAVAILABLE;
log.warn("CHAT", `[${provider}/${model}] ${errorMsg} (${credentials.retryAfterHuman})`);
return unavailableResponse(status, `[${provider}/${model}] ${errorMsg}`, credentials.retryAfter, credentials.retryAfterHuman);
return unavailableResponse(status, `[${provider}/${model}] ${errorMsg}`, credentials.retryAfter, credentials.retryAfterHuman, lastHeaders);
}
if (excludeConnectionIds.size === 0) {
log.warn("AUTH", `No active credentials for provider: ${provider}`);
return errorResponse(HTTP_STATUS.NOT_FOUND, `No active credentials for provider: ${provider}`);
}
log.warn("CHAT", "No more accounts available", { provider });
return errorResponse(lastStatus || HTTP_STATUS.SERVICE_UNAVAILABLE, lastError || "All accounts unavailable");
return errorResponse(lastStatus || HTTP_STATUS.SERVICE_UNAVAILABLE, lastError || "All accounts unavailable", lastHeaders);
}
// Account selection shown in the unified "▶" line (acc:...)
@@ -331,6 +333,7 @@ async function handleSingleModelChat(body, modelStr, clientRawRequest = null, re
excludeConnectionIds.add(credentials.connectionId);
lastError = result.error;
lastStatus = result.status;
lastHeaders = upstreamResponseHeaders(result.response?.headers);
continue;
}

View File

@@ -0,0 +1,80 @@
import { describe, it, expect, beforeEach, vi } from "vitest";
import { mergeAnthropicBeta } from "open-sse/providers/shared.js";
import { upstreamResponseHeaders } from "open-sse/utils/upstreamHeaders.js";
import { createErrorResult, unavailableResponse } from "open-sse/utils/error.js";
const betaFlags = (headers) => (headers["Anthropic-Beta"] || "").split(",").map((s) => s.trim()).filter(Boolean);
describe("mergeAnthropicBeta", () => {
it("unions and dedupes comma lists, ignoring blanks", () => {
expect(mergeAnthropicBeta("a,b", " b , c ,", undefined, "")).toBe("a,b,c");
});
});
describe("DefaultExecutor.buildHeaders() forwards client anthropic-beta", () => {
let DefaultExecutor;
beforeEach(async () => {
vi.resetModules();
({ DefaultExecutor } = await import("open-sse/executors/default.js"));
});
it("keeps unknown client flags alongside the pinned set on claude", () => {
const executor = new DefaultExecutor("claude");
const rawHeaders = { "anthropic-beta": "safeguards-2026-09-01,context-1m-2025-08-07" };
const flags = betaFlags(executor.buildHeaders({ apiKey: "k", rawHeaders }, true, undefined, "claude-opus-5"));
expect(flags).toContain("safeguards-2026-09-01");
expect(flags).toContain("context-1m-2025-08-07");
expect(flags).toContain("context-management-2025-06-27");
expect(new Set(flags).size).toBe(flags.length);
});
it("forwards client flags on anthropic-compatible Claude models", () => {
const executor = new DefaultExecutor("anthropic-compatible-custom");
const creds = { apiKey: "k", rawHeaders: { "anthropic-beta": "safeguards-2026-09-01" }, providerSpecificData: { baseUrl: "https://gw.example.com/v1" } };
const flags = betaFlags(executor.buildHeaders(creds, true, undefined, "claude-sonnet-5"));
expect(flags).toContain("safeguards-2026-09-01");
expect(flags).not.toContain("claude-code-20250219");
});
it("forwards client flags on the anthropic provider", () => {
const executor = new DefaultExecutor("anthropic");
const flags = betaFlags(executor.buildHeaders({ apiKey: "k", rawHeaders: { "anthropic-beta": "safeguards-2026-09-01" } }, true, undefined, "claude-sonnet-5"));
expect(flags).toContain("safeguards-2026-09-01");
});
});
describe("upstream response header forwarding", () => {
const upstream = new Headers({
"retry-after": "12",
"x-should-retry": "false",
"anthropic-ratelimit-unified-status": "rejected",
"anthropic-ratelimit-unified-reset": "1790000000",
"set-cookie": "secret=1",
"content-length": "99",
});
it("picks only retry and ratelimit headers", () => {
expect(upstreamResponseHeaders(upstream)).toEqual({
"retry-after": "12",
"x-should-retry": "false",
"anthropic-ratelimit-unified-status": "rejected",
"anthropic-ratelimit-unified-reset": "1790000000",
});
expect(upstreamResponseHeaders(undefined)).toEqual({});
});
it("attaches them to error results", () => {
const { response } = createErrorResult(429, "limited", undefined, upstreamResponseHeaders(upstream));
expect(response.headers.get("x-should-retry")).toBe("false");
expect(response.headers.get("anthropic-ratelimit-unified-status")).toBe("rejected");
expect(response.headers.get("set-cookie")).toBeNull();
});
it("keeps the gateway retry-after on all-accounts-limited responses", () => {
const retryAt = new Date(Date.now() + 30000).toISOString();
const res = unavailableResponse(503, "busy", retryAt, "30s", upstreamResponseHeaders(upstream));
expect(Number(res.headers.get("retry-after"))).toBeGreaterThan(20);
expect(res.headers.get("anthropic-ratelimit-unified-reset")).toBe("1790000000");
});
});