From b6a142d3d782b831343415a38dbcfa2df7c63c79 Mon Sep 17 00:00:00 2001 From: JUN Date: Sat, 26 Sep 2026 19:19:46 +0900 Subject: [PATCH 1/3] docs(devlog): record the bug-train roadmap after batch 6 --- .../_plan/260926_bug_train_6/010_roadmap.md | 25 +++++++++++++++++++ 1 file changed, 25 insertions(+) create mode 100644 devlog/_plan/260926_bug_train_6/010_roadmap.md diff --git a/devlog/_plan/260926_bug_train_6/010_roadmap.md b/devlog/_plan/260926_bug_train_6/010_roadmap.md new file mode 100644 index 00000000000..da676d9613d --- /dev/null +++ b/devlog/_plan/260926_bug_train_6/010_roadmap.md @@ -0,0 +1,25 @@ +# Bug-PR merge train after batch 6 — roadmap + +Batch 6 landed as #5918 (`76b26a0881`). Open bug-labelled PRs re-queried at that head: +#5921, #5920, #5916, #5915, #5914, #5911, #5831, #5800, #5782, #5539, #5497, #4222. + +## Decisions + +| PR | Decision | Reason | +|---|---|---| +| #5914, #5882, #5849 | Close as landed | Carried in #5918. Close each with the same thank-you note earlier batches used (author credit, merge SHA, what changed on top). | +| #5916 (OAuth 429 send budget, #5880) | **Carry in batch 7 with a fix** | Real starvation bug and meaningful tests (4 cases red on `dev`, 86 focused pass). The independent security review failed on one point: the allowance scales with roster size with no fixed ceiling, so one request can fan out to 3 × N sends. Batch 7 adds a hard per-request ceiling that roster size cannot raise, then re-runs the security review. | +| #5497 (FastWire tier authority for relays) | Evaluate in batch 7 P | Bug-labelled, 357 lines, but it adds a provider config field and fails hygiene. Carry only if the failure is mechanical and the field is opt-in with no default change; otherwise leave it. | +| #5831 (main-lock recovery from two-window WHAM) | NEEDS_HUMAN | @Ingwannu holds approval for an owner decision on whether an omitted short window counts as proof no window exists. That is a policy call for the owner. | +| #5911, #5915 (Meta Muse OAuth continuations) | Leave | Two overlapping drafts on the same OAuth surface; #5911 is the maintainer's own active draft with an ADR. Consolidation is theirs to decide. | +| #5539 (effort wire mapper on unpinned routes) | Evaluate a narrowed carry in batch 7 P | The failure is real: strict upstreams answer `400 Invalid option` for `minimal`/`ultra` sent to a provider with no configured ladder. The native Chat half reverses tests that deliberately preserve the caller's spelling on unpinned routes, so it stays out. The Responses half (`mapRoutedResponsesReasoningEffort` for unconfigured providers) touches no existing assertion; carry it alone if focused tests show no other behavior change. | +| #5800 (Claude Agent SDK harness) | Leave | Adds a runtime dependency (`package.json`, `bun.lock`) and replaces a provider; needs dependency and maintainer security review beyond a bug train. | +| #5782 (Windows manual stops + bridge) | Leave | 9.3k lines, two feature lines, conflicting with `dev`. | +| #4222 (side-chat cache) | Leave | Experimental feature behind a setting, not a bug fix. | +| #5920, #5921 (desktop quota rows, widget reload) | Leave | Opened by the owner account minutes ago with their own devlog plan; another task owns them. | + +## Next cycles + +- wp3 — batch 7: `codex/bug-train-7` from `dev`; carry #5916 plus the ceiling fix; decide #5497 and the narrowed #5539 at P. Security review of the + final diff by an independent reviewer before merge. Exact-head CI, then `--admin --match-head-commit`. +- Closing the loop: re-query open bug PRs and confirm every one left open has a row above. From 1102a5f06726e1ae343c15e7baa4387bc6fec242 Mon Sep 17 00:00:00 2001 From: codingbo Date: Sat, 26 Sep 2026 19:23:08 +0900 Subject: [PATCH 2/3] fix(server): balance multi-account OAuth 429 retries with send budget (#5916) Carried from #5916 as one squashed commit. Closes #5880 once on dev. Co-authored-by: codingbo --- .../docs/reference/configuration/providers.md | 14 +- scripts/test-layout/layout.json | 1 + src/adapters/google-http.ts | 8 +- src/server/inference/context.ts | 27 +- src/server/responses/adapter-continuation.ts | 4 +- src/server/responses/adapter-dispatch.ts | 4 +- src/server/responses/passthrough-dispatch.ts | 4 +- src/server/responses/request-send-budget.ts | 23 +- src/server/responses/request-transport.ts | 14 + src/server/responses/run-turn-execution.ts | 4 +- src/server/responses/sidecar-execution.ts | 4 +- structure/transports/inventory.md | 2 +- structure/transports/responses-failover.md | 15 +- structure/transports/responses-spend.md | 13 + tests/fixtures/test-layout-expected.json | 1 + tests/oauth/generic-oauth-failover.test.ts | 2 +- tests/oauth/oauth-account-attribution.test.ts | 2 +- tests/server/inference-send-budget.test.ts | 133 +++++++- ...oogle-antigravity-oauth-429-budget.test.ts | 304 ++++++++++++++++++ 19 files changed, 545 insertions(+), 34 deletions(-) create mode 100644 tests/server/server-google-antigravity-oauth-429-budget.test.ts diff --git a/docs-site/src/content/docs/reference/configuration/providers.md b/docs-site/src/content/docs/reference/configuration/providers.md index e7e154cf659..ae5415a84d8 100644 --- a/docs-site/src/content/docs/reference/configuration/providers.md +++ b/docs-site/src/content/docs/reference/configuration/providers.md @@ -865,10 +865,16 @@ The Codex pool and the Anthropic pool are excluded and keep their own rotation; changes neither. A provider with a single stored account is a strict no-op, and no cooldown is recorded for it. -On a 429 the failed account is cooled using `Retry-After` when present (capped at 15 minutes) -or a default backoff, and the request is replayed on the next eligible account, up to three -rotations per request. An account flagged for reauthentication is never selected. Cooldowns are -process-local, so a restart forgets them. +Before dispatch, generic OAuth snapshots the eligible roster. On a 429 the failed account is cooled +using `Retry-After` when present (capped at 15 minutes) or a default backoff, and the request is +replayed on the next account selected from the live roster. The stable rotation ceiling is +`max(3, eligibleCount - 1)` per request; live selection still filters cooldowns, and an account +flagged for reauthentication is never selected. Cooldowns are process-local, so a restart forgets +them. + +When at least two accounts are eligible, the ingress-owned default send allowance covers up to three +sends per eligible account. A single eligible account keeps the existing base allowance of three and +total allowance of four. Explicit caller ceilings and combo scopes keep their existing limits. Rotation carries the alternate account's **full** credential snapshot, not just its bearer, so a provider that pairs routing metadata with its token — Antigravity's Cloud Code Assist project id, diff --git a/scripts/test-layout/layout.json b/scripts/test-layout/layout.json index 2eae1c41a3a..c10d490a8e3 100644 --- a/scripts/test-layout/layout.json +++ b/scripts/test-layout/layout.json @@ -1634,6 +1634,7 @@ "server-combo-reasoning-replay-eligibility.test.ts": "server", "server-combo-zero-output-failover.test.ts": "server", "server-google-antigravity-oauth-401-replay.test.ts": "server", + "server-google-antigravity-oauth-429-budget.test.ts": "server", "server-gui-bundle-freshness.test.ts": "server", "server-images-bodyless-content-length.test.ts": "server", "server-images.test.ts": "server", diff --git a/src/adapters/google-http.ts b/src/adapters/google-http.ts index bbebb36cb61..8d1c5df78e5 100644 --- a/src/adapters/google-http.ts +++ b/src/adapters/google-http.ts @@ -14,10 +14,10 @@ import { retryBackoffDelayMs, sleepWithAbort, SendBudgetExhaustedError, + TRANSIENT_RETRY_MAX_ATTEMPTS, isConnectionResetError, } from "../lib/upstream-retry"; -const GOOGLE_RETRY_ATTEMPTS = 3; const GOOGLE_RETRY_BASE_MS = 250; const GOOGLE_RETRY_MAX_MS = 2_000; @@ -60,7 +60,7 @@ export async function fetchGoogleWithRetry( let retryDelayMs = 0; let sendClass: SendClass = "transient"; let recovery: AttemptRecoveryKind | undefined; - for (let attempt = 0; attempt < GOOGLE_RETRY_ATTEMPTS; attempt++) { + for (let attempt = 0; attempt < TRANSIENT_RETRY_MAX_ATTEMPTS; attempt++) { if (ctx.abortSignal?.aborted) throw abortError(ctx.abortSignal); try { const res = await send({ url: activeRequest.url, sendClass, recovery, @@ -105,7 +105,7 @@ export async function fetchGoogleWithRetry( continue; } } - if (!retryableGoogleStatus(res.status) || attempt === GOOGLE_RETRY_ATTEMPTS - 1) { + if (!retryableGoogleStatus(res.status) || attempt === TRANSIENT_RETRY_MAX_ATTEMPTS - 1) { return ctx.returnRawErrors ? res : normalizeFinalGoogleError(label, res, ctx.abortSignal); } // A 429 may be a transient rate limit (retry) or hard quota exhaustion (do NOT retry — @@ -141,7 +141,7 @@ export async function fetchGoogleWithRetry( throw err; } lastError = err; - if (attempt === GOOGLE_RETRY_ATTEMPTS - 1) throw err; + if (attempt === TRANSIENT_RETRY_MAX_ATTEMPTS - 1) throw err; sendClass = "transient"; recovery = isConnectionResetError(err) ? "connection-reset" : undefined; retryDelayMs = retryBackoffDelayMs(attempt, { diff --git a/src/server/inference/context.ts b/src/server/inference/context.ts index 904d2300f38..08d26ada00b 100644 --- a/src/server/inference/context.ts +++ b/src/server/inference/context.ts @@ -1,7 +1,15 @@ -import { createRequestExecutionBudget, type RequestExecutionBudget } from "../../lib/request-execution-budget"; +import { CODEX_TEXT_GUARDED_BUDGET_POLICY, createRequestExecutionBudget, type RequestExecutionBudget } from "../../lib/request-execution-budget"; +import { TRANSIENT_RETRY_MAX_ATTEMPTS, type TransientSendBudget } from "../../lib/upstream-retry"; import type { RequestLogContext } from "../request-log"; import { attachRequestSpendTracker } from "../responses/request-spend"; +// Only ingress owns an expandable policy. Exact caller budgets and combo-derived scopes +// never enter this map, even when their policy happens to equal the default profile. +const ingressPolicies = new WeakMap(); + /** * The one construction of an ingress-owned send budget: the default guarded policy, no logical * request id, and this request's spend tracker as the send observer. Attaching the tracker @@ -11,5 +19,20 @@ export function createInferenceSendBudget( req: Pick, logCtx: RequestLogContext, ): RequestExecutionBudget { - return createRequestExecutionBudget(undefined, undefined, attachRequestSpendTracker(req, logCtx)); + const policy = { ...CODEX_TEXT_GUARDED_BUDGET_POLICY }; + const budget = createRequestExecutionBudget(policy, undefined, attachRequestSpendTracker(req, logCtx)); + ingressPolicies.set(budget, policy); + return budget; +} + +/** Fund each account's normal transient ladder once, before the first physical send. */ +export function expandInferenceOAuthSendBudget(budget: TransientSendBudget | undefined, accounts: number): void { + if (!budget) return; + const policy = ingressPolicies.get(budget); + if (!policy || budget.used !== 0) return; + ingressPolicies.delete(budget); + if (accounts < 2 || !Number.isSafeInteger(accounts)) return; + const sends = accounts * TRANSIENT_RETRY_MAX_ATTEMPTS; + policy.baseSendAllowance = Math.max(policy.baseSendAllowance, sends); + policy.maxTotalModelSends = Math.max(policy.maxTotalModelSends, sends); } diff --git a/src/server/responses/adapter-continuation.ts b/src/server/responses/adapter-continuation.ts index 2ee13b076ee..c3d57e66294 100644 --- a/src/server/responses/adapter-continuation.ts +++ b/src/server/responses/adapter-continuation.ts @@ -39,7 +39,6 @@ import { formatAnthropicProviderForLog, } from "../../oauth/anthropic-routing"; import { - GENERIC_OAUTH_MAX_FAILOVERS_PER_REQUEST, hasEligibleGenericOAuthFailoverTarget, isGenericOAuthFailoverEnabled, rotateGenericOAuthAccountOn429, @@ -84,6 +83,7 @@ export function createAdapterContinuations( | "commitResolvedOAuthSelection" | "genericFailoverAccountId" | "genericFailovers" + | "genericFailoverLimit" | "applyFailoverSnapshot" | "replayOAuthCredentialSnapshot" | "noteRoutedAttemptSend" @@ -407,7 +407,7 @@ export function createAdapterContinuations( response.status === 429 && transportState.genericFailoverAccountId && !isNonReplayableResponse(response) - && transportState.genericFailovers < GENERIC_OAUTH_MAX_FAILOVERS_PER_REQUEST + && transportState.genericFailovers < transportState.genericFailoverLimit && isGenericOAuthFailoverEnabled(config, route.providerName) ) { // Intersection with the shared request budget. The continuation loop re-sends the diff --git a/src/server/responses/adapter-dispatch.ts b/src/server/responses/adapter-dispatch.ts index e6f0c742e85..abb20f9940f 100644 --- a/src/server/responses/adapter-dispatch.ts +++ b/src/server/responses/adapter-dispatch.ts @@ -54,7 +54,6 @@ import { formatAnthropicProviderForLog, } from "../../oauth/anthropic-routing"; import { - GENERIC_OAUTH_MAX_FAILOVERS_PER_REQUEST, isGenericOAuthFailoverEnabled, rotateGenericOAuthAccountOn429, failoverAccountSnapshot, @@ -124,6 +123,7 @@ export async function prepareAdapterExchange( | "commitResolvedOAuthSelection" | "genericFailoverAccountId" | "genericFailovers" + | "genericFailoverLimit" | "applyFailoverSnapshot" | "noteRoutedAttemptSend" >, @@ -856,7 +856,7 @@ export async function prepareAdapterExchange( while ( upstreamResponse.status === 429 && transportState.genericFailoverAccountId - && transportState.genericFailovers < GENERIC_OAUTH_MAX_FAILOVERS_PER_REQUEST + && transportState.genericFailovers < transportState.genericFailoverLimit && isGenericOAuthFailoverEnabled(config, route.providerName) ) { // Intersection with the shared request budget. This arm re-sends through diff --git a/src/server/responses/passthrough-dispatch.ts b/src/server/responses/passthrough-dispatch.ts index 333d5fef893..41bdbfe931d 100644 --- a/src/server/responses/passthrough-dispatch.ts +++ b/src/server/responses/passthrough-dispatch.ts @@ -138,7 +138,6 @@ import type { OAuthAccessSnapshot } from "../../oauth"; import { publicOAuthAuthenticationErrorMessage } from "../../oauth"; import { resolveCopilotApiBaseUrl } from "../../oauth/github-copilot"; import { - GENERIC_OAUTH_MAX_FAILOVERS_PER_REQUEST, hasEligibleGenericOAuthFailoverTarget, isGenericOAuthFailoverEnabled, rotateGenericOAuthAccountOn429, @@ -194,6 +193,7 @@ export async function preparePassthroughExchange( | "refreshResolvedOAuthSelection" | "replayOAuthCredentialSnapshot" | "genericFailovers" + | "genericFailoverLimit" | "applyFailoverSnapshot" | "noteRoutedAttemptSend" | "selectionIsCurrent" @@ -1319,7 +1319,7 @@ export async function preparePassthroughExchange( // have run and would cool down an account that refused nothing. && !isNonReplayableResponse(upstreamResponse) && transportState.genericFailoverAccountId - && transportState.genericFailovers < GENERIC_OAUTH_MAX_FAILOVERS_PER_REQUEST + && transportState.genericFailovers < transportState.genericFailoverLimit && isGenericOAuthFailoverEnabled(config, route.providerName) ) { // The roster cap above is one half of the bound; the request's shared budget is the diff --git a/src/server/responses/request-send-budget.ts b/src/server/responses/request-send-budget.ts index 33f9b3da212..4f27131b94d 100644 --- a/src/server/responses/request-send-budget.ts +++ b/src/server/responses/request-send-budget.ts @@ -200,7 +200,9 @@ export function createResponsesSendBudget( options: { allowFinalRecoveryReserve?: boolean } = {}, ): { attempts: number; permit?: SingleUseDispatchPermit } => { const base = remainingTransientSendBudget(cap); - if (base > 0) return { attempts: base }; + // The hop already paid for this leg's first send. Include it in the helper's total + // attempts without charging it again, or the final account loses one transient attempt. + if (base > 0) return { attempts: Math.min(cap, base + (pendingHopPermit ? 1 : 0)) }; // A provider-configured transient total is an exact physical-send ceiling. Once it is // exhausted, the request-wide recovery reserve must not silently widen it. The default stays // permissive so unconfigured providers retain the guarded profile's fourth recovery send. @@ -217,8 +219,8 @@ export function createResponsesSendBudget( /** * One credential hop of this logical request, admitted by the INTERSECTION of two bounds. * - * `GENERIC_OAUTH_MAX_FAILOVERS_PER_REQUEST` and `ANTHROPIC_POOL_MAX_FAILOVERS_PER_REQUEST` - * stay exactly as they are: they bound rotation within one credential roster. What neither + * The snapshotted generic OAuth roster cap and `ANTHROPIC_POOL_MAX_FAILOVERS_PER_REQUEST` + * bound rotation within one credential roster. What neither * can see is everything else this request already sent, so three hops layered on a spent * budget still reached upstream three more times. A hop now happens only when its own layer * cap AND the shared budget both permit it, and the smaller of the two wins. @@ -245,7 +247,13 @@ export function createResponsesSendBudget( countedExternally = false, ): { allowed: boolean; permit?: SingleUseDispatchPermit } => { if (!isRequestExecutionBudget(sendBudget)) return { allowed: true }; - const decision = sendBudget.reserveDispatch({ sendClass, targetKey, countedExternally }); + const decision = sendBudget.reserveDispatch({ + sendClass, + // A same-provider credential hop is not a model/endpoint transition. Its diagnostic + // label must not replace the physical target used by the adapter's next retry. + targetKey: sendClass === "auth-recovery" ? sendBudget.lastTargetKey ?? targetKey : targetKey, + countedExternally, + }); return decision.allowed ? { allowed: true, permit: decision.permit } : { allowed: false }; }; /** @@ -322,6 +330,13 @@ function adapterDispatchBudgetView( // already paid does not make an unsafe replay safe, so that check stays with the budget. if (intent.replaySafe !== false) { const hopPermit = hop.claimHopPermit(); + if (hopPermit && intent.targetKey !== budget.lastTargetKey) { + // The rotated credential can select a different regional endpoint. Replace the + // provisional booking synchronously so the actual destination obeys transition + // limits, while the physical send is still charged only once. + hopPermit.release(); + return budget.reserveDispatch({ ...intent, sendClass: hopPermit.sendClass }); + } // Confirmed here rather than in `use()`: the adapter reserves immediately before it // opens the transport, which is the same boundary the hop's own confirmation uses. // A permit some other leg already settled returns false, and this falls through to a diff --git a/src/server/responses/request-transport.ts b/src/server/responses/request-transport.ts index 2b302402f38..2ade1c7ebae 100644 --- a/src/server/responses/request-transport.ts +++ b/src/server/responses/request-transport.ts @@ -28,11 +28,15 @@ import { UnsupportedOAuthProviderError, } from "../../oauth"; import { + GENERIC_OAUTH_MAX_FAILOVERS_PER_REQUEST, + eligibleFailoverAccounts, forgetGenericFailoverRoster, isGenericFailoverProvider, preferredInitialAccount, noteGenericPoolSelection, } from "../../oauth/generic-account-failover"; +import { classifyModelFamilyForQuota } from "../../oauth/account-quota-rank"; +import { expandInferenceOAuthSendBudget } from "../inference/context"; import { stampOAuthAccountLabel, usesApiKeyAccount } from "../../providers/label"; import { resolveProviderTransport } from "../../providers/xai-transport"; import { resolveCopilotApiBaseUrl } from "../../oauth/github-copilot"; @@ -721,7 +725,17 @@ export async function prepareResponsesTransport( ); } + // Freeze the request ceiling before any 429 writes cooldowns. Selection still reads + // live eligibility on every hop; cooled accounts cannot shorten this request's allowance. + const genericRosterSize = genericFailoverAccountId ? new Set([ + genericFailoverAccountId, + ...eligibleFailoverAccounts(route.providerName, Date.now(), classifyModelFamilyForQuota(route.providerName, route.modelId)), + ]).size : 0; + const genericFailoverLimit = Math.max(GENERIC_OAUTH_MAX_FAILOVERS_PER_REQUEST, genericRosterSize - 1); + if (genericFailoverAccountId) expandInferenceOAuthSendBudget(options.sendBudget, genericRosterSize); + return { + genericFailoverLimit, isOAuth401ReplayProvider, get sentOAuthSnapshot(): OAuthAccessSnapshot | undefined { return sentOAuthSnapshot; diff --git a/src/server/responses/run-turn-execution.ts b/src/server/responses/run-turn-execution.ts index 762dbe3eec9..5659c0970b9 100644 --- a/src/server/responses/run-turn-execution.ts +++ b/src/server/responses/run-turn-execution.ts @@ -22,7 +22,6 @@ import { normalizeDeclaredToolName, type AdapterEvent, type OcxProviderContinuat import { adapterFailureFromMessage, SEND_BUDGET_EXHAUSTED_CODE } from "../../lib/errors"; import { SendBudgetExhaustedError, markResponseNonReplayable } from "../../lib/upstream-retry"; import { - GENERIC_OAUTH_MAX_FAILOVERS_PER_REQUEST, hasEligibleGenericOAuthFailoverTarget, isGenericOAuthFailoverEnabled, rotateGenericOAuthAccountOn429, @@ -92,6 +91,7 @@ export async function executeResponsesRunTurn( | "replayOAuthCredentialSnapshot" | "genericFailoverAccountId" | "genericFailovers" + | "genericFailoverLimit" | "applyFailoverSnapshot" | "resolveSelectionAdapter" | "adapter" @@ -348,7 +348,7 @@ export async function executeResponsesRunTurn( if ( status !== 429 || !transportState.genericFailoverAccountId - || transportState.genericFailovers >= GENERIC_OAUTH_MAX_FAILOVERS_PER_REQUEST + || transportState.genericFailovers >= transportState.genericFailoverLimit || !isGenericOAuthFailoverEnabled(config, route.providerName) ) return false; // Intersection with the request's shared budget: the roster bound above answers "may this diff --git a/src/server/responses/sidecar-execution.ts b/src/server/responses/sidecar-execution.ts index 24ffc9f126c..4c7318de6e3 100644 --- a/src/server/responses/sidecar-execution.ts +++ b/src/server/responses/sidecar-execution.ts @@ -20,7 +20,6 @@ import type { ProviderAdapter } from "../../adapters/base"; import type { OcxParsedRequest } from "../../types"; import { rotateProviderTransportOn429, rateLimitRetryPolicyFor } from "../../providers/key-failover"; import { - GENERIC_OAUTH_MAX_FAILOVERS_PER_REQUEST, isGenericOAuthFailoverEnabled, rotateGenericOAuthAccountOn429, failoverAccountSnapshot, @@ -59,6 +58,7 @@ export async function executeResponsesSidecars( | "adapter" | "genericFailoverAccountId" | "genericFailovers" + | "genericFailoverLimit" | "applyFailoverSnapshot" | "anthropicPoolAccountId" | "anthropicPoolFailovers" @@ -180,7 +180,7 @@ export async function executeResponsesSidecars( // excludes it), so its sidecar 429s died on this guard before the Anthropic arm below // could ever be considered. transportState.genericFailoverAccountId - && transportState.genericFailovers < GENERIC_OAUTH_MAX_FAILOVERS_PER_REQUEST + && transportState.genericFailovers < transportState.genericFailoverLimit && isGenericOAuthFailoverEnabled(config, route.providerName) ) { // Intersection with the request's shared budget. The sidecar replay is dispatched by the diff --git a/structure/transports/inventory.md b/structure/transports/inventory.md index 6e08a4b846e..a122e562990 100644 --- a/structure/transports/inventory.md +++ b/structure/transports/inventory.md @@ -30,7 +30,7 @@ surface is listed here so a maintainer can find the owner without grepping: | Azure OpenAI Responses | `src/adapters/azure.ts` | Deployment-shaped URLs on top of the Responses contract. | | Responses custom-tool preview | `src/bridge/sse.ts`, `src/server/responses-custom-tool-repair.ts`, `src/responses/progressive-freeform-input.ts`, `src/responses/freeform-wrapper-scan.ts` | Direct adapter events and routed function restoration share one progressive wrapper decoder over one bounded JSON classification of the prefix, so property order and escaped key spellings preview as the wrapper completion unwraps them. Fence-shaped `exec`/`apply_patch` prefixes stay held until authoritative completion normalization, while each caller retains its own patch-envelope and byte-budget policy. | | Meta Muse Responses tool names | `src/responses/muse-tool-name-alias.ts`, `src/adapters/openai-responses.ts` | `api.meta.ai` only: function names over 64 characters or containing characters outside `[a-zA-Z0-9_-]` become collision-safe wire aliases and are restored before the client sees them. | -| Google / Vertex / Antigravity | `src/adapters/google.ts`, `src/adapters/google-http.ts`, `src/adapters/google-wire-compiler.ts`, `src/adapters/google-tool-schema.ts`, `src/adapters/google-truncation.ts`, `src/adapters/google-errors.ts`, `src/adapters/google-antigravity-wire.ts`, `src/adapters/google-antigravity-replay.ts`, `src/adapters/google-wire-shape.ts` | Vertex and Antigravity install a Google-family `fetchResponse` and so own their retry policy, while AI Studio Gemini leaves it undefined and uses the default server fetch path. The Google-family wrapper reuses shared abort/deadline helpers, upstream error normalization, and policy-aware wire-body repair: strict initial schema loss sends nothing, while strict repair withholding returns the original 400 without a changed send. The final compiler produces the [content-free tool-schema loss contract](../providers/google.md#google-tool-schema-loss-reporting). `google-wire-shape.ts` remains diagnostic-only. | +| Google / Vertex / Antigravity | `src/adapters/google.ts`, `src/adapters/google-http.ts`, `src/adapters/google-wire-compiler.ts`, `src/adapters/google-tool-schema.ts`, `src/adapters/google-truncation.ts`, `src/adapters/google-errors.ts`, `src/adapters/google-antigravity-wire.ts`, `src/adapters/google-antigravity-replay.ts`, `src/adapters/google-wire-shape.ts` | Vertex and Antigravity install a Google-family `fetchResponse` and so own their retry policy, while AI Studio Gemini leaves it undefined and uses the default server fetch path. Its transient ladder uses the shared `TRANSIENT_RETRY_MAX_ATTEMPTS` constant, also used to fund each eligible generic OAuth account in the [request send budget](responses-spend.md#credential-hop-reservations). The Google-family wrapper reuses shared abort/deadline helpers, upstream error normalization, and policy-aware wire-body repair: strict initial schema loss sends nothing, while strict repair withholding returns the original 400 without a changed send. The final compiler produces the [content-free tool-schema loss contract](../providers/google.md#google-tool-schema-loss-reporting). `google-wire-shape.ts` remains diagnostic-only. | | Mimo Free | `src/adapters/mimo-free.ts` | Client identity and JWT handling are transport-local; the per-install client id lives in the opencodex state root. | | Anthropic image ingress | `src/adapters/anthropic-image-guard.ts`, `src/adapters/anthropic-image-normalize.ts`, `src/adapters/anthropic-image-codec.ts` | Oversized or unsupported images are normalized or rejected before reaching upstream. An image's ladder position is pinned to its own fixed-size digest of content and canonical media type rather than recomputed from recency each request (#4532); appending a newer image therefore cannot demote and re-encode older images and bust Anthropic's prompt prefix cache. Position bytes count toward the shared retained-memory budget but are pinned — the shared evictor clears normalization cache slots first and never drops a position mid-request; positions shrink only through the store's own entry-count cap at request end. Unseen images still take the age-tier pyramid's first position, the total byte budget still binds, and a 413 `tierBias` retry still applies. Recorded positions only move down the ladder, so the store is monotonic. | | Adapter execution support | `src/adapters/run-turn-queue.ts`, `src/adapters/tool-catalog-nudge.ts`, `src/adapters/identity.ts`, `src/adapters/image.ts`, `src/adapters/upstream-http-error.ts` | Shared machinery: turn ordering, tool-catalog nudging, client fingerprinting, image conversion, upstream error normalization. | diff --git a/structure/transports/responses-failover.md b/structure/transports/responses-failover.md index 76c41755933..9bb15c44f5a 100644 --- a/structure/transports/responses-failover.md +++ b/structure/transports/responses-failover.md @@ -474,10 +474,11 @@ Dashboard Fast-row persistence and client refresh follow the [Fast selector rows ## Account refusal and rotation boundaries -Native Responses uses the existing pre-stream OAuth HTTP-429 account rotation: account quorum, -cooldown and the three-rotation request cap remain in force, the complete credential/transport/replay -identity is refreshed, and usage is attributed to the serving account. Single-account installs do not -retry; a missing alternate credential preserves the original error. Organization or project exhaustion +Native Responses uses the existing pre-stream OAuth HTTP-429 account rotation: account quorum and +cooldown remain in force, while generic OAuth uses the stable snapshot ceiling described below. The +complete credential/transport/replay identity is refreshed, and usage is attributed to the serving +account. Single-account installs do not rotate; a missing alternate credential preserves the original +error while transient recovery remains available. Organization or project exhaustion allows an initial alternate attempt because the response does not identify the refusing scope. After resolving an alternate, organization-level retry is withheld only when both credentials have the same known workspace account id. Stored Pool/main-pool alternates supply that id directly; a request-owned @@ -493,6 +494,12 @@ Send-budget refusal is attributed as a withheld rotation only when a model-famil check confirms from the live roster that at least two accounts exist and an alternate is not currently cooled. That check applies no cooldown and advances no rotation. +Generic OAuth snapshots its eligible roster before dispatch. Its request rotation ceiling is +`max(3, eligibleCount - 1)`; the live picker still filters cooldowns, so the snapshot supplies the +stable ceiling without making a cooled account eligible. Same-provider auth recovery keeps the last +physical target, rather than a diagnostic key, and a real send is charged once even when recovery +rebuilds the request. + Precommit Codex model refusals use bounded account recovery for HTTP `detail` and WebSocket-projected `error.message` bodies. Only an exact HTTP 400 refusal naming the requested or wire model establishes denial evidence; ordinary malformed requests and committed stream errors do not authorize another diff --git a/structure/transports/responses-spend.md b/structure/transports/responses-spend.md index c749375a83b..b4a67211d19 100644 --- a/structure/transports/responses-spend.md +++ b/structure/transports/responses-spend.md @@ -11,6 +11,19 @@ budget before it knows whether a rotation is even possible, because the reservat `reserveDispatch` spends, `permit.use()` only confirms which leg sent, and `permit.release()` is idempotent and a no-op once used. Every ladder therefore owes the budget an answer on every exit. +Generic OAuth snapshots the eligible roster at request ingress before dispatch. Its rotation ceiling +is `max(3, eligibleCount - 1)`; live selection still removes accounts in cooldown, so the snapshot +sets the number of possible moves without making a cooled account selectable. Only the ingress-owned +default execution budget expands when at least two accounts are eligible: its base and total ceilings +cover up to `TRANSIENT_RETRY_MAX_ATTEMPTS` sends per eligible account (currently three). A single +eligible account keeps the existing base ceiling of three and total ceiling of four. Explicit caller +ceilings and combo-derived scopes keep their existing limits. + +A helper recovery's prepaid hop remains part of the full leg attempt allowance. An adapter-owned +pending hop reconciles the actual target synchronously: a regional endpoint change refunds the old +reservation and re-reserves the new transition, so the transition cap is enforced without charging +or sending twice. Diagnostic key changes alone are not transitions; a real destination change is. + The hop pays for a replay that some *other* layer dispatches, so which layer settles the reservation follows the dispatcher, not the ladder. A helper-routed replay reports the same physical send back through `onSendsConsumed`; that is what `countedExternally: true` names, and the diff --git a/tests/fixtures/test-layout-expected.json b/tests/fixtures/test-layout-expected.json index d5272cbc96e..22679a1b04f 100644 --- a/tests/fixtures/test-layout-expected.json +++ b/tests/fixtures/test-layout-expected.json @@ -1460,6 +1460,7 @@ "server-combo-reasoning-replay-eligibility.test.ts": "server", "server-combo-zero-output-failover.test.ts": "server", "server-google-antigravity-oauth-401-replay.test.ts": "server", + "server-google-antigravity-oauth-429-budget.test.ts": "server", "server-gui-bundle-freshness.test.ts": "server", "server-images-bodyless-content-length.test.ts": "server", "server-images.test.ts": "server", diff --git a/tests/oauth/generic-oauth-failover.test.ts b/tests/oauth/generic-oauth-failover.test.ts index 2260a8248ea..2384d971973 100644 --- a/tests/oauth/generic-oauth-failover.test.ts +++ b/tests/oauth/generic-oauth-failover.test.ts @@ -352,7 +352,7 @@ describe("sidecar on429 wiring", () => { // The gate is a POSITIVE else-if, not an early return: an early bare return here made the // Anthropic arm below unreachable, because Anthropic never has a genericFailoverAccountId. expect(body).toContain("genericFailoverAccountId"); - expect(body).toContain("genericFailovers < GENERIC_OAUTH_MAX_FAILOVERS_PER_REQUEST"); + expect(body).toContain("genericFailovers < transportState.genericFailoverLimit"); expect(body).toContain("isGenericOAuthFailoverEnabled(config, route.providerName)"); // Anthropic's pool is excluded from generic failover, so it needs its own arm here or a 429 diff --git a/tests/oauth/oauth-account-attribution.test.ts b/tests/oauth/oauth-account-attribution.test.ts index 68b444a983e..05fed3917b0 100644 --- a/tests/oauth/oauth-account-attribution.test.ts +++ b/tests/oauth/oauth-account-attribution.test.ts @@ -263,7 +263,7 @@ describe("Responses per-account attribution for non-Codex OAuth", () => { }); }); - test.each([[1, 1], [5, 4]])("native Responses with %i accounts stays within %i sends on repeated 429", async (accounts, expectedSends) => { + test.each([[1, 1], [5, 5]])("native Responses with %i accounts visits each once on repeated 429 (%i sends)", async (accounts, expectedSends) => { await withHome(async () => { clearGenericFailoverHealth(); for (let index = 0; index < accounts; index++) { diff --git a/tests/server/inference-send-budget.test.ts b/tests/server/inference-send-budget.test.ts index e53bdfd4a06..0e3ca1597a9 100644 --- a/tests/server/inference-send-budget.test.ts +++ b/tests/server/inference-send-budget.test.ts @@ -1,14 +1,103 @@ -import { describe, expect, test } from "bun:test"; +import { afterEach, beforeEach, describe, expect, test } from "bun:test"; +import { mkdtempSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { removeTreeWithRetry } from "../helpers/remove-tree"; +import { acquireOwnedSpendHome } from "../helpers/owned-spend-home"; import { CODEX_TEXT_GUARDED_BUDGET_POLICY } from "../../src/lib/request-execution-budget"; -import { createInferenceSendBudget } from "../../src/server/inference/context"; +import { createInferenceSendBudget, expandInferenceOAuthSendBudget } from "../../src/server/inference/context"; import type { RequestLogContext } from "../../src/server/request-log"; +import { createRequestExecutionBudget, deriveRequestExecutionBudget } from "../../src/lib/request-execution-budget"; +import { fetchWithTransientRetry, TRANSIENT_RETRY_MAX_ATTEMPTS } from "../../src/lib/upstream-retry"; +import { budgetOwner } from "../helpers/send-budget-owner"; + +let home: string; +let originalHome: string | undefined; +let releaseHome: () => void; +beforeEach(() => { + originalHome = process.env.OPENCODEX_HOME; + home = mkdtempSync(join(tmpdir(), "ocx-inference-budget-")); + process.env.OPENCODEX_HOME = home; + releaseHome = acquireOwnedSpendHome(); +}); +afterEach(() => { + releaseHome(); + if (originalHome === undefined) delete process.env.OPENCODEX_HOME; + else process.env.OPENCODEX_HOME = originalHome; + removeTreeWithRetry(home); +}); describe("createInferenceSendBudget", () => { + test.each([0, 1])("an adapter credential hop enforces a real endpoint transition (limit=%i)", limit => { + const budget = createRequestExecutionBudget({ + ...CODEX_TEXT_GUARDED_BUDGET_POLICY, + maxTargetTransitions: limit, + maxAlternateTargetSends: limit, + }); + const { owner, dispose } = budgetOwner(budget); + try { + const initial = owner.adapterDispatchBudget!.reserveDispatch({ sendClass: "initial", targetKey: "https://region-a.example/" }); + if (initial.allowed) initial.permit.use(); + const hop = owner.reserveCredentialHop("auth-recovery", "provider|model|oauth-429", true); + owner.pendingHopPermit = hop.permit; + const replay = owner.adapterDispatchBudget!.reserveDispatch({ sendClass: "transient", targetKey: "https://region-b.example/" }); + expect(replay.allowed).toBe(limit === 1); + if (replay.allowed) replay.permit.use(); + expect(budget.used).toBe(1 + limit); + expect(budget.targetTransitions).toBe(limit); + expect(budget.alternateTargetSends).toBe(limit); + expect(budget.lastTargetKey).toBe(limit ? "https://region-b.example/" : "https://region-a.example/"); + } finally { dispose(); } + }); + + test("the last helper-driven account gets its prepaid send plus remaining retries", async () => { + const budget = createInferenceSendBudget(new Request("http://localhost/v1/responses"), { model: "m", provider: "p" }); + expandInferenceOAuthSendBudget(budget, 4); + const { owner, dispose } = budgetOwner(budget); + try { + owner.noteTransientSends(9); + const hop = owner.reserveCredentialHop("auth-recovery", "provider|model|oauth-429", true); + owner.pendingHopPermit = hop.permit; + const allowance = owner.recoverySendAllowance(3, "auth-recovery", "provider|model|oauth-429"); + let sends = 0; + const response = await fetchWithTransientRetry(async () => { + if (sends === 0) hop.permit?.use(); + sends++; + return new Response("", { status: sends === 3 ? 200 : 503 }); + }, { attempts: allowance.attempts, onSendsConsumed: owner.noteTransientSends }); + expect(response.status).toBe(200); + expect(sends).toBe(3); + expect(budget.used).toBe(12); + expect(budget.reserveDispatch({ sendClass: "auth-recovery", targetKey: "provider|model|oauth-429" }).allowed).toBe(false); + } finally { dispose(); } + }); + + test("same-provider auth recovery keeps the adapter's physical target", () => { + const budget = createRequestExecutionBudget(); + const { owner, dispose } = budgetOwner(budget); + try { + const targetKey = "https://daily-cloudcode-pa.googleapis.com/v1internal:generateContent"; + const initial = owner.adapterDispatchBudget!.reserveDispatch({ sendClass: "initial", targetKey }); + expect(initial.allowed).toBe(true); + if (initial.allowed) initial.permit.use(); + const hop = owner.reserveCredentialHop("auth-recovery", "google-antigravity|model|oauth-429"); + expect(hop.allowed).toBe(true); + owner.pendingHopPermit = hop.permit; + const replay = owner.adapterDispatchBudget!.reserveDispatch({ sendClass: "initial", targetKey }); + expect(replay.allowed).toBe(true); + if (replay.allowed) replay.permit.use(); + const retry = owner.adapterDispatchBudget!.reserveDispatch({ sendClass: "transient", targetKey }); + expect(retry.allowed).toBe(true); + expect(budget.lastTargetKey).toBe(targetKey); + expect(budget.targetTransitions).toBe(0); + expect(budget.used).toBe(3); + } finally { dispose(); } + }); test("mints a default-policy holder and parks this request's spend tracker on the log", () => { const logCtx: RequestLogContext = { model: "m", provider: "p" }; const req = new Request("http://localhost/v1/responses", { method: "POST" }); const budget = createInferenceSendBudget(req, logCtx); - expect(budget.policy).toBe(CODEX_TEXT_GUARDED_BUDGET_POLICY); + expect(budget.policy).toEqual(CODEX_TEXT_GUARDED_BUDGET_POLICY); expect(typeof budget.logicalRequestId).toBe("string"); expect(budget.used).toBe(0); expect(logCtx.spendTracker).toBeDefined(); @@ -21,4 +110,42 @@ describe("createInferenceSendBudget", () => { expect(a).not.toBe(b); expect(a.logicalRequestId).not.toBe(b.logicalRequestId); }); + + test("roster expansion is bounded and does not change another request or explicit scopes", () => { + const req = new Request("http://localhost/v1/responses"); + const budget = createInferenceSendBudget(req, { model: "m", provider: "p" }); + const other = createInferenceSendBudget(req, { model: "m", provider: "p" }); + const exact = createRequestExecutionBudget(); + const child = deriveRequestExecutionBudget(budget, { ...budget.policy }); + expandInferenceOAuthSendBudget(budget, 4); + expandInferenceOAuthSendBudget(exact, 4); + expandInferenceOAuthSendBudget(child, 4); + expect(budget.policy.baseSendAllowance).toBe(4 * TRANSIENT_RETRY_MAX_ATTEMPTS); + expect(budget.policy.maxTotalModelSends).toBe(4 * TRANSIENT_RETRY_MAX_ATTEMPTS); + for (const unexpanded of [other, exact, child]) { + expect(unexpanded.policy).toEqual(CODEX_TEXT_GUARDED_BUDGET_POLICY); + } + // Later cooldowns or logins cannot shrink or replenish this request's snapshot. + expandInferenceOAuthSendBudget(budget, 2); + expandInferenceOAuthSendBudget(budget, 8); + for (let i = 0; i < 4 * TRANSIENT_RETRY_MAX_ATTEMPTS; i++) { + const send = budget.reserveDispatch({ sendClass: "transient", targetKey: "physical-url" }); + expect(send.allowed).toBe(true); + if (send.allowed) send.permit.use(); + } + expect(budget.reserveDispatch({ sendClass: "auth-recovery", targetKey: "physical-url" }).allowed).toBe(false); + expect(child.used).toBe(budget.used); + }); + + test("one credential retains the default and a spent request cannot expand", () => { + const req = new Request("http://localhost/v1/responses"); + const single = createInferenceSendBudget(req, { model: "m", provider: "p" }); + expandInferenceOAuthSendBudget(single, 1); + expect(single.policy).toEqual(CODEX_TEXT_GUARDED_BUDGET_POLICY); + const started = createInferenceSendBudget(req, { model: "m", provider: "p" }); + const send = started.reserveDispatch({ sendClass: "initial", targetKey: "physical-url" }); + if (send.allowed) send.permit.use(); + expandInferenceOAuthSendBudget(started, 4); + expect(started.policy).toEqual(CODEX_TEXT_GUARDED_BUDGET_POLICY); + }); }); diff --git a/tests/server/server-google-antigravity-oauth-429-budget.test.ts b/tests/server/server-google-antigravity-oauth-429-budget.test.ts new file mode 100644 index 00000000000..58c98264275 --- /dev/null +++ b/tests/server/server-google-antigravity-oauth-429-budget.test.ts @@ -0,0 +1,304 @@ +import { afterEach, beforeEach, describe, expect, spyOn, test } from "bun:test"; +import { mkdtempSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { saveConfig } from "../../src/config"; +import { clearGenericFailoverHealth } from "../../src/oauth/generic-account-failover"; +import { getAccountSet, saveCredential, setActiveAccount } from "../../src/oauth/store"; +import { handleResponses } from "../../src/server/responses"; +import type { OcxConfig } from "../../src/types"; +import { installIsolatedCodexHome, type IsolatedCodexHome } from "../helpers/isolated-codex-home"; +import { acquireOwnedSpendHome } from "../helpers/owned-spend-home"; +import { removeTreeWithRetry } from "../helpers/remove-tree"; +import * as retry from "../../src/lib/upstream-retry"; +import { createRequestExecutionBudget, type RequestExecutionBudgetPolicy } from "../../src/lib/request-execution-budget"; + +const DAILY_API_BASE = "https://daily-cloudcode-pa.googleapis.com"; + +let testDir = ""; +let previousHome: string | undefined; +let isolatedCodexHome: IsolatedCodexHome | null = null; +let originalFetch: typeof fetch; +let releaseSpendHome: (() => void) | undefined; +let sleepSpy: ReturnType | undefined; + +const takeSpendHome = (): void => { + releaseSpendHome ??= acquireOwnedSpendHome(); +}; + +beforeEach(() => { + originalFetch = globalThis.fetch; + previousHome = process.env.OPENCODEX_HOME; + isolatedCodexHome = installIsolatedCodexHome("ocx-google-429-codex-"); + testDir = mkdtempSync(join(tmpdir(), "ocx-google-429-")); + process.env.OPENCODEX_HOME = testDir; + takeSpendHome(); + clearGenericFailoverHealth(); + sleepSpy = spyOn(retry, "sleepWithAbort").mockImplementation(async () => {}); +}); + +afterEach(() => { + sleepSpy?.mockRestore(); + sleepSpy = undefined; + releaseSpendHome?.(); + releaseSpendHome = undefined; + clearGenericFailoverHealth(); + globalThis.fetch = originalFetch; + if (previousHome === undefined) delete process.env.OPENCODEX_HOME; + else process.env.OPENCODEX_HOME = previousHome; + isolatedCodexHome?.restore(); + isolatedCodexHome = null; + if (testDir) removeTreeWithRetry(testDir); +}); + +function antigravityConfig(): OcxConfig { + return { + port: 0, + hostname: "127.0.0.1", + defaultProvider: "google-antigravity", + providers: { + "google-antigravity": { + adapter: "google", + baseUrl: DAILY_API_BASE, + authMode: "oauth", + googleMode: "cloud-code-assist", + project: "initial-project-id", + models: ["gemini-3.8-flash"], + }, + }, + } as OcxConfig; +} + +function jsonSuccessBody(text: string): Record { + return { + response: { + candidates: [{ + content: { + role: "model", + parts: [{ text }], + }, + finishReason: "STOP", + }], + usageMetadata: { + promptTokenCount: 5, + candidatesTokenCount: 3, + totalTokenCount: 8, + }, + }, + }; +} + +function transient429ErrorBody(): Record { + return { + error: { + code: 429, + message: "Resource has been exhausted: rate limit exceeded.", + status: "RESOURCE_EXHAUSTED", + }, + }; +} + +function hardQuota429ErrorBody(): Record { + return { + error: { + code: 429, + message: "Quota exceeded for quota metric ...", + status: "RESOURCE_EXHAUSTED", + }, + }; +} + +function createResponsesRequest(bodyOverrides: Record = {}): Request { + return new Request("http://127.0.0.1/v1/responses", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ + model: "google-antigravity/gemini-3.8-flash", + input: "hello", + stream: false, + ...bodyOverrides, + }), + }); +} + +async function seedAntigravityAccounts(count: number): Promise> { + for (let i = 1; i <= count; i++) { + await saveCredential("google-antigravity", { + access: `token-${i}`, + refresh: `refresh-${i}`, + expires: Date.now() + 3_600_000, + accountId: `account-${i}`, + projectId: `project-${i}`, + }); + } + const accounts = getAccountSet("google-antigravity")!.accounts; + await setActiveAccount("google-antigravity", accounts[0]!.id); + return accounts.map((a, idx) => ({ + id: a.id, + auth: `Bearer token-${idx + 1}`, + project: `project-${idx + 1}`, + })); +} + +function installAntigravityFetchMock( + handler: (info: { auth: string; project: string; sendIndex: number }) => Response | Promise, +): void { + let sendIndex = 0; + globalThis.fetch = (async (input, init) => { + const url = input instanceof Request ? input.url : String(input); + const parsedUrl = new URL(url); + + if (parsedUrl.origin === DAILY_API_BASE + && ["/v1internal:streamGenerateContent", "/v1internal:generateContent"].includes(parsedUrl.pathname)) { + sendIndex += 1; + const auth = new Headers(init?.headers).get("authorization") ?? ""; + let project = ""; + if (typeof init?.body === "string") { + try { + const parsed = JSON.parse(init.body) as { project?: string }; + project = parsed.project ?? ""; + } catch { /* ignore */ } + } + return handler({ auth, project, sendIndex }); + } + + if (parsedUrl.hostname === "127.0.0.1" || parsedUrl.hostname === "localhost") return originalFetch(input, init); + throw new Error(`Unexpected external request: ${url}`); + }) as typeof fetch; +} + +describe("Google Antigravity OAuth 429 retry and multi-account budget (#5880)", () => { + test.each([4, 5])("%i accounts each receive three transient sends before terminal 429", async accountCount => { + const accounts = await seedAntigravityAccounts(accountCount); + const cfg = antigravityConfig(); + saveConfig(cfg); + + const observedSends: Array<{ auth: string; project: string }> = []; + installAntigravityFetchMock(({ auth, project }) => { + observedSends.push({ auth, project }); + return new Response(JSON.stringify(transient429ErrorBody()), { + status: 429, + headers: { "content-type": "application/json" }, + }); + }); + + const res = await handleResponses( + createResponsesRequest(), + cfg, + { model: "gemini-3.8-flash", provider: "google-antigravity" }, + ); + + expect(res.status).toBe(429); + + expect(await res.text()).toContain("rate_limit_exceeded"); + // Four accounts pin 3,3,3,3; five also proves the snapshot can exceed the old hop cap. + expect(observedSends).toHaveLength(accountCount * 3); + + for (let acctIdx = 0; acctIdx < accountCount; acctIdx++) { + const sendsForAccount = observedSends.slice(acctIdx * 3, (acctIdx + 1) * 3); + expect(sendsForAccount).toHaveLength(3); + for (const send of sendsForAccount) { + expect(send).toEqual({ auth: accounts[acctIdx]!.auth, project: accounts[acctIdx]!.project }); + } + } + }); + + test.each([2, 3])("single account succeeds attempt %i on transient 429", async successAttempt => { + const accounts = await seedAntigravityAccounts(1); + const cfg = antigravityConfig(); + saveConfig(cfg); + + const observedSends: Array<{ auth: string; project: string }> = []; + installAntigravityFetchMock(({ auth, project, sendIndex }) => { + observedSends.push({ auth, project }); + if (sendIndex < successAttempt) { + return new Response(JSON.stringify(transient429ErrorBody()), { + status: 429, + headers: { "content-type": "application/json" }, + }); + } + return new Response(JSON.stringify(jsonSuccessBody(`success on attempt ${successAttempt}`)), { + status: 200, + headers: { "content-type": "application/json" }, + }); + }); + + const res = await handleResponses( + createResponsesRequest(), + cfg, + { model: "gemini-3.8-flash", provider: "google-antigravity" }, + ); + + expect(res.status).toBe(200); + const body = await res.json() as Record; + expect(JSON.stringify(body)).toContain(`success on attempt ${successAttempt}`); + expect(observedSends).toHaveLength(successAttempt); + for (const send of observedSends) { + expect(send).toEqual({ auth: accounts[0]!.auth, project: accounts[0]!.project }); + } + }); + + test.each([1, 4])("hard quota with %i accounts does not waste transient retries", async accountCount => { + const accounts = await seedAntigravityAccounts(accountCount); + const cfg = antigravityConfig(); + saveConfig(cfg); + + const observedSends: Array<{ auth: string; project: string }> = []; + installAntigravityFetchMock(({ auth, project }) => { + observedSends.push({ auth, project }); + return new Response(JSON.stringify(hardQuota429ErrorBody()), { + status: 429, + headers: { "content-type": "application/json" }, + }); + }); + + const res = await handleResponses( + createResponsesRequest(), + cfg, + { model: "gemini-3.8-flash", provider: "google-antigravity" }, + ); + + expect(res.status).toBe(429); + // Hard quota must return immediately without burning transient retry attempts + expect(observedSends).toEqual(accounts.map(({ auth, project }) => ({ auth, project }))); + await res.text(); + }); + + test.each([2, 4, 7])("explicit %i-send caller ceiling remains unchanged", async ceiling => { + await seedAntigravityAccounts(4); + const cfg = antigravityConfig(); + saveConfig(cfg); + + const observedSends: Array<{ auth: string; project: string }> = []; + installAntigravityFetchMock(({ auth, project }) => { + observedSends.push({ auth, project }); + return new Response(JSON.stringify(transient429ErrorBody()), { + status: 429, + headers: { "content-type": "application/json" }, + }); + }); + + const customPolicy: RequestExecutionBudgetPolicy = { + maxTotalModelSends: ceiling, + baseSendAllowance: ceiling, + finalRecoveryAllowance: 0, + maxAlternateTargetSends: 0, + maxTargetTransitions: 0, + }; + const customBudget = createRequestExecutionBudget(customPolicy); + + const res = await handleResponses( + createResponsesRequest(), + cfg, + { model: "gemini-3.8-flash", provider: "google-antigravity" }, + { sendBudget: customBudget }, + ); + + expect(res.status).toBe(429); + expect(observedSends).toHaveLength(ceiling); + expect(customBudget.used).toBe(ceiling); + expect(customBudget.policy).toBe(customPolicy); + expect(customBudget.targetTransitions).toBe(0); + await res.text(); + }); +}); From 32be95be753836c6c7885663b41d146b52d9da5e Mon Sep 17 00:00:00 2001 From: JUN Date: Sat, 26 Sep 2026 19:24:03 +0900 Subject: [PATCH 3/3] fix(oauth): cap one request at six funded accounts in the OAuth 429 budget The #5916 allowance funded three sends for every eligible account and raised the hop limit to roster size minus one, with no fixed ceiling: a large roster let one request make 3 x N upstream sends and push a 429 storm across the pool. GENERIC_OAUTH_MAX_ACCOUNTS_PER_REQUEST (6) now clamps the roster snapshot in request-transport.ts and again inside expandInferenceOAuthSendBudget, so the default ingress ceiling is 18 sends whatever the roster size. Explicit caller budgets are unchanged. New case: eight accounts make exactly 18 sends over the first six (24 without the clamp). Docs and structure state the cap; the batch plan is added. --- devlog/_plan/260926_bug_train_6/020_batch7.md | 45 +++++++++++++++++++ .../docs/reference/configuration/providers.md | 5 ++- src/oauth/generic-account-failover.ts | 8 ++++ src/server/inference/context.ts | 9 +++- src/server/responses/request-transport.ts | 7 ++- structure/transports/responses-failover.md | 2 +- structure/transports/responses-spend.md | 6 ++- ...oogle-antigravity-oauth-429-budget.test.ts | 29 ++++++++++++ 8 files changed, 102 insertions(+), 9 deletions(-) create mode 100644 devlog/_plan/260926_bug_train_6/020_batch7.md diff --git a/devlog/_plan/260926_bug_train_6/020_batch7.md b/devlog/_plan/260926_bug_train_6/020_batch7.md new file mode 100644 index 00000000000..3f0a3772415 --- /dev/null +++ b/devlog/_plan/260926_bug_train_6/020_batch7.md @@ -0,0 +1,45 @@ +# Batch 7 — plan + +Previous D (wp2): roadmap locked in `010_roadmap.md`; direction: batch 7 carries #5916 with a fixed per-request ceiling +and decides #5497 and a narrowed #5539 here. + +## P decisions on the open rows + +- #5497 — leave. It adds a public provider field (`responseTierAuthoritative`) with docs in three locales and new cost + and usage-log handling (the canonical ChatGPT exception is kept); that is a product/config review, not a merge-train carry. +- #5539 (Responses half) — leave. For an unconfigured `openai-responses` provider (API-key OpenAI included) it would fold + `minimal` to `low`, which changes requests for OpenAI models that accept `minimal`. The strict-gateway 400 is + fixable today by declaring the provider's `reasoningEfforts` ladder. + +## Carry + +- #5916 (@codingbooo, closes #5880): squash onto `codex/bug-train-7` from `dev` `76b26a0881`, author + `Co-authored-by`. + +## Integration fix (security review blocker) + +The generic OAuth send allowance is `accounts × TRANSIENT_RETRY_MAX_ATTEMPTS` and the hop limit is `accounts - 1`, both +from an uncapped roster snapshot. Add `GENERIC_OAUTH_MAX_ACCOUNTS_PER_REQUEST = 6` in +`src/oauth/generic-account-failover.ts` and apply it in two places: + +1. `src/server/responses/request-transport.ts`: `budgetAccounts = Math.min(genericRosterSize, MAX)`; + `genericFailoverLimit = Math.max(GENERIC_OAUTH_MAX_FAILOVERS_PER_REQUEST, budgetAccounts - 1)`; + `expandInferenceOAuthSendBudget(options.sendBudget, budgetAccounts)`. +2. `src/server/inference/context.ts` `expandInferenceOAuthSendBudget`: clamp `accounts` to the same constant, so any + future caller cannot raise the ceiling either. + +Result: the default ingress OAuth ceiling is 6 × 3 = 18 physical sends per request however many accounts are enrolled, +versus 4 before #5916. Explicitly supplied caller budgets are left unchanged. +The PR's 4- and 5-account cases keep their expectations. + +Regression: in `tests/server/server-google-antigravity-oauth-429-budget.test.ts` (or a sibling if the file-size cap binds), +eight accounts with transient 429 everywhere: exactly 18 sends across exactly the first 6 accounts, then terminal 429. +Must fail without the clamp (24 sends / 8 accounts). + +Docs/structure: update the #5916 text in `structure/transports/responses-failover.md`, `responses-spend.md` and +`docs-site/.../providers.md` to state the ceiling. + +## Check + +tsc, structure:check, privacy:scan, the #5916 test files plus `tests/oauth/generic-oauth-failover.test.ts`, +`tests/oauth/oauth-account-attribution.test.ts`, layout and file-size guards. Independent security re-review of the final +diff. Exact-head hosted CI, then `--admin --match-head-commit`; close #5916 with the batch note. diff --git a/docs-site/src/content/docs/reference/configuration/providers.md b/docs-site/src/content/docs/reference/configuration/providers.md index ae5415a84d8..e45d68b4bbb 100644 --- a/docs-site/src/content/docs/reference/configuration/providers.md +++ b/docs-site/src/content/docs/reference/configuration/providers.md @@ -868,12 +868,13 @@ recorded for it. Before dispatch, generic OAuth snapshots the eligible roster. On a 429 the failed account is cooled using `Retry-After` when present (capped at 15 minutes) or a default backoff, and the request is replayed on the next account selected from the live roster. The stable rotation ceiling is -`max(3, eligibleCount - 1)` per request; live selection still filters cooldowns, and an account +`max(3, min(eligibleCount, 6) - 1)` per request; live selection still filters cooldowns, and an account flagged for reauthentication is never selected. Cooldowns are process-local, so a restart forgets them. When at least two accounts are eligible, the ingress-owned default send allowance covers up to three -sends per eligible account. A single eligible account keeps the existing base allowance of three and +sends per eligible account, counting at most six accounts, so one request +makes no more than 18 sends however many accounts are enrolled. A single eligible account keeps the existing base allowance of three and total allowance of four. Explicit caller ceilings and combo scopes keep their existing limits. Rotation carries the alternate account's **full** credential snapshot, not just its bearer, so a diff --git a/src/oauth/generic-account-failover.ts b/src/oauth/generic-account-failover.ts index 904791f6af2..20301911731 100644 --- a/src/oauth/generic-account-failover.ts +++ b/src/oauth/generic-account-failover.ts @@ -40,6 +40,14 @@ import type { OcxConfig, OcxProviderConfig } from "../types"; /** Cap same-request rotations so a short Retry-After cannot spin. Mirrors the Anthropic bound. */ export const GENERIC_OAUTH_MAX_FAILOVERS_PER_REQUEST = 3; +/** + * Most accounts one request may fund and rotate through, whatever the roster size. Each funded + * account gets the normal transient ladder, so this fixes the default ingress ceiling at + * 6 × TRANSIENT_RETRY_MAX_ATTEMPTS physical sends: enrolling more accounts cannot turn one request + * into a pool-wide 429 storm. + */ +export const GENERIC_OAUTH_MAX_ACCOUNTS_PER_REQUEST = 6; + const DEFAULT_COOLDOWN_MS = 60_000; /** diff --git a/src/server/inference/context.ts b/src/server/inference/context.ts index 08d26ada00b..8883a0edef3 100644 --- a/src/server/inference/context.ts +++ b/src/server/inference/context.ts @@ -1,5 +1,6 @@ import { CODEX_TEXT_GUARDED_BUDGET_POLICY, createRequestExecutionBudget, type RequestExecutionBudget } from "../../lib/request-execution-budget"; import { TRANSIENT_RETRY_MAX_ATTEMPTS, type TransientSendBudget } from "../../lib/upstream-retry"; +import { GENERIC_OAUTH_MAX_ACCOUNTS_PER_REQUEST } from "../../oauth/generic-account-failover"; import type { RequestLogContext } from "../request-log"; import { attachRequestSpendTracker } from "../responses/request-spend"; @@ -25,14 +26,18 @@ export function createInferenceSendBudget( return budget; } -/** Fund each account's normal transient ladder once, before the first physical send. */ +/** + * Fund each account's normal transient ladder once, before the first physical send. The account + * count is clamped here as well as at the caller, so no caller can raise the ceiling past + * GENERIC_OAUTH_MAX_ACCOUNTS_PER_REQUEST × TRANSIENT_RETRY_MAX_ATTEMPTS. + */ export function expandInferenceOAuthSendBudget(budget: TransientSendBudget | undefined, accounts: number): void { if (!budget) return; const policy = ingressPolicies.get(budget); if (!policy || budget.used !== 0) return; ingressPolicies.delete(budget); if (accounts < 2 || !Number.isSafeInteger(accounts)) return; - const sends = accounts * TRANSIENT_RETRY_MAX_ATTEMPTS; + const sends = Math.min(accounts, GENERIC_OAUTH_MAX_ACCOUNTS_PER_REQUEST) * TRANSIENT_RETRY_MAX_ATTEMPTS; policy.baseSendAllowance = Math.max(policy.baseSendAllowance, sends); policy.maxTotalModelSends = Math.max(policy.maxTotalModelSends, sends); } diff --git a/src/server/responses/request-transport.ts b/src/server/responses/request-transport.ts index 2ade1c7ebae..9ee1dab7f67 100644 --- a/src/server/responses/request-transport.ts +++ b/src/server/responses/request-transport.ts @@ -28,6 +28,7 @@ import { UnsupportedOAuthProviderError, } from "../../oauth"; import { + GENERIC_OAUTH_MAX_ACCOUNTS_PER_REQUEST, GENERIC_OAUTH_MAX_FAILOVERS_PER_REQUEST, eligibleFailoverAccounts, forgetGenericFailoverRoster, @@ -727,12 +728,14 @@ export async function prepareResponsesTransport( // Freeze the request ceiling before any 429 writes cooldowns. Selection still reads // live eligibility on every hop; cooled accounts cannot shorten this request's allowance. + // The snapshot is clamped: a larger roster must not raise one request's hops or sends. const genericRosterSize = genericFailoverAccountId ? new Set([ genericFailoverAccountId, ...eligibleFailoverAccounts(route.providerName, Date.now(), classifyModelFamilyForQuota(route.providerName, route.modelId)), ]).size : 0; - const genericFailoverLimit = Math.max(GENERIC_OAUTH_MAX_FAILOVERS_PER_REQUEST, genericRosterSize - 1); - if (genericFailoverAccountId) expandInferenceOAuthSendBudget(options.sendBudget, genericRosterSize); + const fundedAccounts = Math.min(genericRosterSize, GENERIC_OAUTH_MAX_ACCOUNTS_PER_REQUEST); + const genericFailoverLimit = Math.max(GENERIC_OAUTH_MAX_FAILOVERS_PER_REQUEST, fundedAccounts - 1); + if (genericFailoverAccountId) expandInferenceOAuthSendBudget(options.sendBudget, fundedAccounts); return { genericFailoverLimit, diff --git a/structure/transports/responses-failover.md b/structure/transports/responses-failover.md index 9bb15c44f5a..8d4cbb7c3b2 100644 --- a/structure/transports/responses-failover.md +++ b/structure/transports/responses-failover.md @@ -495,7 +495,7 @@ check confirms from the live roster that at least two accounts exist and an alte cooled. That check applies no cooldown and advances no rotation. Generic OAuth snapshots its eligible roster before dispatch. Its request rotation ceiling is -`max(3, eligibleCount - 1)`; the live picker still filters cooldowns, so the snapshot supplies the +`max(3, min(eligibleCount, GENERIC_OAUTH_MAX_ACCOUNTS_PER_REQUEST) - 1)` (the cap is six); the live picker still filters cooldowns, so the snapshot supplies the stable ceiling without making a cooled account eligible. Same-provider auth recovery keeps the last physical target, rather than a diagnostic key, and a real send is charged once even when recovery rebuilds the request. diff --git a/structure/transports/responses-spend.md b/structure/transports/responses-spend.md index b4a67211d19..8f81cfe15a3 100644 --- a/structure/transports/responses-spend.md +++ b/structure/transports/responses-spend.md @@ -12,10 +12,12 @@ budget before it knows whether a rotation is even possible, because the reservat idempotent and a no-op once used. Every ladder therefore owes the budget an answer on every exit. Generic OAuth snapshots the eligible roster at request ingress before dispatch. Its rotation ceiling -is `max(3, eligibleCount - 1)`; live selection still removes accounts in cooldown, so the snapshot +is `max(3, min(eligibleCount, 6) - 1)`; live selection still removes accounts in cooldown, so the snapshot sets the number of possible moves without making a cooled account selectable. Only the ingress-owned default execution budget expands when at least two accounts are eligible: its base and total ceilings -cover up to `TRANSIENT_RETRY_MAX_ATTEMPTS` sends per eligible account (currently three). A single +cover up to `TRANSIENT_RETRY_MAX_ATTEMPTS` sends per eligible account (currently three), for at most +`GENERIC_OAUTH_MAX_ACCOUNTS_PER_REQUEST` (six) accounts, so the default ingress ceiling is 18 sends +whatever the roster size. A single eligible account keeps the existing base ceiling of three and total ceiling of four. Explicit caller ceilings and combo-derived scopes keep their existing limits. diff --git a/tests/server/server-google-antigravity-oauth-429-budget.test.ts b/tests/server/server-google-antigravity-oauth-429-budget.test.ts index 58c98264275..3fbaab1c2e9 100644 --- a/tests/server/server-google-antigravity-oauth-429-budget.test.ts +++ b/tests/server/server-google-antigravity-oauth-429-budget.test.ts @@ -203,6 +203,35 @@ describe("Google Antigravity OAuth 429 retry and multi-account budget (#5880)", } }); + test("a roster larger than the per-request account cap funds only the cap", async () => { + // Eight enrolled accounts must not turn one request into 24 sends: the default ingress + // ceiling is GENERIC_OAUTH_MAX_ACCOUNTS_PER_REQUEST (6) accounts x 3 transient sends. + const accounts = await seedAntigravityAccounts(8); + const cfg = antigravityConfig(); + saveConfig(cfg); + + const observedSends: Array<{ auth: string; project: string }> = []; + installAntigravityFetchMock(({ auth, project }) => { + observedSends.push({ auth, project }); + return new Response(JSON.stringify(transient429ErrorBody()), { + status: 429, + headers: { "content-type": "application/json" }, + }); + }); + + const res = await handleResponses( + createResponsesRequest(), + cfg, + { model: "gemini-3.8-flash", provider: "google-antigravity" }, + ); + + expect(res.status).toBe(429); + expect(observedSends).toHaveLength(18); + const usedAuth = new Set(observedSends.map(send => send.auth)); + expect(usedAuth.size).toBe(6); + for (const account of accounts.slice(6)) expect(usedAuth.has(account.auth)).toBe(false); + }); + test.each([2, 3])("single account succeeds attempt %i on transient 429", async successAttempt => { const accounts = await seedAntigravityAccounts(1); const cfg = antigravityConfig();