From dc198dff1fa5d3a8a232d2d56b3525ab2b760345 Mon Sep 17 00:00:00 2001 From: DavidArthurCole Date: Sat, 26 Sep 2026 11:07:46 +0700 Subject: [PATCH] feat(claude): merge client anthropic-beta flags and forward rate-limit headers --- open-sse/executors/default.js | 7 +- open-sse/handlers/chatCore.js | 3 +- .../handlers/chatCore/nonStreamingHandler.js | 3 +- .../handlers/chatCore/streamingHandler.js | 3 +- open-sse/providers/shared.js | 5 ++ open-sse/utils/error.js | 15 ++-- open-sse/utils/upstreamHeaders.js | 12 +++ src/sse/handlers/chat.js | 7 +- .../anthropic-gateway-passthrough.test.js | 80 +++++++++++++++++++ 9 files changed, 122 insertions(+), 13 deletions(-) create mode 100644 open-sse/utils/upstreamHeaders.js create mode 100644 tests/unit/anthropic-gateway-passthrough.test.js diff --git a/open-sse/executors/default.js b/open-sse/executors/default.js index 64ad46d0..eadf1797 100644 --- a/open-sse/executors/default.js +++ b/open-sse/executors/default.js @@ -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 diff --git a/open-sse/handlers/chatCore.js b/open-sse/handlers/chatCore.js index 0db37ba1..a6230fc4 100644 --- a/open-sse/handlers/chatCore.js +++ b/open-sse/handlers/chatCore.js @@ -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 }; diff --git a/open-sse/handlers/chatCore/nonStreamingHandler.js b/open-sse/handlers/chatCore/nonStreamingHandler.js index 19018147..514a24b2 100644 --- a/open-sse/handlers/chatCore/nonStreamingHandler.js +++ b/open-sse/handlers/chatCore/nonStreamingHandler.js @@ -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) } }) }; } diff --git a/open-sse/handlers/chatCore/streamingHandler.js b/open-sse/handlers/chatCore/streamingHandler.js index 6751256b..55106fd5 100644 --- a/open-sse/handlers/chatCore/streamingHandler.js +++ b/open-sse/handlers/chatCore/streamingHandler.js @@ -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) } }) }; } diff --git a/open-sse/providers/shared.js b/open-sse/providers/shared.js index 176c82e9..def3171d 100644 --- a/open-sse/providers/shared.js +++ b/open-sse/providers/shared.js @@ -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"; diff --git a/open-sse/utils/error.js b/open-sse/utils/error.js index 315723e3..6992d4a4 100644 --- a/open-sse/utils/error.js +++ b/open-sse/utils/error.js @@ -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) } } ); diff --git a/open-sse/utils/upstreamHeaders.js b/open-sse/utils/upstreamHeaders.js new file mode 100644 index 00000000..49ec3188 --- /dev/null +++ b/open-sse/utils/upstreamHeaders.js @@ -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; +} diff --git a/src/sse/handlers/chat.js b/src/sse/handlers/chat.js index 7aa530d3..172549ea 100644 --- a/src/sse/handlers/chat.js +++ b/src/sse/handlers/chat.js @@ -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; } diff --git a/tests/unit/anthropic-gateway-passthrough.test.js b/tests/unit/anthropic-gateway-passthrough.test.js new file mode 100644 index 00000000..308f6636 --- /dev/null +++ b/tests/unit/anthropic-gateway-passthrough.test.js @@ -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"); + }); +});