From f6fc637923594c969bf76f16920e91d48fe9356f Mon Sep 17 00:00:00 2001 From: Mako-L <9460018+Mako-L@users.noreply.github.com> Date: Thu, 1 Oct 2026 03:57:51 +0300 Subject: [PATCH 1/3] feat: route accounts by service tier --- README_EN.md | 49 ++++++++++++ config/default.yaml | 4 + src/auth/account-lifecycle.ts | 13 +++- src/auth/account-pool.ts | 2 +- src/auth/service-tier-routing.ts | 16 ++++ src/config-schema.ts | 10 +++ src/routes/responses-compact.ts | 7 +- src/routes/shared/account-acquisition.ts | 3 +- src/routes/shared/non-streaming-helpers.ts | 2 +- .../shared/proxy-error-retry-transition.ts | 3 + .../shared/proxy-fallback-account-retry.ts | 4 +- src/routes/shared/proxy-handler.ts | 17 +++- tests/integration/proxy-handler.test.ts | 56 +++++++++++++ tests/unit/auth/plan-routing-acquire.test.ts | 78 +++++++++++++++++++ tests/unit/config-schema.test.ts | 17 ++++ 15 files changed, 270 insertions(+), 11 deletions(-) create mode 100644 src/auth/service-tier-routing.ts diff --git a/README_EN.md b/README_EN.md index 8abe7bc5e..ac19c260f 100644 --- a/README_EN.md +++ b/README_EN.md @@ -705,6 +705,55 @@ When `quota.skip_exhausted: true`, the account pool skips accounts whose cached The skip condition is currently `rate_limit.limit_reached === true`, `secondary_rate_limit.limit_reached === true`, or `code_review_rate_limit.limit_reached === true` in cached quota. If `used_percent` is merely near 100, for example 99%, but upstream has not set `limit_reached`, the proxy may still use that account. Once upstream returns 429, the account is marked `rate_limited`, enters backoff, and the request is retried with another available account. Secondary and code-review windows are removed from cache after their own `reset_at` passes, so an account is not skipped forever on stale quota data. +### Account routing by service tier + +Use `auth.service_tier_routing` to reserve accounts for a service tier. The keys +match the effective `service_tier` exactly; the proxy does not translate marketing +names into tier IDs or infer entitlements. Confirm the upstream tier and account +eligibility before configuring a rule. + +```yaml +model: + service_tier_overrides: + gpt-6.1-sol: ultrafast +auth: + service_tier_routing: + ultrafast: + account_ids: ["reserved-account-entry-id"] + default: + exclude_account_ids: ["reserved-account-entry-id"] + standard: + exclude_account_ids: ["reserved-account-entry-id"] +``` + +This example forces all OAuth GPT-6.1-Sol inference requests to `ultrafast`, +even when a client requests another tier. It reserves one account for ultrafast and +keeps requests using `default` or `standard` on the other accounts. The override +uses the resolved model ID and does not change the model itself. Without a model +override, client tier selection is preserved. Normal Astra requests should omit +the tier unless the upstream explicitly supports the requested tier value. `account_ids` and `exclude_account_ids` use the proxy account entry +`id` shown by `/auth/accounts`, not the upstream ChatGPT account ID or email. + +Alternatively, use `plan_types: ["pro"]` when the reported plan type uniquely +identifies eligible accounts. Account plan metadata may not distinguish Pro +subscription variants; use explicit entry IDs in that case. If multiple fields +are supplied, all restrictions must match, and exclusions take precedence. +Unknown plan types cannot satisfy a `plan_types` restriction. + +Rules are applied after model eligibility, quota and concurrency checks, before +plan priority, session affinity and rotation. Account retries retain the tier. +If all matching accounts are unavailable, the request fails instead of selecting +an excluded account or using the API-key fallback. Compact requests also apply +these account restrictions, without adding a service tier to the compact payload. +Explicitly routed third-party API providers do not use the OAuth account pool and +are outside these rules. + +When no request tier is specified, `model.default_service_tier` is used, falling +back to the key `default`. Unconfigured tiers retain existing selection behavior. +Rules are empty by default. Empty lists and empty rules are rejected to catch +configuration errors. This is account selection policy, not an upstream entitlement +check or a guarantee of a particular response speed. + ### Ollama Bridge Configuration ```yaml diff --git a/config/default.yaml b/config/default.yaml index 13cab7612..99e9da124 100644 --- a/config/default.yaml +++ b/config/default.yaml @@ -18,6 +18,8 @@ model: image_host_model: gpt-5.5 default_reasoning_effort: null default_service_tier: null + # Optional forced service tier per OAuth model, overriding client preferences. + service_tier_overrides: {} default_tools: [] aliases: {} custom_models: [] @@ -37,6 +39,8 @@ auth: refresh_enabled: true refresh_margin_seconds: 300 rotation_strategy: least_used + # Hard account restrictions keyed by effective service_tier; see README_EN.md. + service_tier_routing: {} rate_limit_backoff_seconds: 60 oauth_client_id: app_EMoamEEZ73f0CkXaXp7hrann oauth_auth_endpoint: https://auth.openai.com/oauth/authorize diff --git a/src/auth/account-lifecycle.ts b/src/auth/account-lifecycle.ts index 0ecd7b91d..99fc8209b 100644 --- a/src/auth/account-lifecycle.ts +++ b/src/auth/account-lifecycle.ts @@ -5,6 +5,7 @@ * Uses AccountRegistry for entry access (no circular dep — one-way reference). */ +import { getServiceTierAccountRule } from "./service-tier-routing.js"; import { getConfig } from "../config.js"; import { getModelPlanTypes, isPlanFetched } from "../models/model-store.js"; import { hasReachedCachedQuota } from "./quota-skip.js"; @@ -72,7 +73,7 @@ export class AccountLifecycle { } } - acquire(options?: { model?: string; excludeIds?: string[]; preferredEntryId?: string }): AcquiredAccount | null { + acquire(options?: { model?: string; serviceTier?: string | null; excludeIds?: string[]; preferredEntryId?: string }): AcquiredAccount | null { const nowMs = Date.now(); const now = new Date(nowMs); @@ -117,6 +118,16 @@ export class AccountLifecycle { } } + const rule = getServiceTierAccountRule(options?.serviceTier); + if (rule) { + candidates = candidates.filter((account) => + (!rule.plan_types || (account.planType != null && rule.plan_types.includes(account.planType))) && + (!rule.account_ids || rule.account_ids.includes(account.id)) && + !rule.exclude_account_ids?.includes(account.id), + ); + if (candidates.length === 0) return null; + } + // Tier-based filtering: when configured, restrict to the highest available tier const tierPriority = config.auth.tier_priority; if (tierPriority && tierPriority.length > 0) { diff --git a/src/auth/account-pool.ts b/src/auth/account-pool.ts index c8b896ece..ccf22b3dc 100644 --- a/src/auth/account-pool.ts +++ b/src/auth/account-pool.ts @@ -89,7 +89,7 @@ export class AccountPool { // ── Lifecycle (acquire/release) ─────────────────────────────────── - acquire(options?: { model?: string; excludeIds?: string[]; preferredEntryId?: string }): AcquiredAccount | null { + acquire(options?: { model?: string; serviceTier?: string | null; excludeIds?: string[]; preferredEntryId?: string }): AcquiredAccount | null { return this.lifecycle.acquire(options); } diff --git a/src/auth/service-tier-routing.ts b/src/auth/service-tier-routing.ts new file mode 100644 index 000000000..d47c23ede --- /dev/null +++ b/src/auth/service-tier-routing.ts @@ -0,0 +1,16 @@ +import { getConfig } from "../config.js"; + +/** Missing tiers use the configured default, then the upstream standard tier. */ +export function getServiceTierAccountRule(serviceTier?: string | null) { + const config = getConfig(); + if (!config.auth.service_tier_routing) return undefined; + const tier = serviceTier ?? config.model.default_service_tier ?? "default"; + const rules = config.auth.service_tier_routing; + return Object.hasOwn(rules, tier) ? rules[tier] : undefined; +} + +/** Exact model IDs; only explicit operator overrides supersede the client. */ +export function getModelServiceTierOverride(model: string): string | undefined { + const overrides = getConfig().model?.service_tier_overrides; + return overrides && Object.hasOwn(overrides, model) ? overrides[model] : undefined; +} diff --git a/src/config-schema.ts b/src/config-schema.ts index 4953403bd..57cfab8c2 100644 --- a/src/config-schema.ts +++ b/src/config-schema.ts @@ -92,6 +92,8 @@ export const ConfigSchema = z.object({ .default("gpt-5.5"), default_reasoning_effort: z.string().nullable().default(null), default_service_tier: z.string().nullable().default(null), + /** Forced service tiers for OAuth models, overriding client tier preferences. */ + service_tier_overrides: z.record(z.string().trim().min(1), z.string().trim().min(1)).default({}), default_tools: z.array(z.string().trim().min(1)).default([]), aliases: z.record(z.string(), z.string()).default({}), custom_models: z.array(CustomModelSchema).default([]), @@ -122,6 +124,14 @@ export const ConfigSchema = z.object({ rotation_strategy: z.enum(ROTATION_STRATEGIES).default("least_used"), /** Preferred plan-type ordering for account selection (e.g. ["plus","team","free"]). */ tier_priority: z.array(z.string()).nullable().default(null), + /** Hard account restrictions keyed by the effective service tier. */ + service_tier_routing: z.record(z.string().trim().min(1), z.object({ + plan_types: z.array(z.string().trim().min(1)).min(1).optional(), + account_ids: z.array(z.string().trim().min(1)).min(1).optional(), + exclude_account_ids: z.array(z.string().trim().min(1)).min(1).optional(), + }).strict().refine((rule) => rule.plan_types || rule.account_ids || rule.exclude_account_ids, { + message: "A service-tier rule must contain at least one account restriction", + })).default({}), rate_limit_backoff_seconds: z.number().min(1).default(60), oauth_client_id: z.string().default("app_EMoamEEZ73f0CkXaXp7hrann"), oauth_auth_endpoint: z.string().default("https://auth.openai.com/oauth/authorize"), diff --git a/src/routes/responses-compact.ts b/src/routes/responses-compact.ts index 6909671f1..bb2784eca 100644 --- a/src/routes/responses-compact.ts +++ b/src/routes/responses-compact.ts @@ -2,6 +2,7 @@ * Responses API compact handler — non-streaming JSON proxy for /v1/responses/compact. */ +import { getModelServiceTierOverride } from "../auth/service-tier-routing.js"; import type { Context } from "hono"; import type { StatusCode } from "hono/utils/http-status"; import type { AccountPool } from "../auth/account-pool.js"; @@ -94,6 +95,8 @@ export async function handleCompact( const parsed = parseModelName(rawModel); const modelId = resolveModelId(parsed.modelId); + const serviceTier = getModelServiceTierOverride(modelId) ?? + (typeof body.service_tier === "string" ? body.service_tier : parsed.serviceTier); const compactRequest: CodexCompactRequest = { model: modelId, @@ -158,7 +161,7 @@ export async function handleCompact( const triedEntryIds: string[] = []; const released = new Set(); - const acquired = acquireAccount(accountPool, modelId, undefined, TAG); + const acquired = acquireAccount(accountPool, modelId, undefined, TAG, undefined, serviceTier); if (!acquired) { c.status(503); return c.json(formatResponsesError(503, "No available accounts. All accounts are expired or rate-limited.")); @@ -205,7 +208,7 @@ export async function handleCompact( releaseAccount(accountPool, entryId, annotateUsageCost(modelId, compactImageFailedUsage), released); } - const retry = acquireAccount(accountPool, modelId, triedEntryIds, TAG); + const retry = acquireAccount(accountPool, modelId, triedEntryIds, TAG, undefined, serviceTier); if (!retry) { const status = decision.status as StatusCode; c.status(status); diff --git a/src/routes/shared/account-acquisition.ts b/src/routes/shared/account-acquisition.ts index ed8131b96..0c50b209b 100644 --- a/src/routes/shared/account-acquisition.ts +++ b/src/routes/shared/account-acquisition.ts @@ -18,8 +18,9 @@ export function acquireAccount( excludeIds?: string[], tag?: string, preferredEntryId?: string, + serviceTier?: string | null, ): AcquiredAccount | null { - const acquired = pool.acquire({ model, excludeIds, preferredEntryId }); + const acquired = pool.acquire({ model, excludeIds, preferredEntryId, ...(serviceTier != null ? { serviceTier } : {}) }); if (!acquired && tag) { console.warn(`[${tag}] No available account for model "${model}"`); } diff --git a/src/routes/shared/non-streaming-helpers.ts b/src/routes/shared/non-streaming-helpers.ts index 523cf88cd..2dd069c0d 100644 --- a/src/routes/shared/non-streaming-helpers.ts +++ b/src/routes/shared/non-streaming-helpers.ts @@ -291,7 +291,7 @@ export async function retryNonStreamingEmptyResponse( releaseAccount(accountPool, currentEntryId, annotateUsageCost(req.model, annotateImageGenOutcome(collectErr.usage, req.expectsImageGen)), released); restoreImplicitResumeRequest?.(); - const acquired = acquireAccount(accountPool, req.codexRequest.model, undefined, tag); + const acquired = acquireAccount(accountPool, req.codexRequest.model, undefined, tag, undefined, req.codexRequest.service_tier); if (!acquired) { return { action: "respond", diff --git a/src/routes/shared/proxy-error-retry-transition.ts b/src/routes/shared/proxy-error-retry-transition.ts index 2f7392c2f..8f1701701 100644 --- a/src/routes/shared/proxy-error-retry-transition.ts +++ b/src/routes/shared/proxy-error-retry-transition.ts @@ -31,6 +31,7 @@ export interface ApplyProxyErrorRetryTransitionOptions { accountPool: AccountPool; entryId: string; model: string; + serviceTier?: string | null; triedEntryIds: string[]; tag: string; decision: ErrorAction; @@ -49,6 +50,7 @@ export function applyProxyErrorRetryTransition( accountPool, entryId, model, + serviceTier, triedEntryIds, tag, decision, @@ -79,6 +81,7 @@ export function applyProxyErrorRetryTransition( const fallbackRetry = prepareProxyFallbackAccountRetry({ accountPool, model, + serviceTier, triedEntryIds, tag, decision, diff --git a/src/routes/shared/proxy-fallback-account-retry.ts b/src/routes/shared/proxy-fallback-account-retry.ts index 8b790af27..550232dc3 100644 --- a/src/routes/shared/proxy-fallback-account-retry.ts +++ b/src/routes/shared/proxy-fallback-account-retry.ts @@ -27,6 +27,7 @@ export type ProxyFallbackAccountRetryResult = export interface PrepareProxyFallbackAccountRetryOptions { accountPool: AccountPool; model: string; + serviceTier?: string | null; triedEntryIds: string[]; tag: string; decision: RetryDecision; @@ -41,6 +42,7 @@ export function prepareProxyFallbackAccountRetry( const { accountPool, model, + serviceTier, triedEntryIds, tag, decision, @@ -62,7 +64,7 @@ export function prepareProxyFallbackAccountRetry( return fallbackPlan; } - const retry = acquireAccount(accountPool, model, excludeEntryIds, tag); + const retry = acquireAccount(accountPool, model, excludeEntryIds, tag, undefined, serviceTier); if (!retry) { return { action: "respond", diff --git a/src/routes/shared/proxy-handler.ts b/src/routes/shared/proxy-handler.ts index 916ed1ea2..35b94cee9 100644 --- a/src/routes/shared/proxy-handler.ts +++ b/src/routes/shared/proxy-handler.ts @@ -22,6 +22,7 @@ * - non-streaming-handler.ts — collect / retry response lifecycle */ +import { getServiceTierAccountRule, getModelServiceTierOverride } from "../../auth/service-tier-routing.js"; import { CodexApi, CodexApiError, PreviousResponseWebSocketError } from "../../proxy/codex-api.js"; import { toQuota } from "../../auth/quota-utils.js"; import { markFallbackUsed } from "../../auth/fallback-state.js"; @@ -78,7 +79,9 @@ async function respondNoAccountOrFallback( req: ProxyRequest, fmt: FormatAdapter, ): Promise { - const fallback = options.fallbackUpstream?.get(); + const fallback = getServiceTierAccountRule(req.codexRequest.service_tier) + ? undefined + : options.fallbackUpstream?.get(); if (fallback) { console.log( `[${fmt.tag}] No available OAuth accounts — routing through fallback upstream apikey (${fallback.baseUrl})`, @@ -109,7 +112,9 @@ async function respondProxyErrorOrFallback( message: string, useFormat429?: boolean, ): Promise { - const fallback = options.fallbackUpstream?.get(); + const fallback = getServiceTierAccountRule(req.codexRequest.service_tier) + ? undefined + : options.fallbackUpstream?.get(); if (fallback) { console.log( `[${fmt.tag}] Retry exhausted — routing through fallback upstream apikey (${fallback.baseUrl})`, @@ -128,6 +133,9 @@ async function respondProxyErrorOrFallback( export async function handleProxyRequest(options: HandleProxyRequestOptions): Promise { const { c, accountPool, cookieJar, req, fmt, proxyPool } = options; c.set("logForwarded", true); + const forcedTier = getModelServiceTierOverride(req.codexRequest.model); + if (forcedTier !== undefined) req.codexRequest.service_tier = forcedTier; + const affinityMap = getSessionAffinityMap(); const requestId = c.get("requestId") ?? randomUUID().slice(0, 8); @@ -149,7 +157,7 @@ export async function handleProxyRequest(options: HandleProxyRequestOptions): Pr const verifiedExcludeIds: string[] = []; // Single acquire call — preferredEntryId is a hint, not a hard requirement - let acquired = acquireAccount(accountPool, req.codexRequest.model, undefined, fmt.tag, sessionContext.preferredEntryId ?? undefined); + let acquired = acquireAccount(accountPool, req.codexRequest.model, undefined, fmt.tag, sessionContext.preferredEntryId ?? undefined, req.codexRequest.service_tier); if (!acquired) { return respondNoAccountOrFallback(options, req, fmt); } @@ -188,7 +196,7 @@ export async function handleProxyRequest(options: HandleProxyRequestOptions): Pr return respondNoAccountOrFallback(options, req, fmt); } - acquired = acquireAccount(accountPool, req.codexRequest.model, verifiedExcludeIds, fmt.tag, sessionContext.preferredEntryId ?? undefined); + acquired = acquireAccount(accountPool, req.codexRequest.model, verifiedExcludeIds, fmt.tag, sessionContext.preferredEntryId ?? undefined, req.codexRequest.service_tier); if (!acquired) { return respondNoAccountOrFallback(options, req, fmt); } @@ -523,6 +531,7 @@ export async function handleProxyRequest(options: HandleProxyRequestOptions): Pr const errorRetryTransition = applyProxyErrorRetryTransition({ accountPool, entryId, model: req.codexRequest.model, + serviceTier: req.codexRequest.service_tier, triedEntryIds, tag: fmt.tag, decision, released, restoreImplicitResumeRequest: implicitResume.restore, diff --git a/tests/integration/proxy-handler.test.ts b/tests/integration/proxy-handler.test.ts index 9b0b2f3a9..e24f16c08 100644 --- a/tests/integration/proxy-handler.test.ts +++ b/tests/integration/proxy-handler.test.ts @@ -149,6 +149,8 @@ vi.mock("@src/translation/codex-event-extractor.js", () => { }); // Import after mocks are set up +import { getConfig } from "@src/config.js"; +import type { FallbackUpstreamStore } from "@src/auth/fallback-upstream.js"; import { handleProxyRequest } from "@src/routes/shared/proxy-handler.js"; import { CodexApiError, PreviousResponseWebSocketError } from "@src/proxy/codex-api.js"; import { EmptyResponseError, UpstreamPrematureCloseError } from "@src/translation/codex-event-extractor.js"; @@ -202,6 +204,7 @@ function buildTestApp(opts: { fmt?: ReturnType; req?: ProxyRequest; cookieJar?: unknown; + fallbackUpstream?: FallbackUpstreamStore; }) { const accountPool = opts.accountPool ?? createMockAccountPool(); const fmt = opts.fmt ?? createMockFormatAdapter(); @@ -215,6 +218,7 @@ function buildTestApp(opts: { accountPool: accountPool as never, cookieJar, req: proxyReq, + fallbackUpstream: opts.fallbackUpstream, fmt, }), ); @@ -233,6 +237,58 @@ describe("proxy-handler integration", () => { vi.clearAllMocks(); }); + it.each([undefined, "default", "priority"])("forces the configured model tier over client preference %s", async (tier) => { + vi.mocked(getConfig).mockReturnValueOnce({ + auth: {}, model: { service_tier_overrides: { "gpt-6.1-sol": "ultrafast" } }, + } as never); + const req = createDefaultRequest(); + req.codexRequest.model = "gpt-6.1-sol"; + req.codexRequest.service_tier = tier; + mockCreateResponse = async (request) => { + expect(request.service_tier).toBe("ultrafast"); + return new Response("data: {}\n\n"); + }; + const { app, accountPool } = buildTestApp({ req }); + expect((await app.request("/test", { method: "POST" })).status).toBe(200); + expect(accountPool.acquire).toHaveBeenCalledWith(expect.objectContaining({ serviceTier: "ultrafast" })); + expect(req.codexRequest.service_tier).toBe("ultrafast"); + }); + + it("preserves the requested tier when rotating after a rate limit", async () => { + let count = 0; + mockCreateResponse = () => ++count === 1 + ? Promise.reject(new CodexApiError(429, JSON.stringify({ error: { type: "usage_limit_reached", resets_in_seconds: 60 } }))) + : Promise.resolve(new Response("data: {}\n\n")); + const pool = createMockAccountPool({ acquire: vi.fn() + .mockReturnValueOnce({ entryId: "e1", token: "t1", accountId: "a1" }) + .mockReturnValueOnce({ entryId: "e2", token: "t2", accountId: "a2" }) }); + const req = createDefaultRequest(); + req.codexRequest.service_tier = "ultrafast"; + const { app } = buildTestApp({ accountPool: pool, req }); + expect((await app.request("/test", { method: "POST" })).status).toBe(200); + expect(pool.acquire).toHaveBeenCalledTimes(2); + for (const [options] of pool.acquire.mock.calls) { + expect(options).toMatchObject({ serviceTier: "ultrafast" }); + } + }); + + it("does not escape a tier restriction through the API-key fallback", async () => { + vi.mocked(getConfig).mockReturnValueOnce({ auth: { + service_tier_routing: { ultrafast: { account_ids: ["reserved"] } }, + }, model: {} } as never).mockReturnValueOnce({ auth: { + service_tier_routing: { ultrafast: { account_ids: ["reserved"] } }, + }, model: {} } as never); + const get = vi.fn(() => ({ apiKey: "secret", baseUrl: "https://upstream.invalid" })); + const req = createDefaultRequest(); + req.codexRequest.service_tier = "ultrafast"; + const { app } = buildTestApp({ + req, accountPool: createMockAccountPool({ acquire: vi.fn(() => null) }), + fallbackUpstream: { get } as unknown as FallbackUpstreamStore, + }); + expect((await app.request("/test", { method: "POST" })).status).toBe(503); + expect(get).not.toHaveBeenCalled(); + }); + // 1. No account available it("returns noAccountStatus (503) when no account is available", async () => { const accountPool = createMockAccountPool({ diff --git a/tests/unit/auth/plan-routing-acquire.test.ts b/tests/unit/auth/plan-routing-acquire.test.ts index ac669b0a1..907ebc816 100644 --- a/tests/unit/auth/plan-routing-acquire.test.ts +++ b/tests/unit/auth/plan-routing-acquire.test.ts @@ -194,4 +194,82 @@ describe("account-pool plan-based routing", () => { expect(spark).not.toBeNull(); expect(spark!.token).toBe(jwts.get("plus2")); }); + + function routingPool() { + return createPool( + { accountId: "fast", planType: "pro", email: "fast@test.com" }, + { accountId: "normal", planType: "plus", email: "normal@test.com" }, + ); + } + + it("reserves a plan for ultrafast and uses other plans for default-tier Astra", () => { + setConfigForTesting(createMockConfig({ auth: { service_tier_routing: { + ultrafast: { plan_types: ["pro"] }, default: { plan_types: ["plus"] }, + } } })); + const { pool, jwts } = routingPool(); + expect(pool.acquire({ model: "gpt-6.1-sol", serviceTier: "ultrafast" })?.token).toBe(jwts.get("fast")); + expect(pool.acquire({ model: "gpt-6-astra" })?.token).toBe(jwts.get("normal")); + }); + + it("does not bypass tier restrictions for session affinity or global plan priority", () => { + setConfigForTesting(createMockConfig({ auth: { + tier_priority: ["plus", "pro"], service_tier_routing: { ultrafast: { plan_types: ["pro"] } }, + } })); + const { pool, jwts } = routingPool(); + const preferredEntryId = pool.getAccounts().find(a => a.accountId === "normal")!.id; + expect(pool.acquire({ serviceTier: "ultrafast", preferredEntryId })?.token).toBe(jwts.get("fast")); + }); + + it("fails closed when the only matching account was already tried", () => { + setConfigForTesting(createMockConfig({ auth: { service_tier_routing: { ultrafast: { plan_types: ["pro"] } } } })); + const { pool } = routingPool(); + const first = pool.acquire({ serviceTier: "ultrafast" })!; + pool.release(first.entryId); + expect(pool.acquire({ serviceTier: "ultrafast", excludeIds: [first.entryId] })).toBeNull(); + }); + + it("fails closed when matching accounts reach their concurrency limit", () => { + setConfigForTesting(createMockConfig({ auth: { + max_concurrent_per_account: 1, service_tier_routing: { ultrafast: { plan_types: ["pro"] } }, + } })); + const { pool } = routingPool(); + expect(pool.acquire({ serviceTier: "ultrafast" })).not.toBeNull(); + expect(pool.acquire({ serviceTier: "ultrafast" })).toBeNull(); + }); + + it("intersects service-tier rules with known model eligibility", () => { + setConfigForTesting(createMockConfig({ auth: { service_tier_routing: { ultrafast: { plan_types: ["pro"] } } } })); + vi.mocked(getModelPlanTypes).mockReturnValue(["plus"]); + const { pool } = routingPool(); + expect(pool.acquire({ model: "plus-only", serviceTier: "ultrafast" })).toBeNull(); + }); + + it("uses explicit IDs to distinguish accounts with the same reported plan", () => { + const { pool, jwts } = createPool( + { accountId: "fast", planType: "pro", email: "fast@test.com" }, + { accountId: "normal", planType: "pro", email: "normal@test.com" }, + ); + const id = pool.getAccounts().find(a => a.accountId === "fast")!.id; + setConfigForTesting(createMockConfig({ auth: { service_tier_routing: { + ultrafast: { account_ids: [id] }, default: { exclude_account_ids: [id] }, + } } })); + expect(pool.acquire({ serviceTier: "ultrafast" })?.token).toBe(jwts.get("fast")); + expect(pool.acquire({})?.token).toBe(jwts.get("normal")); + }); + + it("uses the configured default tier and honors explicit overrides", () => { + setConfigForTesting(createMockConfig({ model: { default_service_tier: "ultrafast" }, auth: { + service_tier_routing: { ultrafast: { plan_types: ["pro"] }, default: { plan_types: ["plus"] } }, + } })); + const { pool, jwts } = routingPool(); + expect(pool.acquire({})?.token).toBe(jwts.get("fast")); + expect(pool.acquire({ serviceTier: "default" })?.token).toBe(jwts.get("normal")); + }); + + it("does not restrict unmatched tiers", () => { + setConfigForTesting(createMockConfig({ auth: { service_tier_routing: { ultrafast: { plan_types: ["pro"] } } } })); + const { pool } = createPool({ accountId: "normal", planType: "plus", email: "normal@test.com" }); + expect(pool.acquire({ serviceTier: "priority" })).not.toBeNull(); + }); + }); diff --git a/tests/unit/config-schema.test.ts b/tests/unit/config-schema.test.ts index a7b4fea62..168c15f2b 100644 --- a/tests/unit/config-schema.test.ts +++ b/tests/unit/config-schema.test.ts @@ -292,3 +292,20 @@ describe("FingerprintSchema", () => { expect(result.success).toBe(false); }); }); + + +describe("service-tier routing configuration", () => { + const parse = (rules: unknown) => ConfigSchema.safeParse({ + api: {}, client: {}, model: {}, auth: { service_tier_routing: rules }, server: {}, session: {}, + }); + it("defaults to no restrictions", () => { + const result = parse(undefined); + expect(result.success && result.data.auth.service_tier_routing).toEqual({}); + }); + it("accepts plan and account restrictions", () => { + expect(parse({ ultrafast: { plan_types: ["pro"], account_ids: ["account-1"] } }).success).toBe(true); + }); + it.each([{}, { plan_types: [] }, { account_ids: [""] }, { plans: ["pro"] }])("rejects an invalid rule %j", rule => { + expect(parse({ ultrafast: rule }).success).toBe(false); + }); +}); From 8323c4f2c0b15b9f9aeffaeee32a8f0445aa5ceb Mon Sep 17 00:00:00 2001 From: Mako-L <9460018+Mako-L@users.noreply.github.com> Date: Thu, 1 Oct 2026 04:32:55 +0300 Subject: [PATCH 2/3] feat: fall back to default tier when fast account unavailable --- README_EN.md | 7 ++++- src/auth/account-lifecycle.ts | 9 +++++- src/auth/types.ts | 2 ++ src/config-schema.ts | 1 + src/routes/shared/non-streaming-helpers.ts | 1 + .../shared/proxy-error-retry-transition.ts | 2 ++ .../shared/proxy-fallback-account-retry.ts | 2 ++ src/routes/shared/proxy-handler.ts | 2 ++ tests/integration/proxy-handler.test.ts | 30 +++++++++++++++++++ tests/unit/auth/plan-routing-acquire.test.ts | 24 +++++++++++++++ 10 files changed, 78 insertions(+), 2 deletions(-) diff --git a/README_EN.md b/README_EN.md index ac19c260f..6c28c02aa 100644 --- a/README_EN.md +++ b/README_EN.md @@ -720,6 +720,7 @@ auth: service_tier_routing: ultrafast: account_ids: ["reserved-account-entry-id"] + fallback_to_default: true default: exclude_account_ids: ["reserved-account-entry-id"] standard: @@ -742,7 +743,11 @@ Unknown plan types cannot satisfy a `plan_types` restriction. Rules are applied after model eligibility, quota and concurrency checks, before plan priority, session affinity and rotation. Account retries retain the tier. -If all matching accounts are unavailable, the request fails instead of selecting +With `fallback_to_default: true`, unavailable fast-tier accounts trigger another +selection using the `default` rule and send `service_tier: default` upstream. +This covers missing, disabled, busy, exhausted, and previously tried accounts. +The model and reasoning effort are unchanged. Fallback is one-way and optional. +Without this option, if all matching accounts are unavailable, the request fails instead of selecting an excluded account or using the API-key fallback. Compact requests also apply these account restrictions, without adding a service tier to the compact payload. Explicitly routed third-party API providers do not use the OAuth account pool and diff --git a/src/auth/account-lifecycle.ts b/src/auth/account-lifecycle.ts index 99fc8209b..eb72463dd 100644 --- a/src/auth/account-lifecycle.ts +++ b/src/auth/account-lifecycle.ts @@ -125,7 +125,14 @@ export class AccountLifecycle { (!rule.account_ids || rule.account_ids.includes(account.id)) && !rule.exclude_account_ids?.includes(account.id), ); - if (candidates.length === 0) return null; + if (candidates.length === 0) { + const tier = options?.serviceTier ?? config.model.default_service_tier ?? "default"; + if (rule.fallback_to_default && tier !== "default") { + const fallback = this.acquire({ ...options, serviceTier: "default" }); + return fallback ? { ...fallback, serviceTier: "default" } : null; + } + return null; + } } // Tier-based filtering: when configured, restrict to the highest available tier diff --git a/src/auth/types.ts b/src/auth/types.ts index fe9a70e56..b82d4cd83 100644 --- a/src/auth/types.ts +++ b/src/auth/types.ts @@ -165,6 +165,8 @@ export interface CodexQuota { /** Returned by acquire() */ export interface AcquiredAccount { + /** Set when account selection downgraded the requested service tier. */ + serviceTier?: "default"; entryId: string; token: string; accountId: string | null; diff --git a/src/config-schema.ts b/src/config-schema.ts index 57cfab8c2..4a17ed4e6 100644 --- a/src/config-schema.ts +++ b/src/config-schema.ts @@ -126,6 +126,7 @@ export const ConfigSchema = z.object({ tier_priority: z.array(z.string()).nullable().default(null), /** Hard account restrictions keyed by the effective service tier. */ service_tier_routing: z.record(z.string().trim().min(1), z.object({ + fallback_to_default: z.boolean().optional(), plan_types: z.array(z.string().trim().min(1)).min(1).optional(), account_ids: z.array(z.string().trim().min(1)).min(1).optional(), exclude_account_ids: z.array(z.string().trim().min(1)).min(1).optional(), diff --git a/src/routes/shared/non-streaming-helpers.ts b/src/routes/shared/non-streaming-helpers.ts index 2dd069c0d..3f99aa61c 100644 --- a/src/routes/shared/non-streaming-helpers.ts +++ b/src/routes/shared/non-streaming-helpers.ts @@ -300,6 +300,7 @@ export async function retryNonStreamingEmptyResponse( }; } + if (acquired.serviceTier) req.codexRequest.service_tier = acquired.serviceTier; const nextApi = buildCodexApi( acquired.token, acquired.accountId, diff --git a/src/routes/shared/proxy-error-retry-transition.ts b/src/routes/shared/proxy-error-retry-transition.ts index 8f1701701..ee58ca0fe 100644 --- a/src/routes/shared/proxy-error-retry-transition.ts +++ b/src/routes/shared/proxy-error-retry-transition.ts @@ -24,6 +24,7 @@ export type ProxyErrorRetryTransitionResult = entryId: string; api: CodexApi; prevSlotMs: number | null; + serviceTier?: "default"; modelRetried: boolean; }; @@ -105,6 +106,7 @@ export function applyProxyErrorRetryTransition( entryId: fallbackRetry.entryId, api: fallbackRetry.api, prevSlotMs: fallbackRetry.prevSlotMs, + ...(fallbackRetry.serviceTier ? { serviceTier: fallbackRetry.serviceTier } : {}), modelRetried: nextModelRetried, }; } diff --git a/src/routes/shared/proxy-fallback-account-retry.ts b/src/routes/shared/proxy-fallback-account-retry.ts index 550232dc3..c04e5a646 100644 --- a/src/routes/shared/proxy-fallback-account-retry.ts +++ b/src/routes/shared/proxy-fallback-account-retry.ts @@ -22,6 +22,7 @@ export type ProxyFallbackAccountRetryResult = entryId: string; api: CodexApi; prevSlotMs: number | null; + serviceTier?: "default"; }; export interface PrepareProxyFallbackAccountRetryOptions { @@ -89,5 +90,6 @@ export function prepareProxyFallbackAccountRetry( entryId: retry.entryId, api, prevSlotMs: retry.prevSlotMs, + ...(retry.serviceTier ? { serviceTier: retry.serviceTier } : {}), }; } diff --git a/src/routes/shared/proxy-handler.ts b/src/routes/shared/proxy-handler.ts index 35b94cee9..672aa9951 100644 --- a/src/routes/shared/proxy-handler.ts +++ b/src/routes/shared/proxy-handler.ts @@ -213,6 +213,7 @@ export async function handleProxyRequest(options: HandleProxyRequestOptions): Pr } if (!acquired) return respondNoAccountOrFallback(options, req, fmt); + if (acquired.serviceTier) req.codexRequest.service_tier = acquired.serviceTier; let { entryId } = acquired; // First account this request acquired; later attempts that switch to another // entry (fallback account retry) are marked as fallback in the audit log. @@ -563,6 +564,7 @@ export async function handleProxyRequest(options: HandleProxyRequestOptions): Pr if (decision.action === "retry" && decision.markTransportRetried) { transportRetried = true; } + if (errorRetryTransition.serviceTier) req.codexRequest.service_tier = errorRetryTransition.serviceTier; entryId = errorRetryTransition.entryId; triedEntryIds.push(errorRetryTransition.entryId); codexApi = errorRetryTransition.api; diff --git a/tests/integration/proxy-handler.test.ts b/tests/integration/proxy-handler.test.ts index e24f16c08..691c02e81 100644 --- a/tests/integration/proxy-handler.test.ts +++ b/tests/integration/proxy-handler.test.ts @@ -254,6 +254,36 @@ describe("proxy-handler integration", () => { expect(req.codexRequest.service_tier).toBe("ultrafast"); }); + it("sends default upstream after account selection downgrades ultrafast", async () => { + const req = createDefaultRequest(); + req.codexRequest.service_tier = "ultrafast"; + mockCreateResponse = async (request) => { + expect(request.service_tier).toBe("default"); + return new Response("data: {}\n\n"); + }; + const { app } = buildTestApp({ req, accountPool: createMockAccountPool({ + acquire: vi.fn(() => ({ entryId: "normal", token: "t", accountId: "normal", serviceTier: "default" })), + }) }); + expect((await app.request("/test", { method: "POST" })).status).toBe(200); + }); + + it("downgrades the outgoing tier on a rate-limit account retry", async () => { + const tiers: unknown[] = []; + mockCreateResponse = async (request) => { + tiers.push(request.service_tier); + if (tiers.length === 1) throw new CodexApiError(429, JSON.stringify({ error: { type: "usage_limit_reached", resets_in_seconds: 60 } })); + return new Response("data: {}\n\n"); + }; + const req = createDefaultRequest(); + req.codexRequest.service_tier = "ultrafast"; + const pool = createMockAccountPool({ acquire: vi.fn() + .mockReturnValueOnce({ entryId: "fast", token: "t1", accountId: "a1" }) + .mockReturnValueOnce({ entryId: "normal", token: "t2", accountId: "a2", serviceTier: "default" }) }); + const { app } = buildTestApp({ req, accountPool: pool }); + expect((await app.request("/test", { method: "POST" })).status).toBe(200); + expect(tiers).toEqual(["ultrafast", "default"]); + }); + it("preserves the requested tier when rotating after a rate limit", async () => { let count = 0; mockCreateResponse = () => ++count === 1 diff --git a/tests/unit/auth/plan-routing-acquire.test.ts b/tests/unit/auth/plan-routing-acquire.test.ts index 907ebc816..9ca86f8b3 100644 --- a/tests/unit/auth/plan-routing-acquire.test.ts +++ b/tests/unit/auth/plan-routing-acquire.test.ts @@ -272,4 +272,28 @@ describe("account-pool plan-based routing", () => { expect(pool.acquire({ serviceTier: "priority" })).not.toBeNull(); }); + it("downgrades to default when the fast account is unavailable", () => { + setConfigForTesting(createMockConfig({ auth: { max_concurrent_per_account: 1, service_tier_routing: { + ultrafast: { plan_types: ["pro"], fallback_to_default: true }, + default: { plan_types: ["plus"] }, + } } })); + const { pool, jwts } = routingPool(); + expect(pool.acquire({ serviceTier: "ultrafast" })?.token).toBe(jwts.get("fast")); + const fallback = pool.acquire({ serviceTier: "ultrafast" }); + expect(fallback?.token).toBe(jwts.get("normal")); + expect(fallback?.serviceTier).toBe("default"); + expect(pool.acquire({ serviceTier: "ultrafast" })).toBeNull(); + }); + + it("uses default when no fast-plan account exists, preserving retry exclusions", () => { + setConfigForTesting(createMockConfig({ auth: { service_tier_routing: { + ultrafast: { plan_types: ["pro"], fallback_to_default: true }, + } } })); + const { pool } = createPool({ accountId: "normal", planType: "plus", email: "normal@test.com" }); + const fallback = pool.acquire({ serviceTier: "ultrafast" })!; + expect(fallback.serviceTier).toBe("default"); + pool.release(fallback.entryId); + expect(pool.acquire({ serviceTier: "ultrafast", excludeIds: [fallback.entryId] })).toBeNull(); + }); + }); From 19af9bec874b3b012e39f3183228056e70ce8d53 Mon Sep 17 00:00:00 2001 From: SsuJo_ <1049731887@qq.com> Date: Thu, 8 Oct 2026 21:02:18 +0800 Subject: [PATCH 3/3] fix(auth): preserve tier restriction through fallback --- src/routes/shared/proxy-handler.ts | 23 ++++++++++++-------- tests/integration/proxy-handler.test.ts | 28 +++++++++++++++++++------ 2 files changed, 36 insertions(+), 15 deletions(-) diff --git a/src/routes/shared/proxy-handler.ts b/src/routes/shared/proxy-handler.ts index 672aa9951..fe71816ad 100644 --- a/src/routes/shared/proxy-handler.ts +++ b/src/routes/shared/proxy-handler.ts @@ -78,8 +78,9 @@ async function respondNoAccountOrFallback( options: HandleProxyRequestOptions, req: ProxyRequest, fmt: FormatAdapter, + tierRestricted: boolean, ): Promise { - const fallback = getServiceTierAccountRule(req.codexRequest.service_tier) + const fallback = tierRestricted ? undefined : options.fallbackUpstream?.get(); if (fallback) { @@ -110,9 +111,10 @@ async function respondProxyErrorOrFallback( fmt: FormatAdapter, status: number, message: string, - useFormat429?: boolean, + useFormat429: boolean | undefined, + tierRestricted: boolean, ): Promise { - const fallback = getServiceTierAccountRule(req.codexRequest.service_tier) + const fallback = tierRestricted ? undefined : options.fallbackUpstream?.get(); if (fallback) { @@ -135,7 +137,9 @@ export async function handleProxyRequest(options: HandleProxyRequestOptions): Pr c.set("logForwarded", true); const forcedTier = getModelServiceTierOverride(req.codexRequest.model); if (forcedTier !== undefined) req.codexRequest.service_tier = forcedTier; - + // Keep the original route constraint even if account acquisition downgrades + // service_tier to `default`; API-key fallback must not escape that constraint. + const tierRestricted = Boolean(getServiceTierAccountRule(req.codexRequest.service_tier)); const affinityMap = getSessionAffinityMap(); const requestId = c.get("requestId") ?? randomUUID().slice(0, 8); @@ -159,7 +163,7 @@ export async function handleProxyRequest(options: HandleProxyRequestOptions): Pr // Single acquire call — preferredEntryId is a hint, not a hard requirement let acquired = acquireAccount(accountPool, req.codexRequest.model, undefined, fmt.tag, sessionContext.preferredEntryId ?? undefined, req.codexRequest.service_tier); if (!acquired) { - return respondNoAccountOrFallback(options, req, fmt); + return respondNoAccountOrFallback(options, req, fmt, tierRestricted); } // ── Drift-Defense & Verification Loop ── @@ -168,7 +172,7 @@ export async function handleProxyRequest(options: HandleProxyRequestOptions): Pr const MAX_VERIFY_ATTEMPTS = 5; let verifyAttempts = 0; for (;;) { - if (!acquired) return respondNoAccountOrFallback(options, req, fmt); + if (!acquired) return respondNoAccountOrFallback(options, req, fmt, tierRestricted); const entry = accountPool.getEntry(acquired.entryId); if (entry?.quotaVerifyRequired) { const verifyingEntryId = acquired.entryId; @@ -193,12 +197,12 @@ export async function handleProxyRequest(options: HandleProxyRequestOptions): Pr verifyAttempts++; if (verifyAttempts >= MAX_VERIFY_ATTEMPTS) { console.warn(`[${fmt.tag}] ⚠️ Drift-defense hit MAX_VERIFY_ATTEMPTS (${MAX_VERIFY_ATTEMPTS}). Giving up to avoid excess upstream calls.`); - return respondNoAccountOrFallback(options, req, fmt); + return respondNoAccountOrFallback(options, req, fmt, tierRestricted); } acquired = acquireAccount(accountPool, req.codexRequest.model, verifiedExcludeIds, fmt.tag, sessionContext.preferredEntryId ?? undefined, req.codexRequest.service_tier); if (!acquired) { - return respondNoAccountOrFallback(options, req, fmt); + return respondNoAccountOrFallback(options, req, fmt, tierRestricted); } continue; // Loop back to check the newly acquired account } @@ -212,7 +216,7 @@ export async function handleProxyRequest(options: HandleProxyRequestOptions): Pr break; // Verified or no verification required, proceed! } - if (!acquired) return respondNoAccountOrFallback(options, req, fmt); + if (!acquired) return respondNoAccountOrFallback(options, req, fmt, tierRestricted); if (acquired.serviceTier) req.codexRequest.service_tier = acquired.serviceTier; let { entryId } = acquired; // First account this request acquired; later attempts that switch to another @@ -547,6 +551,7 @@ export async function handleProxyRequest(options: HandleProxyRequestOptions): Pr errorRetryTransition.status, errorRetryTransition.message, errorRetryTransition.useFormat429, + tierRestricted, ); } return respondWithProxyError({ diff --git a/tests/integration/proxy-handler.test.ts b/tests/integration/proxy-handler.test.ts index 691c02e81..a7d17322f 100644 --- a/tests/integration/proxy-handler.test.ts +++ b/tests/integration/proxy-handler.test.ts @@ -302,21 +302,37 @@ describe("proxy-handler integration", () => { } }); - it("does not escape a tier restriction through the API-key fallback", async () => { - vi.mocked(getConfig).mockReturnValueOnce({ auth: { - service_tier_routing: { ultrafast: { account_ids: ["reserved"] } }, - }, model: {} } as never).mockReturnValueOnce({ auth: { + it("does not escape a downgraded tier restriction through API-key fallback", async () => { + vi.mocked(getConfig).mockReturnValue({ auth: { service_tier_routing: { ultrafast: { account_ids: ["reserved"] } }, }, model: {} } as never); + mockCreateResponse = async () => { + throw new CodexApiError(429, JSON.stringify({ error: { type: "usage_limit_reached", resets_in_seconds: 60 } })); + }; const get = vi.fn(() => ({ apiKey: "secret", baseUrl: "https://upstream.invalid" })); const req = createDefaultRequest(); req.codexRequest.service_tier = "ultrafast"; + const pool = createMockAccountPool({ + acquire: vi.fn() + .mockReturnValueOnce({ entryId: "reserved", token: "t", accountId: "reserved", serviceTier: "default" }) + .mockReturnValue(null), + }); + const { app } = buildTestApp({ req, accountPool: pool, + fallbackUpstream: { get } as unknown as FallbackUpstreamStore }); + expect((await app.request("/test", { method: "POST" })).status).toBe(429); + expect(req.codexRequest.service_tier).toBe("default"); + expect(get).not.toHaveBeenCalled(); + }); + + it("keeps API-key fallback available without a tier restriction", async () => { + const get = vi.fn(() => ({ apiKey: "secret", baseUrl: "https://upstream.invalid" })); + const req = createDefaultRequest(); const { app } = buildTestApp({ req, accountPool: createMockAccountPool({ acquire: vi.fn(() => null) }), fallbackUpstream: { get } as unknown as FallbackUpstreamStore, }); - expect((await app.request("/test", { method: "POST" })).status).toBe(503); - expect(get).not.toHaveBeenCalled(); + expect((await app.request("/test", { method: "POST" })).status).toBe(502); + expect(get).toHaveBeenCalledOnce(); }); // 1. No account available