From 6bfd76bf6d27d6f3d26b2fc99d41f4a8bd0bb18e Mon Sep 17 00:00:00 2001 From: JUN Date: Fri, 9 Oct 2026 20:41:46 +0900 Subject: [PATCH 1/7] fix(quota): record main-login responses under the observed main credential (#6800) Upstream quota headers were recorded only for pool and main-pool auth contexts, so a request served with the main login as a plain main context (the caller's own ChatGPT bearer, or the stored main substituted for an admission bearer) never refreshed __main__ in codex-quota-cache.json while pool rows refreshed on every response. Materialization now captures a process-local dispatch proof when the bearer actually sent upstream is the main credential the proxy already observed from its own read. HTTP delivery and the WebSocket observer record the response headers under __main__ only if that proof is still live (same identity and credential generation) at write time, and only for the canonical OpenAI forward provider. Pool health, failover and quarantine are unchanged. --- .../docs/reference/cli/providers-accounts.md | 6 + scripts/test-layout/layout.json | 1 + src/codex/auth-context.ts | 21 +- src/codex/main-account-cache.ts | 20 + src/server/responses/core-codex-account.ts | 21 +- src/server/responses/passthrough-delivery.ts | 14 + structure/providers/openai-accounts.md | 1 + structure/providers/openai-tiers.md | 6 + structure/transports/responses.md | 2 +- tests/fixtures/test-layout-expected.json | 1 + .../responses-main-quota-observation.test.ts | 385 ++++++++++++++++++ 11 files changed, 475 insertions(+), 3 deletions(-) create mode 100644 tests/responses/responses-main-quota-observation.test.ts diff --git a/docs-site/src/content/docs/reference/cli/providers-accounts.md b/docs-site/src/content/docs/reference/cli/providers-accounts.md index 88f997b2e87..efd8f00d9a5 100644 --- a/docs-site/src/content/docs/reference/cli/providers-accounts.md +++ b/docs-site/src/content/docs/reference/cli/providers-accounts.md @@ -407,6 +407,12 @@ with the lock status; Direct provider quota omits an unpublished response and it identities and stale 401/403 replies retain the current cached info and cannot clear or set the current account's reauthentication state. +Responses to requests sent with the identified main credential refresh its cached +usage from their quota headers, whether the proxy substituted the stored credential or +the caller sent the same credential itself. A response is applied only if that +credential is still the observed main credential when it arrives; a caller-owned +credential for another account or workspace never updates the main account's usage. + The persisted option is `"codexMainAccountHardLock"` in OpenCodex's `config.json`. An absent key or `true` means on; only an explicit `false` turns it off, and that is what switching the setting off stores. The default changed here: the policy used to be opt-in and the old switch removed the key diff --git a/scripts/test-layout/layout.json b/scripts/test-layout/layout.json index 0b7e20b7a1f..866a1c9f0ea 100644 --- a/scripts/test-layout/layout.json +++ b/scripts/test-layout/layout.json @@ -1313,6 +1313,7 @@ "openai-provider-option.test.ts": "adapters/openai", "responses-forward-client-headers.test.ts": "responses", "openai-responses-passthrough.test.ts": "responses", + "responses-main-quota-observation.test.ts": "responses", "external-task-input-repair.test.ts": "responses", "openai-responses-summary-none.test.ts": "responses", "responses-forward-output-cap.test.ts": "responses", diff --git a/src/codex/auth-context.ts b/src/codex/auth-context.ts index 84929653161..08f29265d1b 100644 --- a/src/codex/auth-context.ts +++ b/src/codex/auth-context.ts @@ -81,11 +81,13 @@ import { import { captureMainAccountIdentityGeneration, captureMainQuotaWriter, + captureMainQuotaDispatch, getObservedMainQuotaIdentityKey, isMainQuotaWriterLive, matchesMainQuotaCredential, observeMainQuotaCredential, type MainQuotaWriter, + type MainQuotaDispatch, } from "./main-account-cache"; import { CODEX_RESERVE_HELPER_UNSUPPORTED_MESSAGE, isCodexReserveHelperUnsupported, isCodexReserveRequestEligible } from "./loopback-target"; import type { DataPlaneAdmission } from "../server/auth-cors"; @@ -275,7 +277,11 @@ export function previewCodexPoolLineage( } export type CodexAuthContext = - | { kind: "main"; accountId: null; reserveAuthorization?: MainReserveAuthorization } + | { + kind: "main"; accountId: null; reserveAuthorization?: MainReserveAuthorization; + /** Captured for the observed main credential actually selected for upstream. */ + mainQuotaDispatch?: MainQuotaDispatch; + } | { kind: "pool"; accountId: string; @@ -1635,6 +1641,14 @@ export class CodexMainSubstitutionUnavailableError extends Error { } } +/** Dispatch identity comes only from an owned observation of the selected credential. */ +function selectedMainQuotaDispatch(selected: Headers): MainQuotaDispatch | undefined { + const bearer = selected.get("authorization")?.replace(/^Bearer\s+/i, "").trim(); + if (!bearer) return undefined; + const accountId = selected.get("chatgpt-account-id") ?? extractAccountId(undefined, bearer); + return captureMainQuotaDispatch(bearer, accountId, captureConfigGeneration()); +} + /** * Build the upstream auth headers for one Codex turn. * @@ -1654,6 +1668,7 @@ export function materializeCodexUpstreamAuth( ctx: CodexAuthContext, options: CodexAuthMaterializationOptions = {}, ): Headers { + if (ctx.kind === "main") ctx.mainQuotaDispatch = undefined; const selected = new Headers(); for (const name of FORWARD_HEADERS) { const value = headers.get(name); @@ -1695,10 +1710,12 @@ export function materializeCodexUpstreamAuth( observeSelectedMainCredential(stored, writer); assertMainAccountPolicy(options.config); assertMaterializedReserve(selected, ctx, options); + ctx.mainQuotaDispatch = selectedMainQuotaDispatch(selected); return selected; } if (callerMatchesObservedMain(selected)) assertMainAccountPolicy(options.config); assertMaterializedReserve(selected, ctx, options); + if (ctx.kind === "main") ctx.mainQuotaDispatch = selectedMainQuotaDispatch(selected); return selected; } @@ -1755,6 +1772,7 @@ export async function materializeCodexUpstreamAuthAsync( if (ctx.kind !== "main" || options.substituteMainCredential !== true) { return materializeCodexUpstreamAuth(headers, ctx, options); } + ctx.mainQuotaDispatch = undefined; const selected = new Headers(); for (const name of FORWARD_HEADERS) { const value = headers.get(name); @@ -1776,6 +1794,7 @@ export async function materializeCodexUpstreamAuthAsync( assertMainAccountPolicy(options.config); // An opt-in enabled during token refresh must not turn a proof-less context into Reserve. assertMaterializedReserve(selected, ctx, options); + ctx.mainQuotaDispatch = selectedMainQuotaDispatch(selected); return selected; } diff --git a/src/codex/main-account-cache.ts b/src/codex/main-account-cache.ts index 81b93dd5289..3a733f71845 100644 --- a/src/codex/main-account-cache.ts +++ b/src/codex/main-account-cache.ts @@ -73,6 +73,26 @@ export function isMainQuotaWriterLive(writer: MainQuotaWriter): boolean { && writer.identityGeneration === mainAccountIdentityGeneration; } +/** Proof that a dispatch used the observed main credential; process-local, never persisted. */ +export type MainQuotaDispatch = Readonly<{ + writer: MainQuotaWriter; + credentialGeneration: number; + configGeneration: number; +}>; + +export function captureMainQuotaDispatch( + accessToken: string, accountId: string | undefined, configGeneration: number, +): MainQuotaDispatch | undefined { + if (!accountId || !matchesMainQuotaCredential(accessToken, accountId)) return undefined; + const writer = captureMainQuotaWriter(accountId); + return writer ? { writer, credentialGeneration: mainQuotaCredentialGeneration, configGeneration } : undefined; +} + +export function isMainQuotaDispatchLive(dispatch: MainQuotaDispatch): boolean { + return isMainQuotaWriterLive(dispatch.writer) + && dispatch.credentialGeneration === mainQuotaCredentialGeneration; +} + export function getObservedMainQuotaIdentityKey(): string | undefined { return observedMainQuotaIdentityKey; } diff --git a/src/server/responses/core-codex-account.ts b/src/server/responses/core-codex-account.ts index 1864a7f6535..da7c3b3f213 100644 --- a/src/server/responses/core-codex-account.ts +++ b/src/server/responses/core-codex-account.ts @@ -6,6 +6,7 @@ import { computeQuotaCooldown, formatCodexProviderForLog, } from "../../codex/routing"; +import { isMainQuotaDispatchLive, type MainQuotaDispatch } from "../../codex/main-account-cache"; import type { CodexWsQuotaObserver } from "./codex-ws-metadata"; import { isCanonicalOpenAiForwardProvider } from "../../providers/openai-tiers"; import { isCodexAccountGenerationLive } from "../../codex/account-store"; @@ -121,8 +122,26 @@ export function usesCodexForwardPoolAuth( } +/** Live proof that this plain-main response used the observed main credential. */ +export function liveMainQuotaDispatch( + authCtx: CodexAuthContext, provider: OcxProviderConfig, +): MainQuotaDispatch | undefined { + if (authCtx.kind !== "main" || !authCtx.mainQuotaDispatch) return undefined; + if (!isCanonicalOpenAiForwardProvider(provider) + || provider.authMode !== "forward" || provider.adapter !== "openai-responses") return undefined; + return isMainQuotaDispatchLive(authCtx.mainQuotaDispatch) ? authCtx.mainQuotaDispatch : undefined; +} + export function codexWsQuotaObserver(authCtx: CodexAuthContext, provider: OcxProviderConfig, modelId?: string): CodexWsQuotaObserver | undefined { - if (!isCanonicalOpenAiForwardProvider(provider) || !usesCodexForwardPoolAuth(authCtx, provider)) return undefined; + if (!isCanonicalOpenAiForwardProvider(provider)) return undefined; + if (!usesCodexForwardPoolAuth(authCtx, provider)) { + const dispatch = liveMainQuotaDispatch(authCtx, provider); + if (!dispatch) return undefined; + return headers => { + if (!isMainQuotaDispatchLive(dispatch)) return; + applyCapturedCodexQuota(MAIN_CODEX_ACCOUNT_ID, headers, dispatch.configGeneration, dispatch.writer, { modelId }); + }; + } const { accountId, writerGeneration } = authCtx; const credentialGeneration = authCtx.kind === "pool" ? authCtx.generation : undefined; const mainWriter = authCtx.kind === "main-pool" ? authCtx.mainQuotaWriter : undefined; diff --git a/src/server/responses/passthrough-delivery.ts b/src/server/responses/passthrough-delivery.ts index c0b86b87595..e50d18d2024 100644 --- a/src/server/responses/passthrough-delivery.ts +++ b/src/server/responses/passthrough-delivery.ts @@ -24,11 +24,14 @@ import { teeWithBoundedInspection } from "../inspection-tee"; import { codexForwardTerminalOutcomeRecorder, usesCodexForwardPoolAuth, + liveMainQuotaDispatch, codexQuotaOutcomeMeta, codexDenialOutcomeMeta, isFixedCodexAccount, shouldDeferCodexResetDerivedCooldown, } from "./core-codex-account"; +import { isMainQuotaDispatchLive } from "../../codex/main-account-cache"; +import { MAIN_CODEX_ACCOUNT_ID } from "../../codex/account-id"; import type { ResponsesTerminalStatus } from "../../bridge"; import { isCodexWsQuotaObservedResponse, isCodexWsUpstreamResponse } from "./ws-upstream"; import { recordSubagentQuotaFailureForThreadSpawn } from "../../codex/subagent-model-fallback"; @@ -535,6 +538,17 @@ export async function deliverPassthroughResponse( ...(admissionState.authCtx.kind === "pool" ? { credentialGeneration: admissionState.authCtx.generation } : {}), }); } + } else { + const mainDispatch = liveMainQuotaDispatch(admissionState.authCtx, route.provider); + if (mainDispatch && !isCodexWsQuotaObservedResponse(upstreamResponse)) { + const { applyAccountQuotaFromUpstreamHeaders } = await import("../../codex/auth-api"); + // Import yields; same-account token replacement leaves the identity writer live. + // Re-check the credential fence with no await before publication. + if (isMainQuotaDispatchLive(mainDispatch)) { + applyAccountQuotaFromUpstreamHeaders(MAIN_CODEX_ACCOUNT_ID, upstreamResponse.headers, + mainDispatch.configGeneration, mainDispatch.writer, { modelId: route.modelId }); + } + } } // Non-2xx passthrough failures must never reach Codex as an empty body — diff --git a/structure/providers/openai-accounts.md b/structure/providers/openai-accounts.md index c408ce0629e..bf22d8bccc6 100644 --- a/structure/providers/openai-accounts.md +++ b/structure/providers/openai-accounts.md @@ -211,6 +211,7 @@ request-scoped: a translated Claude turn resolves through Pool selection like an main keeps its health, quarantine and refresh-and-classify handling. Only a bearer the client itself supplied is caller-owned and exempt from stored state. Both synchronous and asynchronous stored-main substitution in `src/codex/auth-context.ts` remove a caller account header before copying the stored identity; an absent stored account ID leaves no account header. Caller-owned native Direct authentication retains its existing passthrough behavior. +Plain-main HTTP and WebSocket Responses refresh `__main__` only on the canonical OpenAI forward provider when the sent bearer and effective workspace match the main credential already observed under native ownership, the same equality rule used by the hard lock. `src/codex/auth-context.ts` captures a process-local dispatch proof after materialization; `src/server/responses/passthrough-delivery.ts` and `src/server/responses/core-codex-account.ts` recheck its credential/identity generations before publishing. Same-account token rotation and A→B→A changes reject old responses; WebSocket observers retain the original dispatch proof across frames and their projected HTTP response does not publish twice. Caller-owned requests acquire no physical-main read or Pool health state. `src/providers/openai-sidecar.ts` releases quota-probe ownership on every materialization or usability failure before transferring a resolved context to its caller. Audio reports one terminal upstream outcome after validating the response body; redirects remain diff --git a/structure/providers/openai-tiers.md b/structure/providers/openai-tiers.md index 889bfd3c656..8308f52a147 100644 --- a/structure/providers/openai-tiers.md +++ b/structure/providers/openai-tiers.md @@ -429,6 +429,12 @@ invalidates old evidence. Request-owned bearers are matched only against a crede workspace already observed under native ownership; an unrelated or unmatched keyring credential is not attributed to stored main and introduces no physical-main read. Credential equality tags remain process-local and never enter disk, logs, or management DTOs. +Plain-main HTTP and WebSocket Responses on the canonical OpenAI forward provider refresh cached +main usage under that same credential/workspace match, including stored-main substitution and +an identical caller-owned credential. Materialization captures a process-local dispatch proof; +publication rechecks identity and credential generations, including after an awaited HTTP import +and for every WebSocket frame. A replaced credential, unmatched workspace, or custom destination +cannot publish main usage. Pool health/failover handling stays scoped to Pool contexts. `src/codex/auth-api/main-account-probe.ts` re-reads the bounded stored main credential and rechecks its writer, bearer and generation after body/retry awaits, before publishing main usage, credits, plan, reauth or Reserve state, including terminal 401/403 mutations. An unreadable file diff --git a/structure/transports/responses.md b/structure/transports/responses.md index 3c6b5c7025b..263d2599fff 100644 --- a/structure/transports/responses.md +++ b/structure/transports/responses.md @@ -23,7 +23,7 @@ Canonical forward auth retains its separate fixed credential/metadata allowlist; Retired Codex Spark has no model-specific tool or Responses Lite override; general Lite handling and namespace scrubbing remain shared compatibility behavior. Codex quota/reset evidence follows the -[shared/Reserve policy](../providers/openai-tiers.md#public-provider-contract), including suppression of retired model-derived evidence before shared recovery. +[shared/Reserve policy](../providers/openai-tiers.md#public-provider-contract), including suppression of retired model-derived evidence before shared recovery. Plain-main HTTP/WS quota headers refresh `__main__` only on the canonical OpenAI forward provider when the materialized bearer and workspace match the owned main observation, using the hard-lock credential-match rule. `src/server/responses/passthrough-delivery.ts` rechecks the captured credential generation after its awaited import; `src/server/responses/core-codex-account.ts` retains one dispatch proof per WS observer and rechecks it for each frame. Rotated credentials and unmatched workspaces publish nothing; WS-observed responses skip duplicate HTTP writes, and Pool health/failover gates remain unchanged. ### Credential-bearing HTTP redirects diff --git a/tests/fixtures/test-layout-expected.json b/tests/fixtures/test-layout-expected.json index 1560995bd4d..6baefb6075c 100644 --- a/tests/fixtures/test-layout-expected.json +++ b/tests/fixtures/test-layout-expected.json @@ -808,6 +808,7 @@ "openai-provider-option-startup.test.ts": "adapters/openai", "openai-provider-option-tooling.test.ts": "adapters/openai", "openai-provider-option.test.ts": "adapters/openai", "responses-forward-client-headers.test.ts": "responses", "openai-responses-passthrough.test.ts": "responses", + "responses-main-quota-observation.test.ts": "responses", "external-task-input-repair.test.ts": "responses", "openai-responses-summary-none.test.ts": "responses", "responses-forward-output-cap.test.ts": "responses", "opencode-cli.test.ts": "providers", "opencode-free-provider.test.ts": "providers", diff --git a/tests/responses/responses-main-quota-observation.test.ts b/tests/responses/responses-main-quota-observation.test.ts new file mode 100644 index 00000000000..91941365c26 --- /dev/null +++ b/tests/responses/responses-main-quota-observation.test.ts @@ -0,0 +1,385 @@ +import { afterEach, beforeEach, describe, expect, spyOn, test } from "bun:test"; +import { mkdirSync, mkdtempSync, readFileSync, writeFileSync } from "node:fs"; +import { join } from "node:path"; +import { MAIN_CODEX_ACCOUNT_ID as MAIN } from "../../src/codex/account-id"; +import * as mainCache from "../../src/codex/main-account-cache"; +import { + materializeCodexUpstreamAuth, materializeCodexUpstreamAuthAsync, resolveCodexAuthContext, + type CodexAuthContext, +} from "../../src/codex/auth-context"; +import { resetMainCodexAccountIdentityTrackingForTests } from "../../src/codex/account-lifecycle"; +import { clearAccountNeedsReauth } from "../../src/codex/account-runtime-state"; +import { saveCodexAccountCredential } from "../../src/codex/account-store"; +import { clearAccountQuota, getAccountQuota, getMainPolicyQuota, setAccountQuotaFromParsed } from "../../src/codex/quota"; +import { clearCodexUpstreamHealth, clearThreadAccountMap } from "../../src/codex/routing"; +import { clearPoolRotationState } from "../../src/codex/pool-rotation"; +import { createTranslatorBudget } from "../../src/lib/translator-budget"; +import { captureConfigGeneration } from "../../src/lib/state-store-sweeper"; +import { parseRequest } from "../../src/responses/parser"; +import { codexAccountSelectionForTurn, tryAdmitTurn } from "../../src/server/lifecycle"; +import { deliverPassthroughResponse } from "../../src/server/responses/passthrough-delivery"; +import { codexWsQuotaObserver, retryCodexPoolOnAlternateAccount } from "../../src/server/responses/core-codex-account"; +import { CodexWsMetadata } from "../../src/server/responses/codex-ws-metadata"; +import { markCodexWsResponse, isCodexWsQuotaObservedResponse } from "../../src/server/responses/codex-ws-wire"; +import type { OcxConfig, OcxProviderConfig } from "../../src/types"; +import { fakeChatGptJwt } from "../helpers/agent-task-recovery"; +import { acquireOwnedSpendHome } from "../helpers/owned-spend-home"; +import { removeTreeWithRetry } from "../helpers/remove-tree"; +import { repoPath } from "../helpers/repo-root"; + +const ACCOUNT = "fixture-observed-main"; +const POOL = "fixture-observed-pool"; +const provider: OcxProviderConfig = { + adapter: "openai-responses", authMode: "forward", baseUrl: "https://chatgpt.com/backend-api/codex", +}; +const originalFetch = globalThis.fetch; +let root: string; +let oldOcxHome: string | undefined; +let oldCodexHome: string | undefined; +let bearer: string; +let config: OcxConfig; +let releaseSpend: (() => void) | undefined; +let pendingPersist: { run: () => void; timer: ReturnType } | undefined; +let clock: ReturnType; + +// Exercise the real serializer without a wall-clock race, as main-quota-provenance does. +function installPersistenceClock() { + const nativeTimeout = globalThis.setTimeout; + return spyOn(globalThis, "setTimeout").mockImplementation((( + callback: (...args: unknown[]) => void, delay?: number, ...args: unknown[] + ) => { + if (delay !== 250) return nativeTimeout(callback, delay, ...args); + const timer = nativeTimeout(() => {}, 60_000); + pendingPersist = { run: () => callback(...args), timer }; + return timer; + }) as typeof setTimeout); +} + +function observe(token = bearer, account = ACCOUNT): void { + mainCache.observeMainQuotaIdentity(account); + expect(mainCache.observeMainQuotaCredential(token, account)).toBeDefined(); +} +function caller(token = bearer, account: string | undefined = ACCOUNT): Headers { + const headers = new Headers({ authorization: `Bearer ${token}` }); + if (account !== undefined) headers.set("chatgpt-account-id", account); + return headers; +} +function materialized(headers = caller()): Extract { + const ctx: Extract = { kind: "main", accountId: null }; + materializeCodexUpstreamAuth(headers, ctx, { config, modelId: "gpt-5.5" }); + return ctx; +} +function quotaHeaders(percent = "23"): Headers { + return new Headers({ "x-codex-primary-used-percent": percent, "x-codex-primary-window-minutes": "10080" }); +} +function quotaResponse(headers = quotaHeaders()): Response { + // A real upstream redirect takes the early relay return after quota publication. This keeps + // this delivery fixture independent of unrelated success-body repair and continuation state. + headers = new Headers(headers); + headers.set("location", "https://chatgpt.com/backend-api/codex/responses"); + return new Response(null, { status: 307, headers }); +} +async function deliver(ctx: CodexAuthContext, response = quotaResponse(), selectedProvider = provider): Promise { + type Args = Parameters; + const budget = createTranslatorBudget(); + try { + const result = await deliverPassthroughResponse( + { config, logCtx: { model: "", provider: "" }, options: {}, req: new Request("http://localhost/v1/responses") }, + { authCtx: ctx } as Args[1], + { parsed: parseRequest({ model: "gpt-5.5", input: "hi", stream: false }), + route: { providerName: "openai", modelId: "gpt-5.5", provider: selectedProvider }, + clientRequestedStream: false, translatorBudget: budget, inboundWire: "responses" }, + { requestBindings: undefined } as Args[3], {}, + { plaintextV2AgentMessageToolNames: new Set(), routedMuseToolNameAliases: new Map(), + routedNamespaceToolAliases: new Map(), plaintextV2AgentMessageAliasedToolNames: new Set(), + commitReasoningReplayServingRoute: () => {}, recordTerminalOutcomes: () => {}, + responseCompletionCancelled: () => false } as Args[5], + { upstreamResponse: response, upstream: new AbortController(), connectMs: 1000 } as Args[6], + ); + expect(result.status).toBe(307); + await result.body?.cancel(); + } finally { budget.dispose(); } +} +function frame(metadata: CodexWsMetadata, percent: number): void { + const event = { type: "codex.rate_limits", rate_limits: { primary: { used_percent: percent, window_minutes: 10080 } } }; + expect(metadata.consume(event, JSON.stringify(event).length)).not.toBeNull(); +} + +beforeEach(() => { + mkdirSync(repoPath(".tmp"), { recursive: true }); + root = mkdtempSync(repoPath(".tmp/main-quota-observation-")); + oldOcxHome = process.env.OPENCODEX_HOME; + oldCodexHome = process.env.CODEX_HOME; + process.env.OPENCODEX_HOME = root; + process.env.CODEX_HOME = root; + releaseSpend = acquireOwnedSpendHome(); + clearAccountQuota(); + clearCodexUpstreamHealth(); + clearThreadAccountMap(); + clearPoolRotationState(); + clearAccountNeedsReauth(MAIN); + clearAccountNeedsReauth(POOL); + resetMainCodexAccountIdentityTrackingForTests(); + mainCache.clearMainAccountInfoCache(); + mainCache.observeMainQuotaIdentity("fixture-unobserved"); + bearer = fakeChatGptJwt(ACCOUNT); + config = { providers: { openai: provider }, codexAccounts: [], codexMainAccountHardLock: false } as OcxConfig; + pendingPersist = undefined; + clock = installPersistenceClock(); + globalThis.fetch = (async () => { throw new Error("unexpected network call"); }) as typeof fetch; +}); +afterEach(() => { + releaseSpend?.(); releaseSpend = undefined; + globalThis.fetch = originalFetch; + clearAccountQuota(); + if (pendingPersist) clearTimeout(pendingPersist.timer); + pendingPersist = undefined; + clock.mockRestore(); + clearCodexUpstreamHealth(); clearThreadAccountMap(); clearPoolRotationState(); + clearAccountNeedsReauth(MAIN); clearAccountNeedsReauth(POOL); + resetMainCodexAccountIdentityTrackingForTests(); mainCache.clearMainAccountInfoCache(); + if (oldOcxHome === undefined) delete process.env.OPENCODEX_HOME; + else process.env.OPENCODEX_HOME = oldOcxHome; + if (oldCodexHome === undefined) delete process.env.CODEX_HOME; + else process.env.CODEX_HOME = oldCodexHome; + removeTreeWithRetry(root); +}); + +describe("credential-bound plain-main Responses quota", () => { + test("1: observed caller bearer updates main through HTTP delivery", async () => { + observe(); + const ctx = materialized(); + expect(ctx.mainQuotaDispatch).toBeDefined(); + await deliver(ctx); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(23); + expect(getMainPolicyQuota()?.weeklyPercent).toBe(23); + // Materialization resets an old proof when this context is reused for another bearer. + materializeCodexUpstreamAuth(caller("fixture-unmatched"), ctx, { config }); + expect(ctx.mainQuotaDispatch).toBeUndefined(); + await deliver(ctx, quotaResponse(quotaHeaders("31"))); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(23); + }); + + for (const mode of ["sync", "async"] as const) { + test(`2: ${mode} stored-main substitution captures the sent credential`, async () => { + observe(); + writeFileSync(join(root, "auth.json"), JSON.stringify({ tokens: { access_token: bearer, account_id: ACCOUNT } })); + const ctx: Extract = { kind: "main", accountId: null }; + const options = { config, substituteMainCredential: true }; + const selected = mode === "sync" + ? materializeCodexUpstreamAuth(caller("fixture-admission", "fixture-wrong-workspace"), ctx, options) + : await materializeCodexUpstreamAuthAsync(caller("fixture-admission", "fixture-wrong-workspace"), ctx, options); + // Compare without putting credential material in a failed assertion's output. + expect(selected.get("authorization") === `Bearer ${bearer}`).toBe(true); + expect(selected.get("chatgpt-account-id") === ACCOUNT).toBe(true); + expect(ctx.mainQuotaDispatch).toBeDefined(); + await deliver(ctx); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(23); + }); + } + + test("3: different account, different token and unobserved caller cannot publish", async () => { + const unobserved = materialized(); + expect(unobserved.mainQuotaDispatch).toBeUndefined(); + await deliver(unobserved); + observe(); + for (const headers of [caller(fakeChatGptJwt("fixture-other"), "fixture-other"), caller("fixture-other-token")]) { + const ctx = materialized(headers); + expect(ctx.mainQuotaDispatch).toBeUndefined(); + await deliver(ctx); + } + expect(getAccountQuota(MAIN)).toBeNull(); + expect(getMainPolicyQuota()).toBeNull(); + }); + + test("4: same-account rotation and A to B to A reject old HTTP dispatches", async () => { + observe(); + const rotation = materialized(); + observe("fixture-rotated"); + expect(mainCache.isMainQuotaWriterLive(rotation.mainQuotaDispatch!.writer)).toBe(true); + expect(mainCache.isMainQuotaDispatchLive(rotation.mainQuotaDispatch!)).toBe(false); + await deliver(rotation); + observe(); + const aba = materialized(); + observe("fixture-b", "fixture-account-b"); + observe(); + await deliver(aba); + expect(getAccountQuota(MAIN)).toBeNull(); + expect(getMainPolicyQuota()).toBeNull(); + }); + + test("5: non-canonical, key-auth and other-adapter providers cannot publish", async () => { + observe(); + const ctx = materialized(); + for (const other of [{ ...provider, baseUrl: "https://fixture.test/v1" }, + { ...provider, authMode: "key" as const }, { ...provider, adapter: "openai-chat" }]) { + await deliver(ctx, quotaResponse(), other); + expect(codexWsQuotaObserver(ctx, other)).toBeUndefined(); + } + expect(getAccountQuota(MAIN)).toBeNull(); + }); + + test("6: WS metadata publishes once and its projected HTTP response cannot duplicate", async () => { + observe(); + const ctx = materialized(); + const observer = codexWsQuotaObserver(ctx, provider, "gpt-5.5"); + expect(observer).toBeDefined(); + const metadata = new CodexWsMetadata(observer); + frame(metadata, 19); + const quota = getAccountQuota(MAIN); + expect(quota?.weeklyPercent).toBe(19); + // A distinct projected value detects a second write even within the same millisecond. + const projected = quotaResponse(quotaHeaders("47")); + markCodexWsResponse(projected, true); + expect(isCodexWsQuotaObservedResponse(projected)).toBe(true); + await deliver(ctx, projected); + expect(getAccountQuota(MAIN)).toEqual(quota); + expect(getMainPolicyQuota()?.weeklyPercent).toBe(19); + metadata.finish(); + }); + + test("7: stored pool HTTP and WS still update only their pool row", async () => { + const generation = saveCodexAccountCredential(POOL, { accessToken: "fixture-pool-token", + refreshToken: "fixture-pool-refresh", chatgptAccountId: "fixture-pool-workspace", expiresAt: Date.now() + 3600_000 }); + const ctx: Extract = { kind: "pool", accountId: POOL, + generation, writerGeneration: captureConfigGeneration(), accessToken: "fixture-pool-token", chatgptAccountId: "fixture-pool-workspace" }; + await deliver(ctx); + expect(getAccountQuota(POOL)?.weeklyPercent).toBe(23); + const observer = codexWsQuotaObserver(ctx, provider, "gpt-5.5"); + expect(observer).toBeDefined(); + observer!(quotaHeaders("29")); + expect(getAccountQuota(POOL)?.weeklyPercent).toBe(29); + expect(getAccountQuota(MAIN)).toBeNull(); + }); + + test("8: credential replacement during the HTTP import yield rejects publication", async () => { + observe(); + const ctx = materialized(); + const live = mainCache.isMainQuotaDispatchLive; + let checks = 0; + const fence = spyOn(mainCache, "isMainQuotaDispatchLive").mockImplementation(dispatch => { + const current = live(dispatch); + if (++checks === 1) queueMicrotask(() => observe("fixture-during-import")); + return current; + }); + try { + await deliver(ctx); + expect(checks).toBe(2); + expect(mainCache.isMainQuotaWriterLive(ctx.mainQuotaDispatch!.writer)).toBe(true); + expect(getAccountQuota(MAIN)).toBeNull(); + expect(getMainPolicyQuota()).toBeNull(); + } finally { fence.mockRestore(); } + }); + + test("9: WS closure rejects a second frame after rotation despite mutable context recapture", () => { + observe(); + const ctx = materialized(); + const metadata = new CodexWsMetadata(codexWsQuotaObserver(ctx, provider, "gpt-5.5")); + frame(metadata, 17); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(17); + observe("fixture-new-token"); + materializeCodexUpstreamAuth(caller("fixture-new-token"), ctx, { config }); + expect(mainCache.isMainQuotaDispatchLive(ctx.mainQuotaDispatch!)).toBe(true); + frame(metadata, 39); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(17); + expect(getMainPolicyQuota()?.weeklyPercent).toBe(17); + metadata.finish(); + }); + + test("10: the same bearer for a different workspace cannot publish", async () => { + observe(); + const ctx = materialized(caller(bearer, "fixture-different-workspace")); + expect(ctx.mainQuotaDispatch).toBeUndefined(); + await deliver(ctx); + expect(getAccountQuota(MAIN)).toBeNull(); + expect(getMainPolicyQuota()).toBeNull(); + }); + + test("11: absent and nonnumeric headers do nothing; invalid ranges clamp display but retain policy", async () => { + observe(); + const ctx = materialized(); + for (const headers of [new Headers(), quotaHeaders("not-a-number")]) await deliver(ctx, quotaResponse(headers)); + expect(getAccountQuota(MAIN)).toBeNull(); + expect(getMainPolicyQuota()).toBeNull(); + setAccountQuotaFromParsed(MAIN, { weeklyPercent: 44, shortPercent: 12, shortWindowSeconds: 18_000 }, + undefined, ctx.mainQuotaDispatch!.writer); + const before = getMainPolicyQuota(); + await deliver(ctx, quotaResponse(quotaHeaders("120"))); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(100); + expect(getMainPolicyQuota()).toEqual(before); + await deliver(ctx, quotaResponse(new Headers({ "x-codex-primary-used-percent": "-5", + "x-codex-primary-window-minutes": "300", "x-codex-secondary-used-percent": "25" }))); + expect(getAccountQuota(MAIN)).toMatchObject({ shortPercent: 0, weeklyPercent: 25 }); + // The existing main-pool consumer rejects the entire policy projection of a mixed set. + expect(getMainPolicyQuota()).toEqual(before); + }); + + test("12: real persistence saves fresh main usage without credential or dispatch proof", async () => { + observe(); + const before = Date.now(); + const ctx = materialized(); + await deliver(ctx); + expect(pendingPersist).toBeDefined(); + const pending = pendingPersist!; + pendingPersist = undefined; + clearTimeout(pending.timer); + pending.run(); + const body = readFileSync(join(root, "codex-quota-cache.json"), "utf8"); + const persisted = JSON.parse(body); + expect(persisted.quotas[MAIN].weeklyPercent).toBe(23); + expect(persisted.quotas[MAIN].updatedAt).toBeGreaterThanOrEqual(before); + expect(persisted.quotas[MAIN].updatedAt).toBeLessThanOrEqual(Date.now()); + for (const forbidden of [bearer, ACCOUNT, "bearerHmac", "mainQuotaDispatch", "credentialGeneration", "configGeneration", "identityGeneration"]) + expect(body.includes(forbidden)).toBe(false); + expect(Object.keys(persisted.mainPolicyQuota).sort()).toEqual(["identityKey", "quota"]); + }); + + test("13: actual pool to caller-main retry publishes only the final dispatch proof", async () => { + observe(); + saveCodexAccountCredential(POOL, { accessToken: "fixture-retry-pool", + refreshToken: "fixture-retry-refresh", chatgptAccountId: "fixture-pool-workspace", expiresAt: Date.now() + 3600_000 }); + config = { ...config, activeCodexAccountId: POOL, autoSwitchThreshold: 0, + providers: { openai: { ...provider, codexAccountMode: "pool" } }, codexAccounts: [{ id: POOL, label: "fixture pool" }] }; + const turn = tryAdmitTurn(); + expect(turn).not.toBeNull(); + const budget = createTranslatorBudget(); + let sends = 0; + globalThis.fetch = (async (_input, init) => { + sends++; + const headers = new Headers(init?.headers); + expect(headers.get("authorization") === `Bearer ${bearer}`).toBe(true); + expect(headers.get("chatgpt-account-id") === ACCOUNT).toBe(true); + return quotaResponse(quotaHeaders("32")); + }) as typeof fetch; + try { + const firstAuthCtx = await resolveCodexAuthContext(caller(), config, "pool", { + modelId: "gpt-5.5", requestScopedMainCredential: true, + beginCodexAccountSelection: codexAccountSelectionForTurn(turn!), + }); + expect(firstAuthCtx.kind).toBe("pool"); + if (firstAuthCtx.kind !== "pool") throw new Error("fixture did not select pool"); + const result = await retryCodexPoolOnAlternateAccount({ callerAuthHeaders: caller(), config, firstAuthCtx, + firstResponse: new Response(null, { status: 429, headers: { "retry-after": "60", ...Object.fromEntries(quotaHeaders("91")) } }), + outcomeStatus: 429, route: { providerName: "openai", modelId: "gpt-5.5", provider: config.providers.openai! }, + parsed: parseRequest({ model: "gpt-5.5", input: "hi", stream: false }), logCtx: { model: "", provider: "" }, + options: { translatorBudget: budget, turnAdmissionLease: turn! }, upstream: new AbortController(), + connectMs: 1000, stream: false, httpOnly: true }); + expect(result.kind).toBe("retried"); + if (result.kind !== "retried") throw new Error("fixture did not retry"); + expect(sends).toBe(1); + expect(result.authCtx.kind).toBe("main"); + if (result.authCtx.kind !== "main") throw new Error("fixture did not choose caller main"); + expect(result.authCtx.mainQuotaDispatch).toBeDefined(); + expect(getAccountQuota(POOL)?.weeklyPercent).toBe(91); + expect(getAccountQuota(MAIN)).toBeNull(); + await deliver(result.authCtx, result.upstreamResponse); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(32); + // An older completed dispatch cannot overwrite a newer credential's observation. + observe("fixture-final-replacement"); + const final = materialized(caller("fixture-final-replacement")); + await deliver(final, quotaResponse(quotaHeaders("41"))); + await deliver(result.authCtx, quotaResponse(quotaHeaders("79"))); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(41); + } finally { turn?.release(); budget.dispose(); } + }); +}); From dbc171d78db83fc033e9ce82a99bdbd3a129b9c3 Mon Sep 17 00:00:00 2001 From: JUN Date: Fri, 9 Oct 2026 21:01:59 +0900 Subject: [PATCH 2/7] fix(quota): fence main quota dispatches on credential commits and WS refusals Review follow-up for #6800: - the dispatch proof also captures the process-wide credential mutation epoch, so a native refresh or same-account reauth commit between dispatch and response drops the main quota update; - a WebSocket exchange publishes plain-main quota only through its observer; the HTTP path skips WS upstream responses and precommit refusal projections, whose headers replay the prelude snapshot; - the async materializer clears a reused context's proof before Reserve delegation. --- src/codex/auth-context.ts | 2 +- src/codex/main-account-cache.ts | 10 +- src/server/responses/codex-ws-exchange.ts | 6 +- src/server/responses/codex-ws-wire.ts | 10 ++ src/server/responses/passthrough-delivery.ts | 6 +- src/server/responses/ws-upstream.ts | 2 +- structure/providers/openai-accounts.md | 2 +- structure/providers/openai-tiers.md | 3 + structure/transports/responses.md | 2 +- .../responses-main-quota-observation.test.ts | 152 +++++++++++++++++- 10 files changed, 181 insertions(+), 14 deletions(-) diff --git a/src/codex/auth-context.ts b/src/codex/auth-context.ts index 08f29265d1b..3903e295fde 100644 --- a/src/codex/auth-context.ts +++ b/src/codex/auth-context.ts @@ -1766,13 +1766,13 @@ export async function materializeCodexUpstreamAuthAsync( ctx: CodexAuthContext, options: CodexAuthMaterializationOptions = {}, ): Promise { + if (ctx.kind === "main") ctx.mainQuotaDispatch = undefined; if (requiresReserveAuthorization(options.config, options.modelId, options.admission)) { return materializeReserveUpstreamAuth(headers, ctx, options); } if (ctx.kind !== "main" || options.substituteMainCredential !== true) { return materializeCodexUpstreamAuth(headers, ctx, options); } - ctx.mainQuotaDispatch = undefined; const selected = new Headers(); for (const name of FORWARD_HEADERS) { const value = headers.get(name); diff --git a/src/codex/main-account-cache.ts b/src/codex/main-account-cache.ts index 3a733f71845..e6551df38ff 100644 --- a/src/codex/main-account-cache.ts +++ b/src/codex/main-account-cache.ts @@ -1,4 +1,5 @@ import { createHash, createHmac, randomBytes, timingSafeEqual } from "node:crypto"; +import { codexCredentialMutationEpoch } from "./credential-mutation-epoch"; import type { StoredAccountQuota } from "./quota-types"; import { truncateRetainedUtf8 } from "../lib/admission"; @@ -77,6 +78,7 @@ export function isMainQuotaWriterLive(writer: MainQuotaWriter): boolean { export type MainQuotaDispatch = Readonly<{ writer: MainQuotaWriter; credentialGeneration: number; + credentialMutationEpoch: number; configGeneration: number; }>; @@ -85,12 +87,16 @@ export function captureMainQuotaDispatch( ): MainQuotaDispatch | undefined { if (!accountId || !matchesMainQuotaCredential(accessToken, accountId)) return undefined; const writer = captureMainQuotaWriter(accountId); - return writer ? { writer, credentialGeneration: mainQuotaCredentialGeneration, configGeneration } : undefined; + return writer ? { writer, credentialGeneration: mainQuotaCredentialGeneration, + credentialMutationEpoch: codexCredentialMutationEpoch(), configGeneration } : undefined; } export function isMainQuotaDispatchLive(dispatch: MainQuotaDispatch): boolean { + // Other OpenCodex-owned credential publications also advance this epoch; + // dropping a main quota update after any such publication is the intended safe direction. return isMainQuotaWriterLive(dispatch.writer) - && dispatch.credentialGeneration === mainQuotaCredentialGeneration; + && dispatch.credentialGeneration === mainQuotaCredentialGeneration + && dispatch.credentialMutationEpoch === codexCredentialMutationEpoch(); } export function getObservedMainQuotaIdentityKey(): string | undefined { diff --git a/src/server/responses/codex-ws-exchange.ts b/src/server/responses/codex-ws-exchange.ts index ddb587a1a4d..0944a074c15 100644 --- a/src/server/responses/codex-ws-exchange.ts +++ b/src/server/responses/codex-ws-exchange.ts @@ -9,7 +9,7 @@ import { CODEX_RESPONSES_HTTP_URL, type PreparedCodexWsRequest } from "./codex-w import { CodexWsCorrelation } from "./codex-ws-correlation"; import type { CodexWsSession } from "./codex-ws-session"; import { UPGRADE_DEADLINE_MS, CODEX_WS_LIVENESS_PING_INTERVAL_MS, CODEX_WS_RESPONSE_PRELUDE_TIMEOUT_MS, MAX_CODEX_WS_FRAME_BYTES, - MAX_CODEX_WS_QUEUE_BYTES, markCodexWsResponse, normalizeResponsesWsRelayEvent, closedBeforeTerminalMessage, + MAX_CODEX_WS_QUEUE_BYTES, markCodexWsResponse, markCodexWsRejectionResponse, normalizeResponsesWsRelayEvent, closedBeforeTerminalMessage, codexWsCreateFrameExceedsLimit, codexWsFailureDetail, codexWsPreResponseFailure, markCodexWsStage, codexWsOcxVersion, markCodexWsSocketDeath, type CodexWsFailureStage, type CodexWsStageRecord } from "./codex-ws-wire"; @@ -45,7 +45,8 @@ function rejectionHeaders(source: Record, prelude: Headers): He } } // Reuse the metadata owner's count/value/family budgets and window freshness - // rules, without publishing quota twice. The unmarked HTTP response owns it. + // rules, without publishing quota twice. Pool bookkeeping consumes the HTTP + // projection; plain-main publication belongs only to the WS observer. const projected = new CodexWsMetadata(); try { for (const values of [Object.fromEntries(prelude), source]) { @@ -513,6 +514,7 @@ export function codexWsExchange(options: ExchangeOptions): Promise { cleanup(); try { controller.close(); } catch { /* unused stream already closed */ } session.dispose(); + markCodexWsRejectionResponse(rejection); resolve(rejection); return; } diff --git a/src/server/responses/codex-ws-wire.ts b/src/server/responses/codex-ws-wire.ts index af8539e93d3..dccfbbeff7b 100644 --- a/src/server/responses/codex-ws-wire.ts +++ b/src/server/responses/codex-ws-wire.ts @@ -49,6 +49,7 @@ const WS_CLOSE_MESSAGE_TOO_BIG = 1009; const codexWsUpstreamResponses = new WeakSet(); const quotaObservedResponses = new WeakSet(); +const codexWsRejectionResponses = new WeakSet(); /** Quota arrived directly at its captured account; do not replay old HTTP prelude headers. */ export function isCodexWsQuotaObservedResponse(response: Response): boolean { @@ -61,6 +62,15 @@ export function isCodexWsUpstreamResponse(response: Response): boolean { } +/** Precommit refusal projection; its headers retain the exchange's prelude snapshot. */ +export function markCodexWsRejectionResponse(response: Response): void { + codexWsRejectionResponses.add(response); +} + +export function isCodexWsRejectionResponse(response: Response): boolean { + return codexWsRejectionResponses.has(response); +} + export function markCodexWsResponse(response: Response, observed: boolean): void { codexWsUpstreamResponses.add(response); if (observed) quotaObservedResponses.add(response); diff --git a/src/server/responses/passthrough-delivery.ts b/src/server/responses/passthrough-delivery.ts index e50d18d2024..1ecdfe44733 100644 --- a/src/server/responses/passthrough-delivery.ts +++ b/src/server/responses/passthrough-delivery.ts @@ -33,7 +33,7 @@ import { import { isMainQuotaDispatchLive } from "../../codex/main-account-cache"; import { MAIN_CODEX_ACCOUNT_ID } from "../../codex/account-id"; import type { ResponsesTerminalStatus } from "../../bridge"; -import { isCodexWsQuotaObservedResponse, isCodexWsUpstreamResponse } from "./ws-upstream"; +import { isCodexWsQuotaObservedResponse, isCodexWsUpstreamResponse, isCodexWsRejectionResponse } from "./ws-upstream"; import { recordSubagentQuotaFailureForThreadSpawn } from "../../codex/subagent-model-fallback"; import { recordCodexUpstreamOutcome } from "../../codex/routing"; import { codexProbeLeaseId, codexProbeQuotaScope, codexTransientProbeGrant, releaseCodexAuthContextProbeLease } from "../../codex/auth-context"; @@ -540,7 +540,9 @@ export async function deliverPassthroughResponse( } } else { const mainDispatch = liveMainQuotaDispatch(admissionState.authCtx, route.provider); - if (mainDispatch && !isCodexWsQuotaObservedResponse(upstreamResponse)) { + // The WS observer is the only plain-main publisher for a WebSocket exchange; + // a refusal projection carries the prelude snapshot, not fresh evidence. + if (mainDispatch && !(isCodexWsUpstreamResponse(upstreamResponse) || isCodexWsRejectionResponse(upstreamResponse))) { const { applyAccountQuotaFromUpstreamHeaders } = await import("../../codex/auth-api"); // Import yields; same-account token replacement leaves the identity writer live. // Re-check the credential fence with no await before publication. diff --git a/src/server/responses/ws-upstream.ts b/src/server/responses/ws-upstream.ts index b51f66abc1a..d2569e536af 100644 --- a/src/server/responses/ws-upstream.ts +++ b/src/server/responses/ws-upstream.ts @@ -26,7 +26,7 @@ import { codexWsCreateFrameExceedsLimit } from "./codex-ws-wire"; import { isLoopbackUrl, rewriteWebSocketDial } from "../../plugins/upstream-hooks"; export { CODEX_WS_LIVENESS_PING_INTERVAL_MS, CODEX_WS_RESPONSE_PRELUDE_TIMEOUT_MS, MAX_CODEX_WS_FRAME_BYTES, MAX_CODEX_WS_QUEUE_BYTES, MAX_CODEX_WS_CREATE_FRAME_BYTES, CODEX_WS_CREATE_FRAME_LIMIT_BYTES, codexWsCreateFrameExceedsLimit, - isCodexWsQuotaObservedResponse, isCodexWsUpstreamResponse } from "./codex-ws-wire"; + isCodexWsQuotaObservedResponse, isCodexWsUpstreamResponse, isCodexWsRejectionResponse } from "./codex-ws-wire"; export const MIN_BOUNDED_CODEX_WS_BUN_VERSION = "1.4.0"; /** diff --git a/structure/providers/openai-accounts.md b/structure/providers/openai-accounts.md index bf22d8bccc6..64a68b63fce 100644 --- a/structure/providers/openai-accounts.md +++ b/structure/providers/openai-accounts.md @@ -211,7 +211,7 @@ request-scoped: a translated Claude turn resolves through Pool selection like an main keeps its health, quarantine and refresh-and-classify handling. Only a bearer the client itself supplied is caller-owned and exempt from stored state. Both synchronous and asynchronous stored-main substitution in `src/codex/auth-context.ts` remove a caller account header before copying the stored identity; an absent stored account ID leaves no account header. Caller-owned native Direct authentication retains its existing passthrough behavior. -Plain-main HTTP and WebSocket Responses refresh `__main__` only on the canonical OpenAI forward provider when the sent bearer and effective workspace match the main credential already observed under native ownership, the same equality rule used by the hard lock. `src/codex/auth-context.ts` captures a process-local dispatch proof after materialization; `src/server/responses/passthrough-delivery.ts` and `src/server/responses/core-codex-account.ts` recheck its credential/identity generations before publishing. Same-account token rotation and A→B→A changes reject old responses; WebSocket observers retain the original dispatch proof across frames and their projected HTTP response does not publish twice. Caller-owned requests acquire no physical-main read or Pool health state. +Plain-main HTTP and WebSocket Responses refresh `__main__` only on the canonical OpenAI forward provider when the sent bearer and effective workspace match the main credential already observed under native ownership, the same equality rule used by the hard lock. `src/codex/auth-context.ts` captures a process-local dispatch proof after materialization; `src/server/responses/passthrough-delivery.ts` and `src/server/responses/core-codex-account.ts` recheck its credential/identity generations before publishing. The dispatch also captures the process-wide credential mutation epoch; any OpenCodex-owned credential publication, including native main refresh or same-account reauth before a new quota credential observation, rejects an older dispatch. Publications for other credentials conservatively drop the main update as well. Same-account token rotation and A→B→A changes reject old responses; WebSocket observers retain the original dispatch proof across frames and their HTTP response does not publish twice, including separately marked precommit refusal projections whose prelude headers remain available to Pool replay. Caller-owned requests acquire no physical-main read or Pool health state. `src/providers/openai-sidecar.ts` releases quota-probe ownership on every materialization or usability failure before transferring a resolved context to its caller. Audio reports one terminal upstream outcome after validating the response body; redirects remain diff --git a/structure/providers/openai-tiers.md b/structure/providers/openai-tiers.md index 8308f52a147..ec64480c51b 100644 --- a/structure/providers/openai-tiers.md +++ b/structure/providers/openai-tiers.md @@ -435,6 +435,9 @@ an identical caller-owned credential. Materialization captures a process-local d publication rechecks identity and credential generations, including after an awaited HTTP import and for every WebSocket frame. A replaced credential, unmatched workspace, or custom destination cannot publish main usage. Pool health/failover handling stays scoped to Pool contexts. +The dispatch additionally fences the process-wide credential mutation epoch, so native main +refresh and same-account reauth commits reject an older response before quota observation catches +up. Publications for other credentials also conservatively drop the main update. `src/codex/auth-api/main-account-probe.ts` re-reads the bounded stored main credential and rechecks its writer, bearer and generation after body/retry awaits, before publishing main usage, credits, plan, reauth or Reserve state, including terminal 401/403 mutations. An unreadable file diff --git a/structure/transports/responses.md b/structure/transports/responses.md index 263d2599fff..b316cf1959d 100644 --- a/structure/transports/responses.md +++ b/structure/transports/responses.md @@ -23,7 +23,7 @@ Canonical forward auth retains its separate fixed credential/metadata allowlist; Retired Codex Spark has no model-specific tool or Responses Lite override; general Lite handling and namespace scrubbing remain shared compatibility behavior. Codex quota/reset evidence follows the -[shared/Reserve policy](../providers/openai-tiers.md#public-provider-contract), including suppression of retired model-derived evidence before shared recovery. Plain-main HTTP/WS quota headers refresh `__main__` only on the canonical OpenAI forward provider when the materialized bearer and workspace match the owned main observation, using the hard-lock credential-match rule. `src/server/responses/passthrough-delivery.ts` rechecks the captured credential generation after its awaited import; `src/server/responses/core-codex-account.ts` retains one dispatch proof per WS observer and rechecks it for each frame. Rotated credentials and unmatched workspaces publish nothing; WS-observed responses skip duplicate HTTP writes, and Pool health/failover gates remain unchanged. +[shared/Reserve policy](../providers/openai-tiers.md#public-provider-contract), including suppression of retired model-derived evidence before shared recovery. Plain-main HTTP/WS quota headers refresh `__main__` only on the canonical OpenAI forward provider when the materialized bearer and workspace match the owned main observation, using the hard-lock credential-match rule. `src/server/responses/passthrough-delivery.ts` rechecks the captured credential generation after its awaited import; `src/server/responses/core-codex-account.ts` retains one dispatch proof per WS observer and rechecks it for each frame. Rotated credentials and unmatched workspaces publish nothing; dispatch proofs also reject any OpenCodex-owned credential publication epoch change, including native main refresh or same-account reauth before quota re-observation; WS upstream responses and separately marked precommit refusal projections skip plain-main HTTP quota writes; refusal projections retain prelude headers for Pool replay, and Pool health/failover gates remain unchanged. ### Credential-bearing HTTP redirects diff --git a/tests/responses/responses-main-quota-observation.test.ts b/tests/responses/responses-main-quota-observation.test.ts index 91941365c26..287430423a0 100644 --- a/tests/responses/responses-main-quota-observation.test.ts +++ b/tests/responses/responses-main-quota-observation.test.ts @@ -4,9 +4,11 @@ import { join } from "node:path"; import { MAIN_CODEX_ACCOUNT_ID as MAIN } from "../../src/codex/account-id"; import * as mainCache from "../../src/codex/main-account-cache"; import { - materializeCodexUpstreamAuth, materializeCodexUpstreamAuthAsync, resolveCodexAuthContext, + CodexReserveUnavailableError, materializeCodexUpstreamAuth, materializeCodexUpstreamAuthAsync, resolveCodexAuthContext, type CodexAuthContext, } from "../../src/codex/auth-context"; +import { beginNativeMainReauth, forceRefreshMainAccountToken } from "../../src/codex/main-account"; +import { codexCredentialMutationEpoch } from "../../src/codex/credential-mutation-epoch"; import { resetMainCodexAccountIdentityTrackingForTests } from "../../src/codex/account-lifecycle"; import { clearAccountNeedsReauth } from "../../src/codex/account-runtime-state"; import { saveCodexAccountCredential } from "../../src/codex/account-store"; @@ -20,7 +22,13 @@ import { codexAccountSelectionForTurn, tryAdmitTurn } from "../../src/server/lif import { deliverPassthroughResponse } from "../../src/server/responses/passthrough-delivery"; import { codexWsQuotaObserver, retryCodexPoolOnAlternateAccount } from "../../src/server/responses/core-codex-account"; import { CodexWsMetadata } from "../../src/server/responses/codex-ws-metadata"; -import { markCodexWsResponse, isCodexWsQuotaObservedResponse } from "../../src/server/responses/codex-ws-wire"; +import { codexWsExchange } from "../../src/server/responses/codex-ws-exchange"; +import { CodexWsSession } from "../../src/server/responses/codex-ws-session"; +import { CODEX_RESPONSES_HTTP_URL, CODEX_RESPONSES_WS_URL, prepareCodexWsRequest } from "../../src/server/responses/codex-ws-request"; +import { NATIVE_RESERVE_MODEL } from "../../src/codex/catalog/native-models"; +import { getMainAccountHardLockStatus } from "../../src/codex/main-account-hard-lock"; +import { isCodexWsRejectionResponse } from "../../src/server/responses/ws-upstream"; +import { markCodexWsResponse, isCodexWsQuotaObservedResponse, isCodexWsUpstreamResponse } from "../../src/server/responses/codex-ws-wire"; import type { OcxConfig, OcxProviderConfig } from "../../src/types"; import { fakeChatGptJwt } from "../helpers/agent-task-recovery"; import { acquireOwnedSpendHome } from "../helpers/owned-spend-home"; @@ -96,7 +104,7 @@ async function deliver(ctx: CodexAuthContext, response = quotaResponse(), select responseCompletionCancelled: () => false } as Args[5], { upstreamResponse: response, upstream: new AbortController(), connectMs: 1000 } as Args[6], ); - expect(result.status).toBe(307); + expect(result.status).toBe(response.status); await result.body?.cancel(); } finally { budget.dispose(); } } @@ -208,6 +216,51 @@ describe("credential-bound plain-main Responses quota", () => { expect(getMainPolicyQuota()).toBeNull(); }); + for (const commit of ["refresh", "reauth"] as const) { + for (const transport of ["HTTP", "WS"] as const) { + test(`real native main ${commit} commit fences an older ${transport} quota dispatch`, async () => { + writeFileSync(join(root, "auth.json"), JSON.stringify({ tokens: { + access_token: bearer, refresh_token: "fixture-main-before-publication", account_id: ACCOUNT, + } })); + observe(); + const ctx = materialized(); + expect(ctx.mainQuotaDispatch).toBeDefined(); + const observer = transport === "WS" ? codexWsQuotaObserver(ctx, provider, "gpt-5.5") : undefined; + if (transport === "WS") expect(observer).toBeDefined(); + const epoch = codexCredentialMutationEpoch(); + expect(ctx.mainQuotaDispatch!.credentialMutationEpoch).toBe(epoch); + const generation = mainCache.getMainQuotaCredentialGeneration(); + const replacement = fakeChatGptJwt(ACCOUNT, { fixture_publication: commit }); + if (commit === "refresh") { + let refreshes = 0; + const result = await forceRefreshMainAccountToken(bearer, { + refreshToken: async () => { + refreshes++; + return { access: replacement, refresh: "fixture-main-after-refresh", expires: Date.now() + 3600_000, accountId: ACCOUNT }; + }, + }); + expect(refreshes).toBe(1); + expect(result?.accessToken === replacement).toBe(true); + } else { + const result = await beginNativeMainReauth().commit({ accessToken: replacement, + refreshToken: "fixture-main-after-reauth", idToken: "fixture-main-identity-token", chatgptAccountId: ACCOUNT }); + expect(result.chatgptAccountId).toBe(ACCOUNT); + } + const stored = JSON.parse(readFileSync(join(root, "auth.json"), "utf8")); + expect(stored.tokens.access_token === replacement).toBe(true); + expect(codexCredentialMutationEpoch()).toBeGreaterThan(epoch); + // No subsequent credential observation warmed the main quota generation. + expect(mainCache.getMainQuotaCredentialGeneration()).toBe(generation); + expect(mainCache.isMainQuotaWriterLive(ctx.mainQuotaDispatch!.writer)).toBe(true); + if (transport === "HTTP") await deliver(ctx); + else observer!(quotaHeaders()); + expect(getAccountQuota(MAIN)).toBeNull(); + expect(getMainPolicyQuota()).toBeNull(); + expect(mainCache.isMainQuotaDispatchLive(ctx.mainQuotaDispatch!)).toBe(false); + }); + } + } + test("5: non-canonical, key-auth and other-adapter providers cannot publish", async () => { observe(); const ctx = materialized(); @@ -238,6 +291,97 @@ describe("credential-bound plain-main Responses quota", () => { metadata.finish(); }); + test("a WS-marked unobserved refusal cannot overwrite newer main display or policy quota", async () => { + observe(); + const ctx = materialized(); + const metadata = new CodexWsMetadata(codexWsQuotaObserver(ctx, provider, "gpt-5.5")); + frame(metadata, 17); + frame(metadata, 99); + const display = getAccountQuota(MAIN); + const policy = getMainPolicyQuota(); + const refusal = Response.json({ error: { type: "invalid_request_error", message: "fixture refused create" } }, + { status: 400, headers: quotaHeaders("17") }); + markCodexWsResponse(refusal, false); + expect(isCodexWsUpstreamResponse(refusal)).toBe(true); + expect(isCodexWsQuotaObservedResponse(refusal)).toBe(false); + await deliver(ctx, refusal); + expect(getAccountQuota(MAIN)).toEqual(display); + expect(getMainPolicyQuota()).toEqual(policy); + expect(getMainAccountHardLockStatus({ codexMainAccountHardLock: true }).state).toBe("blocked"); + metadata.finish(); + }); + + test("WS precommit refusal through the real exchange cannot republish stale prelude quota", async () => { + observe(); + const ctx = materialized(); + const observer = codexWsQuotaObserver(ctx, provider, "gpt-5.5"); + expect(observer).toBeDefined(); + const originalSocket = globalThis.WebSocket; + class RefusalSocket extends EventTarget { + readyState = 0; + constructor() { + super(); + queueMicrotask(() => { this.readyState = 1; this.dispatchEvent(new Event("open")); }); + } + send(_text: string): void { + queueMicrotask(() => { + this.dispatchEvent(new MessageEvent("message", { data: JSON.stringify({ + type: "codex.rate_limits", rate_limits: { primary: { used_percent: 17, window_minutes: 10080 } }, + }) })); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(17); + // A concurrent main response publishes newer usage before this exchange refuses. + observer!(quotaHeaders("99")); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(99); + expect(getMainPolicyQuota()?.weeklyPercent).toBe(99); + this.dispatchEvent(new MessageEvent("message", { data: JSON.stringify({ + type: "error", status_code: 400, error: { type: "invalid_request_error", message: "fixture refused create" }, + }) })); + }); + } + close(): void { + if (this.readyState === 3) return; + this.readyState = 3; + this.dispatchEvent(new Event("close")); + } + } + globalThis.WebSocket = RefusalSocket as unknown as typeof WebSocket; + const session = new CodexWsSession(CODEX_RESPONSES_WS_URL, {}); + try { + expect(session.reserve()).toBe(true); + const init = { method: "POST", headers: caller(), body: JSON.stringify({ model: "gpt-5.5", stream: true, input: "hi" }) }; + const prepared = prepareCodexWsRequest(CODEX_RESPONSES_HTTP_URL, init); + expect(prepared).not.toBeNull(); + const refusal = await codexWsExchange({ session, url: CODEX_RESPONSES_HTTP_URL, init, prepared: prepared!, + sseFallback: globalThis.fetch, onQuota: observer, bunVersion: "1.4.0" }); + expect(refusal.status).toBe(400); + expect(isCodexWsRejectionResponse(refusal)).toBe(true); + expect(isCodexWsUpstreamResponse(refusal)).toBe(false); + expect(isCodexWsRejectionResponse(quotaResponse())).toBe(false); + expect(isCodexWsQuotaObservedResponse(refusal)).toBe(false); + expect(refusal.headers.get("x-codex-primary-used-percent")).toBe("17"); + expect(getMainAccountHardLockStatus({ codexMainAccountHardLock: true }).state).toBe("blocked"); + await deliver(ctx, refusal); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(99); + expect(getMainPolicyQuota()?.weeklyPercent).toBe(99); + expect(getMainAccountHardLockStatus({ codexMainAccountHardLock: true }).state).toBe("blocked"); + } finally { + session.dispose(); + globalThis.WebSocket = originalSocket; + } + }); + + test("async Reserve admission failure clears a reused main dispatch proof", async () => { + observe(); + const ctx = materialized(); + expect(ctx.mainQuotaDispatch).toBeDefined(); + expect(mainCache.isMainQuotaDispatchLive(ctx.mainQuotaDispatch!)).toBe(true); + await expect(materializeCodexUpstreamAuthAsync(caller(), ctx, { + config: { ...config, codexDesktopAuthless: true, pausedCodexAccountIds: [MAIN] }, + modelId: NATIVE_RESERVE_MODEL, admission: { source: "loopback" }, + })).rejects.toBeInstanceOf(CodexReserveUnavailableError); + expect(ctx.mainQuotaDispatch).toBeUndefined(); + }); + test("7: stored pool HTTP and WS still update only their pool row", async () => { const generation = saveCodexAccountCredential(POOL, { accessToken: "fixture-pool-token", refreshToken: "fixture-pool-refresh", chatgptAccountId: "fixture-pool-workspace", expiresAt: Date.now() + 3600_000 }); @@ -329,7 +473,7 @@ describe("credential-bound plain-main Responses quota", () => { expect(persisted.quotas[MAIN].weeklyPercent).toBe(23); expect(persisted.quotas[MAIN].updatedAt).toBeGreaterThanOrEqual(before); expect(persisted.quotas[MAIN].updatedAt).toBeLessThanOrEqual(Date.now()); - for (const forbidden of [bearer, ACCOUNT, "bearerHmac", "mainQuotaDispatch", "credentialGeneration", "configGeneration", "identityGeneration"]) + for (const forbidden of [bearer, ACCOUNT, "bearerHmac", "mainQuotaDispatch", "credentialGeneration", "credentialMutationEpoch", "configGeneration", "identityGeneration"]) expect(body.includes(forbidden)).toBe(false); expect(Object.keys(persisted.mainPolicyQuota).sort()).toEqual(["identityKey", "quota"]); }); From 18f8cc01a362fee2eeef4721f288f5765ee352db Mon Sep 17 00:00:00 2001 From: JUN Date: Fri, 9 Oct 2026 21:12:42 +0900 Subject: [PATCH 3/7] fix(quota): keep WS prelude projections out of plain-main quota publication Review follow-up for #6800: a pre-response 502/504 from a WebSocket exchange carries the prelude snapshot just like a precommit refusal, so delivery could replay stale main quota over a newer reading. Both projections now carry one prelude-projection marker that the plain-main HTTP branch skips; real HTTP fallbacks after a failed upgrade still publish. --- src/server/responses/codex-ws-exchange.ts | 5 +- src/server/responses/codex-ws-wire.ts | 13 +- src/server/responses/passthrough-delivery.ts | 6 +- src/server/responses/ws-upstream.ts | 2 +- structure/providers/openai-accounts.md | 2 +- structure/transports/responses.md | 2 +- .../responses-main-quota-observation.test.ts | 184 ++++++++++++------ 7 files changed, 142 insertions(+), 72 deletions(-) diff --git a/src/server/responses/codex-ws-exchange.ts b/src/server/responses/codex-ws-exchange.ts index 0944a074c15..849c7184198 100644 --- a/src/server/responses/codex-ws-exchange.ts +++ b/src/server/responses/codex-ws-exchange.ts @@ -9,7 +9,7 @@ import { CODEX_RESPONSES_HTTP_URL, type PreparedCodexWsRequest } from "./codex-w import { CodexWsCorrelation } from "./codex-ws-correlation"; import type { CodexWsSession } from "./codex-ws-session"; import { UPGRADE_DEADLINE_MS, CODEX_WS_LIVENESS_PING_INTERVAL_MS, CODEX_WS_RESPONSE_PRELUDE_TIMEOUT_MS, MAX_CODEX_WS_FRAME_BYTES, - MAX_CODEX_WS_QUEUE_BYTES, markCodexWsResponse, markCodexWsRejectionResponse, normalizeResponsesWsRelayEvent, closedBeforeTerminalMessage, + MAX_CODEX_WS_QUEUE_BYTES, markCodexWsResponse, markCodexWsPreludeProjection, normalizeResponsesWsRelayEvent, closedBeforeTerminalMessage, codexWsCreateFrameExceedsLimit, codexWsFailureDetail, codexWsPreResponseFailure, markCodexWsStage, codexWsOcxVersion, markCodexWsSocketDeath, type CodexWsFailureStage, type CodexWsStageRecord } from "./codex-ws-wire"; @@ -242,6 +242,7 @@ export function codexWsExchange(options: ExchangeOptions): Promise { // are left out: their channel may already have sent continuation frames on this socket, so // the create frame alone no longer describes the turn. if (socketDied && !nativeControl) markCodexWsSocketDeath(failureResponse, stage); + markCodexWsPreludeProjection(failureResponse); resolve(failureResponse); return; } @@ -514,7 +515,7 @@ export function codexWsExchange(options: ExchangeOptions): Promise { cleanup(); try { controller.close(); } catch { /* unused stream already closed */ } session.dispose(); - markCodexWsRejectionResponse(rejection); + markCodexWsPreludeProjection(rejection); resolve(rejection); return; } diff --git a/src/server/responses/codex-ws-wire.ts b/src/server/responses/codex-ws-wire.ts index dccfbbeff7b..592d2ef32ad 100644 --- a/src/server/responses/codex-ws-wire.ts +++ b/src/server/responses/codex-ws-wire.ts @@ -49,7 +49,7 @@ const WS_CLOSE_MESSAGE_TOO_BIG = 1009; const codexWsUpstreamResponses = new WeakSet(); const quotaObservedResponses = new WeakSet(); -const codexWsRejectionResponses = new WeakSet(); +const codexWsPreludeProjections = new WeakSet(); /** Quota arrived directly at its captured account; do not replay old HTTP prelude headers. */ export function isCodexWsQuotaObservedResponse(response: Response): boolean { @@ -61,14 +61,13 @@ export function isCodexWsUpstreamResponse(response: Response): boolean { return codexWsUpstreamResponses.has(response); } - -/** Precommit refusal projection; its headers retain the exchange's prelude snapshot. */ -export function markCodexWsRejectionResponse(response: Response): void { - codexWsRejectionResponses.add(response); +/** Pre-response projection; its headers retain the exchange's prelude snapshot. */ +export function markCodexWsPreludeProjection(response: Response): void { + codexWsPreludeProjections.add(response); } -export function isCodexWsRejectionResponse(response: Response): boolean { - return codexWsRejectionResponses.has(response); +export function isCodexWsPreludeProjection(response: Response): boolean { + return codexWsPreludeProjections.has(response); } export function markCodexWsResponse(response: Response, observed: boolean): void { diff --git a/src/server/responses/passthrough-delivery.ts b/src/server/responses/passthrough-delivery.ts index 1ecdfe44733..3734fde3e08 100644 --- a/src/server/responses/passthrough-delivery.ts +++ b/src/server/responses/passthrough-delivery.ts @@ -33,7 +33,7 @@ import { import { isMainQuotaDispatchLive } from "../../codex/main-account-cache"; import { MAIN_CODEX_ACCOUNT_ID } from "../../codex/account-id"; import type { ResponsesTerminalStatus } from "../../bridge"; -import { isCodexWsQuotaObservedResponse, isCodexWsUpstreamResponse, isCodexWsRejectionResponse } from "./ws-upstream"; +import { isCodexWsQuotaObservedResponse, isCodexWsUpstreamResponse, isCodexWsPreludeProjection } from "./ws-upstream"; import { recordSubagentQuotaFailureForThreadSpawn } from "../../codex/subagent-model-fallback"; import { recordCodexUpstreamOutcome } from "../../codex/routing"; import { codexProbeLeaseId, codexProbeQuotaScope, codexTransientProbeGrant, releaseCodexAuthContextProbeLease } from "../../codex/auth-context"; @@ -541,8 +541,8 @@ export async function deliverPassthroughResponse( } else { const mainDispatch = liveMainQuotaDispatch(admissionState.authCtx, route.provider); // The WS observer is the only plain-main publisher for a WebSocket exchange; - // a refusal projection carries the prelude snapshot, not fresh evidence. - if (mainDispatch && !(isCodexWsUpstreamResponse(upstreamResponse) || isCodexWsRejectionResponse(upstreamResponse))) { + // a prelude projection carries the prelude snapshot, not fresh evidence. + if (mainDispatch && !(isCodexWsUpstreamResponse(upstreamResponse) || isCodexWsPreludeProjection(upstreamResponse))) { const { applyAccountQuotaFromUpstreamHeaders } = await import("../../codex/auth-api"); // Import yields; same-account token replacement leaves the identity writer live. // Re-check the credential fence with no await before publication. diff --git a/src/server/responses/ws-upstream.ts b/src/server/responses/ws-upstream.ts index d2569e536af..186302d2f1b 100644 --- a/src/server/responses/ws-upstream.ts +++ b/src/server/responses/ws-upstream.ts @@ -26,7 +26,7 @@ import { codexWsCreateFrameExceedsLimit } from "./codex-ws-wire"; import { isLoopbackUrl, rewriteWebSocketDial } from "../../plugins/upstream-hooks"; export { CODEX_WS_LIVENESS_PING_INTERVAL_MS, CODEX_WS_RESPONSE_PRELUDE_TIMEOUT_MS, MAX_CODEX_WS_FRAME_BYTES, MAX_CODEX_WS_QUEUE_BYTES, MAX_CODEX_WS_CREATE_FRAME_BYTES, CODEX_WS_CREATE_FRAME_LIMIT_BYTES, codexWsCreateFrameExceedsLimit, - isCodexWsQuotaObservedResponse, isCodexWsUpstreamResponse, isCodexWsRejectionResponse } from "./codex-ws-wire"; + isCodexWsQuotaObservedResponse, isCodexWsUpstreamResponse, isCodexWsPreludeProjection } from "./codex-ws-wire"; export const MIN_BOUNDED_CODEX_WS_BUN_VERSION = "1.4.0"; /** diff --git a/structure/providers/openai-accounts.md b/structure/providers/openai-accounts.md index 64a68b63fce..0f798896275 100644 --- a/structure/providers/openai-accounts.md +++ b/structure/providers/openai-accounts.md @@ -211,7 +211,7 @@ request-scoped: a translated Claude turn resolves through Pool selection like an main keeps its health, quarantine and refresh-and-classify handling. Only a bearer the client itself supplied is caller-owned and exempt from stored state. Both synchronous and asynchronous stored-main substitution in `src/codex/auth-context.ts` remove a caller account header before copying the stored identity; an absent stored account ID leaves no account header. Caller-owned native Direct authentication retains its existing passthrough behavior. -Plain-main HTTP and WebSocket Responses refresh `__main__` only on the canonical OpenAI forward provider when the sent bearer and effective workspace match the main credential already observed under native ownership, the same equality rule used by the hard lock. `src/codex/auth-context.ts` captures a process-local dispatch proof after materialization; `src/server/responses/passthrough-delivery.ts` and `src/server/responses/core-codex-account.ts` recheck its credential/identity generations before publishing. The dispatch also captures the process-wide credential mutation epoch; any OpenCodex-owned credential publication, including native main refresh or same-account reauth before a new quota credential observation, rejects an older dispatch. Publications for other credentials conservatively drop the main update as well. Same-account token rotation and A→B→A changes reject old responses; WebSocket observers retain the original dispatch proof across frames and their HTTP response does not publish twice, including separately marked precommit refusal projections whose prelude headers remain available to Pool replay. Caller-owned requests acquire no physical-main read or Pool health state. +Plain-main HTTP and WebSocket Responses refresh `__main__` only on the canonical OpenAI forward provider when the sent bearer and effective workspace match the main credential already observed under native ownership, the same equality rule used by the hard lock. `src/codex/auth-context.ts` captures a process-local dispatch proof after materialization; `src/server/responses/passthrough-delivery.ts` and `src/server/responses/core-codex-account.ts` recheck its credential/identity generations before publishing. The dispatch also captures the process-wide credential mutation epoch; any OpenCodex-owned credential publication, including native main refresh or same-account reauth before a new quota credential observation, rejects an older dispatch. Publications for other credentials conservatively drop the main update as well. Same-account token rotation and A→B→A changes reject old responses; WebSocket observers retain the original dispatch proof across frames and their HTTP response does not publish twice, including separately marked pre-response prelude projections (4xx refusals and 502/504 gateway failures). Prelude headers remain available to Pool replay; real HTTP fallback responses still publish through HTTP delivery. Caller-owned requests acquire no physical-main read or Pool health state. `src/providers/openai-sidecar.ts` releases quota-probe ownership on every materialization or usability failure before transferring a resolved context to its caller. Audio reports one terminal upstream outcome after validating the response body; redirects remain diff --git a/structure/transports/responses.md b/structure/transports/responses.md index b316cf1959d..3eb0bdacea9 100644 --- a/structure/transports/responses.md +++ b/structure/transports/responses.md @@ -23,7 +23,7 @@ Canonical forward auth retains its separate fixed credential/metadata allowlist; Retired Codex Spark has no model-specific tool or Responses Lite override; general Lite handling and namespace scrubbing remain shared compatibility behavior. Codex quota/reset evidence follows the -[shared/Reserve policy](../providers/openai-tiers.md#public-provider-contract), including suppression of retired model-derived evidence before shared recovery. Plain-main HTTP/WS quota headers refresh `__main__` only on the canonical OpenAI forward provider when the materialized bearer and workspace match the owned main observation, using the hard-lock credential-match rule. `src/server/responses/passthrough-delivery.ts` rechecks the captured credential generation after its awaited import; `src/server/responses/core-codex-account.ts` retains one dispatch proof per WS observer and rechecks it for each frame. Rotated credentials and unmatched workspaces publish nothing; dispatch proofs also reject any OpenCodex-owned credential publication epoch change, including native main refresh or same-account reauth before quota re-observation; WS upstream responses and separately marked precommit refusal projections skip plain-main HTTP quota writes; refusal projections retain prelude headers for Pool replay, and Pool health/failover gates remain unchanged. +[shared/Reserve policy](../providers/openai-tiers.md#public-provider-contract), including suppression of retired model-derived evidence before shared recovery. Plain-main HTTP/WS quota headers refresh `__main__` only on the canonical OpenAI forward provider when the materialized bearer and workspace match the owned main observation, using the hard-lock credential-match rule. `src/server/responses/passthrough-delivery.ts` rechecks the captured credential generation after its awaited import; `src/server/responses/core-codex-account.ts` retains one dispatch proof per WS observer and rechecks it for each frame. Rotated credentials and unmatched workspaces publish nothing; dispatch proofs also reject any OpenCodex-owned credential publication epoch change, including native main refresh or same-account reauth before quota re-observation; WS upstream responses and separately marked pre-response prelude projections (4xx refusals and 502/504 gateway failures) skip plain-main HTTP quota writes; prelude headers remain available to Pool replay, real HTTP fallbacks still publish, and Pool health/failover gates remain unchanged. ### Credential-bearing HTTP redirects diff --git a/tests/responses/responses-main-quota-observation.test.ts b/tests/responses/responses-main-quota-observation.test.ts index 287430423a0..47871c0af84 100644 --- a/tests/responses/responses-main-quota-observation.test.ts +++ b/tests/responses/responses-main-quota-observation.test.ts @@ -27,8 +27,8 @@ import { CodexWsSession } from "../../src/server/responses/codex-ws-session"; import { CODEX_RESPONSES_HTTP_URL, CODEX_RESPONSES_WS_URL, prepareCodexWsRequest } from "../../src/server/responses/codex-ws-request"; import { NATIVE_RESERVE_MODEL } from "../../src/codex/catalog/native-models"; import { getMainAccountHardLockStatus } from "../../src/codex/main-account-hard-lock"; -import { isCodexWsRejectionResponse } from "../../src/server/responses/ws-upstream"; -import { markCodexWsResponse, isCodexWsQuotaObservedResponse, isCodexWsUpstreamResponse } from "../../src/server/responses/codex-ws-wire"; +import { CODEX_WS_RESPONSE_PRELUDE_TIMEOUT_MS, isCodexWsPreludeProjection } from "../../src/server/responses/ws-upstream"; +import { UPGRADE_DEADLINE_MS, markCodexWsResponse, isCodexWsQuotaObservedResponse, isCodexWsUpstreamResponse } from "../../src/server/responses/codex-ws-wire"; import type { OcxConfig, OcxProviderConfig } from "../../src/types"; import { fakeChatGptJwt } from "../helpers/agent-task-recovery"; import { acquireOwnedSpendHome } from "../helpers/owned-spend-home"; @@ -49,6 +49,8 @@ let config: OcxConfig; let releaseSpend: (() => void) | undefined; let pendingPersist: { run: () => void; timer: ReturnType } | undefined; let clock: ReturnType; +let pendingPreludeTimeout: (() => void) | undefined; +let pendingUpgradeTimeout: (() => void) | undefined; // Exercise the real serializer without a wall-clock race, as main-quota-provenance does. function installPersistenceClock() { @@ -56,6 +58,14 @@ function installPersistenceClock() { return spyOn(globalThis, "setTimeout").mockImplementation((( callback: (...args: unknown[]) => void, delay?: number, ...args: unknown[] ) => { + if (delay === UPGRADE_DEADLINE_MS) { + pendingUpgradeTimeout = () => callback(...args); + return nativeTimeout(() => {}, delay); + } + if (delay === CODEX_WS_RESPONSE_PRELUDE_TIMEOUT_MS) { + pendingPreludeTimeout = () => callback(...args); + return nativeTimeout(() => {}, delay); + } if (delay !== 250) return nativeTimeout(callback, delay, ...args); const timer = nativeTimeout(() => {}, 60_000); pendingPersist = { run: () => callback(...args), timer }; @@ -133,6 +143,8 @@ beforeEach(() => { bearer = fakeChatGptJwt(ACCOUNT); config = { providers: { openai: provider }, codexAccounts: [], codexMainAccountHardLock: false } as OcxConfig; pendingPersist = undefined; + pendingPreludeTimeout = undefined; + pendingUpgradeTimeout = undefined; clock = installPersistenceClock(); globalThis.fetch = (async () => { throw new Error("unexpected network call"); }) as typeof fetch; }); @@ -142,6 +154,8 @@ afterEach(() => { clearAccountQuota(); if (pendingPersist) clearTimeout(pendingPersist.timer); pendingPersist = undefined; + pendingPreludeTimeout = undefined; + pendingUpgradeTimeout = undefined; clock.mockRestore(); clearCodexUpstreamHealth(); clearThreadAccountMap(); clearPoolRotationState(); clearAccountNeedsReauth(MAIN); clearAccountNeedsReauth(POOL); @@ -311,64 +325,120 @@ describe("credential-bound plain-main Responses quota", () => { metadata.finish(); }); - test("WS precommit refusal through the real exchange cannot republish stale prelude quota", async () => { - observe(); - const ctx = materialized(); - const observer = codexWsQuotaObserver(ctx, provider, "gpt-5.5"); - expect(observer).toBeDefined(); - const originalSocket = globalThis.WebSocket; - class RefusalSocket extends EventTarget { - readyState = 0; - constructor() { - super(); - queueMicrotask(() => { this.readyState = 1; this.dispatchEvent(new Event("open")); }); + for (const failure of ["4xx refusal", "socket closure", "prelude timeout", "connect timeout"] as const) { + test(`WS ${failure} through the real exchange cannot republish stale prelude quota`, async () => { + observe(); + const ctx = materialized(); + const observer = codexWsQuotaObserver(ctx, provider, "gpt-5.5"); + expect(observer).toBeDefined(); + const originalSocket = globalThis.WebSocket; + const controller = new AbortController(); + class RefusalSocket extends EventTarget { + readyState = 0; + constructor() { + super(); + queueMicrotask(() => { this.readyState = 1; this.dispatchEvent(new Event("open")); }); + } + send(_text: string): void { + queueMicrotask(() => { + this.dispatchEvent(new MessageEvent("message", { data: JSON.stringify({ + type: "codex.rate_limits", rate_limits: { primary: { used_percent: 17, window_minutes: 10080 } }, + }) })); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(17); + // A concurrent main response publishes newer usage before this exchange fails. + observer!(quotaHeaders("99")); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(99); + expect(getMainPolicyQuota()?.weeklyPercent).toBe(99); + if (failure === "socket closure") this.close(); + else if (failure === "prelude timeout") { + expect(pendingPreludeTimeout).toBeDefined(); + pendingPreludeTimeout!(); + } else if (failure === "connect timeout") controller.abort(new DOMException("fixture deadline", "TimeoutError")); + else this.dispatchEvent(new MessageEvent("message", { data: JSON.stringify({ + type: "error", status_code: 400, error: { type: "invalid_request_error", message: "fixture refused create" }, + }) })); + }); + } + close(): void { + if (this.readyState === 3) return; + this.readyState = 3; + this.dispatchEvent(new Event("close")); + } } - send(_text: string): void { - queueMicrotask(() => { - this.dispatchEvent(new MessageEvent("message", { data: JSON.stringify({ - type: "codex.rate_limits", rate_limits: { primary: { used_percent: 17, window_minutes: 10080 } }, - }) })); - expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(17); - // A concurrent main response publishes newer usage before this exchange refuses. - observer!(quotaHeaders("99")); - expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(99); - expect(getMainPolicyQuota()?.weeklyPercent).toBe(99); - this.dispatchEvent(new MessageEvent("message", { data: JSON.stringify({ - type: "error", status_code: 400, error: { type: "invalid_request_error", message: "fixture refused create" }, - }) })); - }); + globalThis.WebSocket = RefusalSocket as unknown as typeof WebSocket; + const session = new CodexWsSession(CODEX_RESPONSES_WS_URL, {}); + try { + expect(session.reserve()).toBe(true); + const init = { method: "POST", signal: controller.signal, headers: caller(), body: JSON.stringify({ model: "gpt-5.5", stream: true, input: "hi" }) }; + const prepared = prepareCodexWsRequest(CODEX_RESPONSES_HTTP_URL, init); + expect(prepared).not.toBeNull(); + const refusal = await codexWsExchange({ session, url: CODEX_RESPONSES_HTTP_URL, init, prepared: prepared!, + sseFallback: globalThis.fetch, onQuota: observer, bunVersion: "1.4.0" }); + expect(refusal.status).toBe(failure === "4xx refusal" ? 400 : failure === "socket closure" ? 502 : 504); + expect(isCodexWsUpstreamResponse(refusal)).toBe(false); + expect(isCodexWsPreludeProjection(quotaResponse())).toBe(false); + expect(isCodexWsQuotaObservedResponse(refusal)).toBe(false); + expect(refusal.headers.get("x-codex-primary-used-percent")).toBe("17"); + expect(getMainAccountHardLockStatus({ codexMainAccountHardLock: true }).state).toBe("blocked"); + await deliver(ctx, refusal); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(99); + expect(getMainPolicyQuota()?.weeklyPercent).toBe(99); + expect(getMainAccountHardLockStatus({ codexMainAccountHardLock: true }).state).toBe("blocked"); + expect(isCodexWsPreludeProjection(refusal)).toBe(true); + } finally { + session.dispose(); + globalThis.WebSocket = originalSocket; } - close(): void { - if (this.readyState === 3) return; - this.readyState = 3; - this.dispatchEvent(new Event("close")); + }); + } + + for (const failure of ["close", "error", "upgrade timeout", "send exception"] as const) { + test(`real HTTP fallback after WS ${failure} still publishes main quota`, async () => { + observe(); + const ctx = materialized(); + const originalSocket = globalThis.WebSocket; + class UpgradeFailureSocket extends EventTarget { + readyState = 0; + constructor() { + super(); + queueMicrotask(() => { + if (failure === "upgrade timeout") { expect(pendingUpgradeTimeout).toBeDefined(); pendingUpgradeTimeout!(); } + else if (failure === "send exception") { this.readyState = 1; this.dispatchEvent(new Event("open")); } + else if (failure === "close") this.close(); + else this.dispatchEvent(new Event("error")); + }); + } + send(): void { throw new Error("fixture unsent create"); } + close(): void { + if (this.readyState === 3) return; + this.readyState = 3; + this.dispatchEvent(new Event("close")); + } } - } - globalThis.WebSocket = RefusalSocket as unknown as typeof WebSocket; - const session = new CodexWsSession(CODEX_RESPONSES_WS_URL, {}); - try { - expect(session.reserve()).toBe(true); - const init = { method: "POST", headers: caller(), body: JSON.stringify({ model: "gpt-5.5", stream: true, input: "hi" }) }; - const prepared = prepareCodexWsRequest(CODEX_RESPONSES_HTTP_URL, init); - expect(prepared).not.toBeNull(); - const refusal = await codexWsExchange({ session, url: CODEX_RESPONSES_HTTP_URL, init, prepared: prepared!, - sseFallback: globalThis.fetch, onQuota: observer, bunVersion: "1.4.0" }); - expect(refusal.status).toBe(400); - expect(isCodexWsRejectionResponse(refusal)).toBe(true); - expect(isCodexWsUpstreamResponse(refusal)).toBe(false); - expect(isCodexWsRejectionResponse(quotaResponse())).toBe(false); - expect(isCodexWsQuotaObservedResponse(refusal)).toBe(false); - expect(refusal.headers.get("x-codex-primary-used-percent")).toBe("17"); - expect(getMainAccountHardLockStatus({ codexMainAccountHardLock: true }).state).toBe("blocked"); - await deliver(ctx, refusal); - expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(99); - expect(getMainPolicyQuota()?.weeklyPercent).toBe(99); - expect(getMainAccountHardLockStatus({ codexMainAccountHardLock: true }).state).toBe("blocked"); - } finally { - session.dispose(); - globalThis.WebSocket = originalSocket; - } - }); + globalThis.WebSocket = UpgradeFailureSocket as unknown as typeof WebSocket; + const session = new CodexWsSession(CODEX_RESPONSES_WS_URL, {}); + let fallbackCalls = 0; + try { + expect(session.reserve()).toBe(true); + const init = { method: "POST", headers: caller(), body: JSON.stringify({ model: "gpt-5.5", stream: true, input: "hi" }) }; + const prepared = prepareCodexWsRequest(CODEX_RESPONSES_HTTP_URL, init); + expect(prepared).not.toBeNull(); + const response = await codexWsExchange({ session, url: CODEX_RESPONSES_HTTP_URL, init, prepared: prepared!, + sseFallback: (async () => { fallbackCalls++; return quotaResponse(quotaHeaders("37")); }) as typeof fetch, + onQuota: codexWsQuotaObserver(ctx, provider, "gpt-5.5"), bunVersion: "1.4.0" }); + expect(fallbackCalls).toBe(1); + expect(isCodexWsUpstreamResponse(response)).toBe(false); + expect(isCodexWsPreludeProjection(response)).toBe(false); + expect(isCodexWsQuotaObservedResponse(response)).toBe(false); + await deliver(ctx, response); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(37); + expect(getMainPolicyQuota()?.weeklyPercent).toBe(37); + } finally { + session.dispose(); + globalThis.WebSocket = originalSocket; + } + }); + } test("async Reserve admission failure clears a reused main dispatch proof", async () => { observe(); From cecd1921502c64cf43d36014616da372ac3c2152 Mon Sep 17 00:00:00 2001 From: JUN Date: Fri, 9 Oct 2026 21:23:14 +0900 Subject: [PATCH 4/7] fix(quota): let the WS observer claim its main quota dispatch Review follow-up for #6800: response wrappers such as the combo stream preflight replace the committed WebSocket Response, so a response-object marker cannot be the only guard. The plain-main WS observer now claims its dispatch before publishing, and HTTP delivery skips any claimed dispatch. A real HTTP fallback that never saw a WebSocket quota frame still publishes. --- src/codex/main-account-cache.ts | 12 ++ src/server/responses/core-codex-account.ts | 3 +- src/server/responses/passthrough-delivery.ts | 6 +- structure/providers/openai-accounts.md | 2 +- structure/providers/openai-tiers.md | 3 + structure/transports/responses.md | 2 +- .../responses-main-quota-observation.test.ts | 112 ++++++++++++++++++ 7 files changed, 135 insertions(+), 5 deletions(-) diff --git a/src/codex/main-account-cache.ts b/src/codex/main-account-cache.ts index e6551df38ff..41e7a9c5d89 100644 --- a/src/codex/main-account-cache.ts +++ b/src/codex/main-account-cache.ts @@ -82,6 +82,18 @@ export type MainQuotaDispatch = Readonly<{ configGeneration: number; }>; +// WS quota frames publish through their observer; prelude quota can only come from those frames. +// A real HTTP fallback after a failed upgrade never invokes the observer and stays unclaimed. +const wsObservedMainDispatches = new WeakSet(); + +export function claimMainQuotaDispatchForWs(dispatch: MainQuotaDispatch): void { + wsObservedMainDispatches.add(dispatch); +} + +export function isMainQuotaDispatchWsClaimed(dispatch: MainQuotaDispatch): boolean { + return wsObservedMainDispatches.has(dispatch); +} + export function captureMainQuotaDispatch( accessToken: string, accountId: string | undefined, configGeneration: number, ): MainQuotaDispatch | undefined { diff --git a/src/server/responses/core-codex-account.ts b/src/server/responses/core-codex-account.ts index da7c3b3f213..a138b13209d 100644 --- a/src/server/responses/core-codex-account.ts +++ b/src/server/responses/core-codex-account.ts @@ -6,7 +6,7 @@ import { computeQuotaCooldown, formatCodexProviderForLog, } from "../../codex/routing"; -import { isMainQuotaDispatchLive, type MainQuotaDispatch } from "../../codex/main-account-cache"; +import { claimMainQuotaDispatchForWs, isMainQuotaDispatchLive, type MainQuotaDispatch } from "../../codex/main-account-cache"; import type { CodexWsQuotaObserver } from "./codex-ws-metadata"; import { isCanonicalOpenAiForwardProvider } from "../../providers/openai-tiers"; import { isCodexAccountGenerationLive } from "../../codex/account-store"; @@ -138,6 +138,7 @@ export function codexWsQuotaObserver(authCtx: CodexAuthContext, provider: OcxPro const dispatch = liveMainQuotaDispatch(authCtx, provider); if (!dispatch) return undefined; return headers => { + claimMainQuotaDispatchForWs(dispatch); if (!isMainQuotaDispatchLive(dispatch)) return; applyCapturedCodexQuota(MAIN_CODEX_ACCOUNT_ID, headers, dispatch.configGeneration, dispatch.writer, { modelId }); }; diff --git a/src/server/responses/passthrough-delivery.ts b/src/server/responses/passthrough-delivery.ts index 3734fde3e08..e17f93da8f2 100644 --- a/src/server/responses/passthrough-delivery.ts +++ b/src/server/responses/passthrough-delivery.ts @@ -30,7 +30,7 @@ import { isFixedCodexAccount, shouldDeferCodexResetDerivedCooldown, } from "./core-codex-account"; -import { isMainQuotaDispatchLive } from "../../codex/main-account-cache"; +import { isMainQuotaDispatchLive, isMainQuotaDispatchWsClaimed } from "../../codex/main-account-cache"; import { MAIN_CODEX_ACCOUNT_ID } from "../../codex/account-id"; import type { ResponsesTerminalStatus } from "../../bridge"; import { isCodexWsQuotaObservedResponse, isCodexWsUpstreamResponse, isCodexWsPreludeProjection } from "./ws-upstream"; @@ -542,7 +542,9 @@ export async function deliverPassthroughResponse( const mainDispatch = liveMainQuotaDispatch(admissionState.authCtx, route.provider); // The WS observer is the only plain-main publisher for a WebSocket exchange; // a prelude projection carries the prelude snapshot, not fresh evidence. - if (mainDispatch && !(isCodexWsUpstreamResponse(upstreamResponse) || isCodexWsPreludeProjection(upstreamResponse))) { + // The dispatch claim is authoritative because downstream wrappers can replace the Response. + if (mainDispatch && !isMainQuotaDispatchWsClaimed(mainDispatch) + && !(isCodexWsUpstreamResponse(upstreamResponse) || isCodexWsPreludeProjection(upstreamResponse))) { const { applyAccountQuotaFromUpstreamHeaders } = await import("../../codex/auth-api"); // Import yields; same-account token replacement leaves the identity writer live. // Re-check the credential fence with no await before publication. diff --git a/structure/providers/openai-accounts.md b/structure/providers/openai-accounts.md index 0f798896275..7bd72b05ca0 100644 --- a/structure/providers/openai-accounts.md +++ b/structure/providers/openai-accounts.md @@ -211,7 +211,7 @@ request-scoped: a translated Claude turn resolves through Pool selection like an main keeps its health, quarantine and refresh-and-classify handling. Only a bearer the client itself supplied is caller-owned and exempt from stored state. Both synchronous and asynchronous stored-main substitution in `src/codex/auth-context.ts` remove a caller account header before copying the stored identity; an absent stored account ID leaves no account header. Caller-owned native Direct authentication retains its existing passthrough behavior. -Plain-main HTTP and WebSocket Responses refresh `__main__` only on the canonical OpenAI forward provider when the sent bearer and effective workspace match the main credential already observed under native ownership, the same equality rule used by the hard lock. `src/codex/auth-context.ts` captures a process-local dispatch proof after materialization; `src/server/responses/passthrough-delivery.ts` and `src/server/responses/core-codex-account.ts` recheck its credential/identity generations before publishing. The dispatch also captures the process-wide credential mutation epoch; any OpenCodex-owned credential publication, including native main refresh or same-account reauth before a new quota credential observation, rejects an older dispatch. Publications for other credentials conservatively drop the main update as well. Same-account token rotation and A→B→A changes reject old responses; WebSocket observers retain the original dispatch proof across frames and their HTTP response does not publish twice, including separately marked pre-response prelude projections (4xx refusals and 502/504 gateway failures). Prelude headers remain available to Pool replay; real HTTP fallback responses still publish through HTTP delivery. Caller-owned requests acquire no physical-main read or Pool health state. +Plain-main HTTP and WebSocket Responses refresh `__main__` only on the canonical OpenAI forward provider when the sent bearer and effective workspace match the main credential already observed under native ownership, the same equality rule used by the hard lock. `src/codex/auth-context.ts` captures a process-local dispatch proof after materialization; `src/server/responses/passthrough-delivery.ts` and `src/server/responses/core-codex-account.ts` recheck its credential/identity generations before publishing. The dispatch also captures the process-wide credential mutation epoch; any OpenCodex-owned credential publication, including native main refresh or same-account reauth before a new quota credential observation, rejects an older dispatch. Publications for other credentials conservatively drop the main update as well. Same-account token rotation and A→B→A changes reject old responses; WebSocket observers retain the original dispatch proof across frames and claim it on every invocation before checking liveness. This process-local claim prevents plain-main HTTP publication even when a downstream stream wrapper replaces the Response; response markers remain an additional guard, including separately marked pre-response prelude projections (4xx refusals and 502/504 gateway failures). Prelude headers remain available to Pool replay; real HTTP fallback responses still publish through HTTP delivery. Caller-owned requests acquire no physical-main read or Pool health state. `src/providers/openai-sidecar.ts` releases quota-probe ownership on every materialization or usability failure before transferring a resolved context to its caller. Audio reports one terminal upstream outcome after validating the response body; redirects remain diff --git a/structure/providers/openai-tiers.md b/structure/providers/openai-tiers.md index ec64480c51b..3c4ac9b54bc 100644 --- a/structure/providers/openai-tiers.md +++ b/structure/providers/openai-tiers.md @@ -438,6 +438,9 @@ cannot publish main usage. Pool health/failover handling stays scoped to Pool co The dispatch additionally fences the process-wide credential mutation epoch, so native main refresh and same-account reauth commits reject an older response before quota observation catches up. Publications for other credentials also conservatively drop the main update. +Every plain-main WS observer invocation claims its dispatch before checking liveness, preventing +HTTP publication from replacement Responses. A failed-upgrade HTTP fallback with no WS quota +frames remains unclaimed and publishes normally; response markers remain an additional guard. `src/codex/auth-api/main-account-probe.ts` re-reads the bounded stored main credential and rechecks its writer, bearer and generation after body/retry awaits, before publishing main usage, credits, plan, reauth or Reserve state, including terminal 401/403 mutations. An unreadable file diff --git a/structure/transports/responses.md b/structure/transports/responses.md index 3eb0bdacea9..b923b473e93 100644 --- a/structure/transports/responses.md +++ b/structure/transports/responses.md @@ -23,7 +23,7 @@ Canonical forward auth retains its separate fixed credential/metadata allowlist; Retired Codex Spark has no model-specific tool or Responses Lite override; general Lite handling and namespace scrubbing remain shared compatibility behavior. Codex quota/reset evidence follows the -[shared/Reserve policy](../providers/openai-tiers.md#public-provider-contract), including suppression of retired model-derived evidence before shared recovery. Plain-main HTTP/WS quota headers refresh `__main__` only on the canonical OpenAI forward provider when the materialized bearer and workspace match the owned main observation, using the hard-lock credential-match rule. `src/server/responses/passthrough-delivery.ts` rechecks the captured credential generation after its awaited import; `src/server/responses/core-codex-account.ts` retains one dispatch proof per WS observer and rechecks it for each frame. Rotated credentials and unmatched workspaces publish nothing; dispatch proofs also reject any OpenCodex-owned credential publication epoch change, including native main refresh or same-account reauth before quota re-observation; WS upstream responses and separately marked pre-response prelude projections (4xx refusals and 502/504 gateway failures) skip plain-main HTTP quota writes; prelude headers remain available to Pool replay, real HTTP fallbacks still publish, and Pool health/failover gates remain unchanged. +[shared/Reserve policy](../providers/openai-tiers.md#public-provider-contract), including suppression of retired model-derived evidence before shared recovery. Plain-main HTTP/WS quota headers refresh `__main__` only on the canonical OpenAI forward provider when the materialized bearer and workspace match the owned main observation, using the hard-lock credential-match rule. `src/server/responses/passthrough-delivery.ts` rechecks the captured credential generation after its awaited import; `src/server/responses/core-codex-account.ts` retains one dispatch proof per WS observer and rechecks it for each frame. Rotated credentials and unmatched workspaces publish nothing; dispatch proofs also reject any OpenCodex-owned credential publication epoch change, including native main refresh or same-account reauth before quota re-observation; each plain-main WS observer invocation claims its captured dispatch before liveness checks, and this claim is the authoritative HTTP publication guard across Response replacement; WS upstream responses and separately marked pre-response prelude projections (4xx refusals and 502/504 gateway failures) skip plain-main HTTP quota writes; prelude headers remain available to Pool replay, real HTTP fallbacks still publish, and Pool health/failover gates remain unchanged. ### Credential-bearing HTTP redirects diff --git a/tests/responses/responses-main-quota-observation.test.ts b/tests/responses/responses-main-quota-observation.test.ts index 47871c0af84..fcf0bb8db56 100644 --- a/tests/responses/responses-main-quota-observation.test.ts +++ b/tests/responses/responses-main-quota-observation.test.ts @@ -19,6 +19,8 @@ import { createTranslatorBudget } from "../../src/lib/translator-budget"; import { captureConfigGeneration } from "../../src/lib/state-store-sweeper"; import { parseRequest } from "../../src/responses/parser"; import { codexAccountSelectionForTurn, tryAdmitTurn } from "../../src/server/lifecycle"; +import { handleResponses } from "../../src/server/responses"; +import { BOUNDED_WS_RUNTIME } from "../helpers/ws-upstream-fixtures"; import { deliverPassthroughResponse } from "../../src/server/responses/passthrough-delivery"; import { codexWsQuotaObserver, retryCodexPoolOnAlternateAccount } from "../../src/server/responses/core-codex-account"; import { CodexWsMetadata } from "../../src/server/responses/codex-ws-metadata"; @@ -427,12 +429,14 @@ describe("credential-bound plain-main Responses quota", () => { sseFallback: (async () => { fallbackCalls++; return quotaResponse(quotaHeaders("37")); }) as typeof fetch, onQuota: codexWsQuotaObserver(ctx, provider, "gpt-5.5"), bunVersion: "1.4.0" }); expect(fallbackCalls).toBe(1); + expect(mainCache.isMainQuotaDispatchWsClaimed(ctx.mainQuotaDispatch!)).toBe(false); expect(isCodexWsUpstreamResponse(response)).toBe(false); expect(isCodexWsPreludeProjection(response)).toBe(false); expect(isCodexWsQuotaObservedResponse(response)).toBe(false); await deliver(ctx, response); expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(37); expect(getMainPolicyQuota()?.weeklyPercent).toBe(37); + expect(mainCache.isMainQuotaDispatchWsClaimed(ctx.mainQuotaDispatch!)).toBe(false); } finally { session.dispose(); globalThis.WebSocket = originalSocket; @@ -440,6 +444,114 @@ describe("credential-bound plain-main Responses quota", () => { }); } + test("full handler preflight for encrypted function output retains newer WS quota", async () => { + observe(); + const newerObserver = codexWsQuotaObserver(materialized(), provider, "gpt-5.5"); + expect(newerObserver).toBeDefined(); + const originalSocket = globalThis.WebSocket; + const sockets: PreflightSocket[] = []; + const opaqueBytes = Buffer.alloc(73, 1); + opaqueBytes[0] = 0x80; + const opaqueOutput = opaqueBytes.toString("base64").replace(/\+/g, "-").replace(/\//g, "_"); + class PreflightSocket extends EventTarget { + readyState = 0; + sent: string[] = []; + constructor() { + super(); + sockets.push(this); + queueMicrotask(() => { this.readyState = 1; this.dispatchEvent(new Event("open")); }); + } + send(text: string): void { + this.sent.push(text); + queueMicrotask(() => { + const emit = (payload: unknown) => this.dispatchEvent(new MessageEvent("message", { data: JSON.stringify(payload) })); + emit({ type: "codex.rate_limits", rate_limits: { primary: { used_percent: 17, window_minutes: 10080 } } }); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(17); + emit({ type: "response.created", response: { id: "fixture-preflight", status: "in_progress", output: [] } }); + // The committed response has a 17% prelude; another dispatch publishes 99%. + newerObserver!(quotaHeaders("99")); + expect(getMainPolicyQuota()?.weeklyPercent).toBe(99); + emit({ type: "response.output_text.delta", delta: "fixture answer", output_index: 0, content_index: 0 }); + emit({ type: "response.completed", response: { id: "fixture-preflight", status: "completed", output: [] } }); + }); + } + close(): void { + if (this.readyState === 3) return; + this.readyState = 3; + this.dispatchEvent(new Event("close")); + } + } + globalThis.WebSocket = PreflightSocket as unknown as typeof WebSocket; + try { + const response = await handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", headers: new Headers({ ...Object.fromEntries(caller()), "content-type": "application/json" }), + body: JSON.stringify({ model: "gpt-5.5", stream: true, input: [ + { type: "function_call", call_id: "fixture-call", name: "fixture_tool", arguments: "{}" }, + { type: "function_call_output", call_id: "fixture-call", + output: [{ type: "encrypted_content", encrypted_content: opaqueOutput }] }, + ] }), + }), { ...config, defaultProvider: "openai", streamMode: "legacy-tee", + providers: { openai: { ...provider, codexAccountMode: "direct" } } }, { model: "", provider: "" }, + { codexWsRuntimeIdentity: BOUNDED_WS_RUNTIME }); + expect(response.status).toBe(200); + expect(sockets).toHaveLength(1); + const sentOutput = JSON.parse(sockets[0]!.sent[0]!).input.find((item: { type?: string }) => item.type === "function_call_output"); + expect(sentOutput.output[0].type).toBe("encrypted_content"); + // The preflight replay is a fresh Response; no response-object marker survives. + expect(isCodexWsUpstreamResponse(response)).toBe(false); + expect(isCodexWsPreludeProjection(response)).toBe(false); + const text = await response.text(); + expect(text).toContain("fixture answer"); + expect(text).toContain("response.completed"); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(99); + expect(getMainPolicyQuota()?.weeklyPercent).toBe(99); + expect(getMainAccountHardLockStatus({ codexMainAccountHardLock: true }).state).toBe("blocked"); + } finally { + for (const socket of sockets) socket.close(); + globalThis.WebSocket = originalSocket; + } + }); + + test("an observed WS dispatch rejects an unmarked HTTP-shaped quota snapshot", async () => { + observe(); + const ctx = materialized(); + const observer = codexWsQuotaObserver(ctx, provider, "gpt-5.5"); + expect(observer).toBeDefined(); + observer!(quotaHeaders("17")); + observer!(quotaHeaders("99")); + expect(mainCache.isMainQuotaDispatchWsClaimed(ctx.mainQuotaDispatch!)).toBe(true); + const response = quotaResponse(quotaHeaders("17")); + expect(isCodexWsUpstreamResponse(response)).toBe(false); + expect(isCodexWsPreludeProjection(response)).toBe(false); + await deliver(ctx, response); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(99); + expect(getMainPolicyQuota()?.weeklyPercent).toBe(99); + expect(getMainAccountHardLockStatus({ codexMainAccountHardLock: true }).state).toBe("blocked"); + }); + + test("every observer invocation claims its captured dispatch before checking liveness", async () => { + observe(); + const ctx = materialized(); + const dispatch = ctx.mainQuotaDispatch!; + const observer = codexWsQuotaObserver(ctx, provider, "gpt-5.5"); + expect(observer).toBeDefined(); + expect(mainCache.isMainQuotaDispatchWsClaimed(dispatch)).toBe(false); + mainCache.observeMainQuotaCredential("fixture-replaced-before-frame", ACCOUNT); + expect(mainCache.isMainQuotaDispatchLive(dispatch)).toBe(false); + observer!(new Headers()); + expect(mainCache.isMainQuotaDispatchWsClaimed(dispatch)).toBe(true); + expect(getAccountQuota(MAIN)).toBeNull(); + // Reusing the context cannot transfer an old observer's claim to a new dispatch. + observe(); + materializeCodexUpstreamAuth(caller(), ctx, { config, modelId: "gpt-5.5" }); + expect(ctx.mainQuotaDispatch).not.toBe(dispatch); + expect(mainCache.isMainQuotaDispatchWsClaimed(ctx.mainQuotaDispatch!)).toBe(false); + observer!(quotaHeaders("99")); + await deliver(ctx, quotaResponse(quotaHeaders("37"))); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(37); + expect(getMainPolicyQuota()?.weeklyPercent).toBe(37); + }); + test("async Reserve admission failure clears a reused main dispatch proof", async () => { observe(); const ctx = materialized(); From 5d96e500ec6a1292e7459cc384c457ddb0bb8b24 Mon Sep 17 00:00:00 2001 From: JUN Date: Fri, 9 Oct 2026 21:30:56 +0900 Subject: [PATCH 5/7] fix(quota): give an HTTP replacement send its own main quota dispatch Review follow-up for #6800: the ambiguous-reset HTTP replacement reused the dispatch a failed WebSocket attempt had already claimed, so its fresh quota could not publish. The replacement now renews the dispatch into a new, unclaimed object that copies every captured fence unchanged; the failed attempt keeps the claimed one. --- src/codex/main-account-cache.ts | 5 + src/server/responses/passthrough-dispatch.ts | 4 + structure/providers/openai-accounts.md | 2 +- structure/providers/openai-tiers.md | 5 +- structure/transports/responses-failover.md | 2 +- structure/transports/responses.md | 2 +- .../responses-main-quota-observation.test.ts | 92 ++++++++++++++++++- 7 files changed, 107 insertions(+), 5 deletions(-) diff --git a/src/codex/main-account-cache.ts b/src/codex/main-account-cache.ts index 41e7a9c5d89..3a46180e672 100644 --- a/src/codex/main-account-cache.ts +++ b/src/codex/main-account-cache.ts @@ -94,6 +94,11 @@ export function isMainQuotaDispatchWsClaimed(dispatch: MainQuotaDispatch): boole return wsObservedMainDispatches.has(dispatch); } +/** Give a replacement physical attempt its own quota ownership without recapturing credential fences. */ +export function renewMainQuotaDispatchForAttempt(dispatch: MainQuotaDispatch): MainQuotaDispatch { + return { ...dispatch }; +} + export function captureMainQuotaDispatch( accessToken: string, accountId: string | undefined, configGeneration: number, ): MainQuotaDispatch | undefined { diff --git a/src/server/responses/passthrough-dispatch.ts b/src/server/responses/passthrough-dispatch.ts index a2b423b0a0e..bdca37aed80 100644 --- a/src/server/responses/passthrough-dispatch.ts +++ b/src/server/responses/passthrough-dispatch.ts @@ -12,6 +12,7 @@ import type { ResponsesTransport } from "./request-transport"; import type { ResponsesEffects } from "./response-effects"; import type { ResponsesSendBudget } from "./request-send-budget"; import { transientSendCapFor } from "./request-send-budget"; +import { renewMainQuotaDispatchForAttempt } from "../../codex/main-account-cache"; import { isCanonicalOpenAiForwardProvider } from "../../providers/openai-tiers"; import { isLocalUpstream } from "../../lib/local-upstream"; import { codexSafetyBufferingFilterOptions, terminalStatusFromParsed } from "../relay"; @@ -802,6 +803,9 @@ export async function preparePassthroughExchange( const sendAmbiguousReplacement = ( signal: AbortSignal = upstream.signal, ): Promise => { + if (admissionState.authCtx.kind === "main" && admissionState.authCtx.mainQuotaDispatch) { + admissionState.authCtx.mainQuotaDispatch = renewMainQuotaDispatchForAttempt(admissionState.authCtx.mainQuotaDispatch); + } const report = transientSendReporter(); let started = false; const run = () => fetchWithHeaderTimeout( diff --git a/structure/providers/openai-accounts.md b/structure/providers/openai-accounts.md index 7bd72b05ca0..5e716baef5f 100644 --- a/structure/providers/openai-accounts.md +++ b/structure/providers/openai-accounts.md @@ -211,7 +211,7 @@ request-scoped: a translated Claude turn resolves through Pool selection like an main keeps its health, quarantine and refresh-and-classify handling. Only a bearer the client itself supplied is caller-owned and exempt from stored state. Both synchronous and asynchronous stored-main substitution in `src/codex/auth-context.ts` remove a caller account header before copying the stored identity; an absent stored account ID leaves no account header. Caller-owned native Direct authentication retains its existing passthrough behavior. -Plain-main HTTP and WebSocket Responses refresh `__main__` only on the canonical OpenAI forward provider when the sent bearer and effective workspace match the main credential already observed under native ownership, the same equality rule used by the hard lock. `src/codex/auth-context.ts` captures a process-local dispatch proof after materialization; `src/server/responses/passthrough-delivery.ts` and `src/server/responses/core-codex-account.ts` recheck its credential/identity generations before publishing. The dispatch also captures the process-wide credential mutation epoch; any OpenCodex-owned credential publication, including native main refresh or same-account reauth before a new quota credential observation, rejects an older dispatch. Publications for other credentials conservatively drop the main update as well. Same-account token rotation and A→B→A changes reject old responses; WebSocket observers retain the original dispatch proof across frames and claim it on every invocation before checking liveness. This process-local claim prevents plain-main HTTP publication even when a downstream stream wrapper replaces the Response; response markers remain an additional guard, including separately marked pre-response prelude projections (4xx refusals and 502/504 gateway failures). Prelude headers remain available to Pool replay; real HTTP fallback responses still publish through HTTP delivery. Caller-owned requests acquire no physical-main read or Pool health state. +Plain-main HTTP and WebSocket Responses refresh `__main__` only on the canonical OpenAI forward provider when the sent bearer and effective workspace match the main credential already observed under native ownership, the same equality rule used by the hard lock. `src/codex/auth-context.ts` captures a process-local dispatch proof after materialization; `src/server/responses/passthrough-delivery.ts` and `src/server/responses/core-codex-account.ts` recheck its credential/identity generations before publishing. The dispatch also captures the process-wide credential mutation epoch; any OpenCodex-owned credential publication, including native main refresh or same-account reauth before a new quota credential observation, rejects an older dispatch. Publications for other credentials conservatively drop the main update as well. Same-account token rotation and A→B→A changes reject old responses; WebSocket observers retain the original dispatch proof across frames and claim it on every invocation before checking liveness. This process-local claim belongs to one physical attempt and prevents plain-main HTTP publication even when a downstream stream wrapper replaces the Response; response markers remain an additional guard, including separately marked pre-response prelude projections (4xx refusals and 502/504 gateway failures). Prelude headers remain available to Pool replay; real HTTP fallback responses still publish through HTTP delivery. An operator-granted HTTP replacement gets a new unclaimed dispatch object with every captured credential and config fence copied unchanged; the failed WS observer retains the old object. Caller-owned requests acquire no physical-main read or Pool health state. `src/providers/openai-sidecar.ts` releases quota-probe ownership on every materialization or usability failure before transferring a resolved context to its caller. Audio reports one terminal upstream outcome after validating the response body; redirects remain diff --git a/structure/providers/openai-tiers.md b/structure/providers/openai-tiers.md index 3c4ac9b54bc..2a44c52135d 100644 --- a/structure/providers/openai-tiers.md +++ b/structure/providers/openai-tiers.md @@ -439,8 +439,11 @@ The dispatch additionally fences the process-wide credential mutation epoch, so refresh and same-account reauth commits reject an older response before quota observation catches up. Publications for other credentials also conservatively drop the main update. Every plain-main WS observer invocation claims its dispatch before checking liveness, preventing -HTTP publication from replacement Responses. A failed-upgrade HTTP fallback with no WS quota +HTTP publication from stream-wrapper replacement Responses. A failed-upgrade HTTP fallback with no WS quota frames remains unclaimed and publishes normally; response markers remain an additional guard. +An operator-granted HTTP replacement renews only the dispatch object identity, copying all captured +credential and config fences unchanged. Its unclaimed attempt can publish only while those original +fences remain live; the failed WS observer retains its old claimed object. `src/codex/auth-api/main-account-probe.ts` re-reads the bounded stored main credential and rechecks its writer, bearer and generation after body/retry awaits, before publishing main usage, credits, plan, reauth or Reserve state, including terminal 401/403 mutations. An unreadable file diff --git a/structure/transports/responses-failover.md b/structure/transports/responses-failover.md index cc26a8ed31d..5967a7e6bc1 100644 --- a/structure/transports/responses-failover.md +++ b/structure/transports/responses-failover.md @@ -409,7 +409,7 @@ the public server reference already documents. The 504 and a drop after the resp are never replaced. Only the 502 of a socket that closed or errored before any Responses event may be replaced over HTTP, when the provider opted into `retryOnReset` (#4191). That replacement claims from the request's one allowance; if it resets before its head, that is the pre-header row -again and may use a configured second replacement, otherwise it settles as the refusal. +again and may use a configured second replacement, otherwise it settles as the refusal. For plain-main quota, each physical HTTP replacement renews the dispatch object identity while copying every original credential and config fence unchanged. The failed WS observer retains its claimed object, so only the live replacement attempt may publish fresh HTTP headers. This reclassification is the recorded behaviour change: before it, the pre-header refusal borrowed `upstream_closed_before_response` and its 502, which multiplied the duplicate send diff --git a/structure/transports/responses.md b/structure/transports/responses.md index b923b473e93..81ea11f63b1 100644 --- a/structure/transports/responses.md +++ b/structure/transports/responses.md @@ -23,7 +23,7 @@ Canonical forward auth retains its separate fixed credential/metadata allowlist; Retired Codex Spark has no model-specific tool or Responses Lite override; general Lite handling and namespace scrubbing remain shared compatibility behavior. Codex quota/reset evidence follows the -[shared/Reserve policy](../providers/openai-tiers.md#public-provider-contract), including suppression of retired model-derived evidence before shared recovery. Plain-main HTTP/WS quota headers refresh `__main__` only on the canonical OpenAI forward provider when the materialized bearer and workspace match the owned main observation, using the hard-lock credential-match rule. `src/server/responses/passthrough-delivery.ts` rechecks the captured credential generation after its awaited import; `src/server/responses/core-codex-account.ts` retains one dispatch proof per WS observer and rechecks it for each frame. Rotated credentials and unmatched workspaces publish nothing; dispatch proofs also reject any OpenCodex-owned credential publication epoch change, including native main refresh or same-account reauth before quota re-observation; each plain-main WS observer invocation claims its captured dispatch before liveness checks, and this claim is the authoritative HTTP publication guard across Response replacement; WS upstream responses and separately marked pre-response prelude projections (4xx refusals and 502/504 gateway failures) skip plain-main HTTP quota writes; prelude headers remain available to Pool replay, real HTTP fallbacks still publish, and Pool health/failover gates remain unchanged. +[shared/Reserve policy](../providers/openai-tiers.md#public-provider-contract), including suppression of retired model-derived evidence before shared recovery. Plain-main HTTP/WS quota headers refresh `__main__` only on the canonical OpenAI forward provider when the materialized bearer and workspace match the owned main observation, using the hard-lock credential-match rule. `src/server/responses/passthrough-delivery.ts` rechecks the captured credential generation after its awaited import; `src/server/responses/core-codex-account.ts` retains one dispatch proof per WS observer and rechecks it for each frame. Rotated credentials and unmatched workspaces publish nothing; dispatch proofs also reject any OpenCodex-owned credential publication epoch change, including native main refresh or same-account reauth before quota re-observation; each plain-main WS observer invocation claims its captured dispatch before liveness checks, and this claim is the authoritative HTTP publication guard across Response replacement; WS upstream responses and separately marked pre-response prelude projections (4xx refusals and 502/504 gateway failures) skip plain-main HTTP quota writes; prelude headers remain available to Pool replay, real HTTP fallbacks still publish; each operator-granted HTTP replacement renews the dispatch object with its original credential and config fences unchanged, and Pool health/failover gates remain unchanged. ### Credential-bearing HTTP redirects diff --git a/tests/responses/responses-main-quota-observation.test.ts b/tests/responses/responses-main-quota-observation.test.ts index fcf0bb8db56..7e75c16da80 100644 --- a/tests/responses/responses-main-quota-observation.test.ts +++ b/tests/responses/responses-main-quota-observation.test.ts @@ -8,7 +8,7 @@ import { type CodexAuthContext, } from "../../src/codex/auth-context"; import { beginNativeMainReauth, forceRefreshMainAccountToken } from "../../src/codex/main-account"; -import { codexCredentialMutationEpoch } from "../../src/codex/credential-mutation-epoch"; +import { advanceCodexCredentialMutationEpoch, codexCredentialMutationEpoch } from "../../src/codex/credential-mutation-epoch"; import { resetMainCodexAccountIdentityTrackingForTests } from "../../src/codex/account-lifecycle"; import { clearAccountNeedsReauth } from "../../src/codex/account-runtime-state"; import { saveCodexAccountCredential } from "../../src/codex/account-store"; @@ -444,6 +444,96 @@ describe("credential-bound plain-main Responses quota", () => { }); } + test("full handler HTTP replacement after WS quota publishes under a new physical attempt", async () => { + observe(); + const originalSocket = globalThis.WebSocket; + const sockets: FailedAttemptSocket[] = []; + class FailedAttemptSocket extends EventTarget { + readyState = 0; + sent: string[] = []; + constructor() { + super(); + sockets.push(this); + queueMicrotask(() => { this.readyState = 1; this.dispatchEvent(new Event("open")); }); + } + send(text: string): void { + this.sent.push(text); + queueMicrotask(() => { + this.dispatchEvent(new MessageEvent("message", { data: JSON.stringify({ + type: "codex.rate_limits", rate_limits: { primary: { used_percent: 17, window_minutes: 10080 } }, + }) })); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(17); + expect(getMainPolicyQuota()?.weeklyPercent).toBe(17); + this.close(); + }); + } + close(): void { + if (this.readyState === 3) return; + this.readyState = 3; + this.dispatchEvent(new Event("close")); + } + } + globalThis.WebSocket = FailedAttemptSocket as unknown as typeof WebSocket; + let httpCalls = 0; + globalThis.fetch = (async (_input: RequestInfo | URL, init?: RequestInit) => { + httpCalls++; + expect(JSON.parse(String(init?.body))).toMatchObject({ model: "gpt-5.5", store: false }); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(17); + const headers = quotaHeaders("99"); + headers.set("content-type", "text/event-stream"); + return new Response(`event: response.completed\ndata: ${JSON.stringify({ + type: "response.completed", response: { id: "fixture-http-replacement", status: "completed", output: [] }, + })}\n\n`, { status: 200, headers }); + }) as typeof fetch; + try { + const response = await handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", headers: new Headers({ ...Object.fromEntries(caller()), "content-type": "application/json" }), + body: JSON.stringify({ model: "gpt-5.5", stream: true, store: false, input: "fixture self-contained turn" }), + }), { ...config, defaultProvider: "openai", + providers: { openai: { ...provider, codexAccountMode: "direct", retryOnReset: {} } } }, + { model: "", provider: "" }, { codexWsRuntimeIdentity: BOUNDED_WS_RUNTIME }); + expect(response.status).toBe(200); + expect(await response.text()).toContain("fixture-http-replacement"); + expect(sockets).toHaveLength(1); + expect(sockets[0]!.sent).toHaveLength(1); + expect(sockets[0]!.readyState).toBe(3); + expect(httpCalls).toBe(1); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(99); + expect(getMainPolicyQuota()?.weeklyPercent).toBe(99); + expect(getMainAccountHardLockStatus({ codexMainAccountHardLock: true }).state).toBe("blocked"); + } finally { + for (const socket of sockets) socket.close(); + globalThis.WebSocket = originalSocket; + } + }); + + for (const mutation of ["credential observation", "credential publication epoch"] as const) { + test(`attempt renewal preserves retired fences after ${mutation}`, async () => { + observe(); + const ctx = materialized(); + const original = ctx.mainQuotaDispatch!; + const observer = codexWsQuotaObserver(ctx, provider, "gpt-5.5"); + expect(observer).toBeDefined(); + observer!(quotaHeaders("17")); + expect(mainCache.isMainQuotaDispatchWsClaimed(original)).toBe(true); + if (mutation === "credential observation") mainCache.observeMainQuotaCredential("fixture-next-credential", ACCOUNT); + else advanceCodexCredentialMutationEpoch(); + const renewed = mainCache.renewMainQuotaDispatchForAttempt(original); + expect(renewed).not.toBe(original); + expect(renewed).toEqual(original); + expect(renewed.writer).toBe(original.writer); + expect(mainCache.isMainQuotaDispatchWsClaimed(renewed)).toBe(false); + expect(mainCache.isMainQuotaDispatchWsClaimed(original)).toBe(true); + expect(mainCache.isMainQuotaDispatchLive(renewed)).toBe(false); + ctx.mainQuotaDispatch = renewed; + observer!(quotaHeaders("99")); + expect(mainCache.isMainQuotaDispatchWsClaimed(renewed)).toBe(false); + await deliver(ctx, quotaResponse(quotaHeaders("99"))); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(17); + expect(getMainPolicyQuota()?.weeklyPercent).toBe(17); + }); + } + test("full handler preflight for encrypted function output retains newer WS quota", async () => { observe(); const newerObserver = codexWsQuotaObserver(materialized(), provider, "gpt-5.5"); From 44661352e63c1431a6e47e9061789ac5e92d61bc Mon Sep 17 00:00:00 2001 From: JUN Date: Fri, 9 Oct 2026 21:45:33 +0900 Subject: [PATCH 6/7] fix(quota): scope main quota dispatch ownership to each WS observer Review follow-up for #6800: a recovery send after a failed WebSocket attempt reused the dispatch that attempt had claimed, suppressing fresh HTTP evidence. Each plain-main WebSocket observer now renews the context's dispatch into its own copy and closes over it, so ownership belongs to one physical attempt while every captured credential fence stays unchanged. --- src/server/responses/core-codex-account.ts | 7 +- structure/providers/openai-accounts.md | 2 +- structure/providers/openai-tiers.md | 6 +- structure/transports/responses-failover.md | 2 +- structure/transports/responses.md | 2 +- .../responses-main-quota-observation.test.ts | 113 +++++++++++++++++- 6 files changed, 122 insertions(+), 10 deletions(-) diff --git a/src/server/responses/core-codex-account.ts b/src/server/responses/core-codex-account.ts index a138b13209d..8fa6dafb1ed 100644 --- a/src/server/responses/core-codex-account.ts +++ b/src/server/responses/core-codex-account.ts @@ -6,7 +6,7 @@ import { computeQuotaCooldown, formatCodexProviderForLog, } from "../../codex/routing"; -import { claimMainQuotaDispatchForWs, isMainQuotaDispatchLive, type MainQuotaDispatch } from "../../codex/main-account-cache"; +import { claimMainQuotaDispatchForWs, isMainQuotaDispatchLive, renewMainQuotaDispatchForAttempt, type MainQuotaDispatch } from "../../codex/main-account-cache"; import type { CodexWsQuotaObserver } from "./codex-ws-metadata"; import { isCanonicalOpenAiForwardProvider } from "../../providers/openai-tiers"; import { isCodexAccountGenerationLive } from "../../codex/account-store"; @@ -135,8 +135,9 @@ export function liveMainQuotaDispatch( export function codexWsQuotaObserver(authCtx: CodexAuthContext, provider: OcxProviderConfig, modelId?: string): CodexWsQuotaObserver | undefined { if (!isCanonicalOpenAiForwardProvider(provider)) return undefined; if (!usesCodexForwardPoolAuth(authCtx, provider)) { - const dispatch = liveMainQuotaDispatch(authCtx, provider); - if (!dispatch) return undefined; + const captured = liveMainQuotaDispatch(authCtx, provider); + if (!captured || authCtx.kind !== "main") return undefined; + const dispatch = authCtx.mainQuotaDispatch = renewMainQuotaDispatchForAttempt(captured); return headers => { claimMainQuotaDispatchForWs(dispatch); if (!isMainQuotaDispatchLive(dispatch)) return; diff --git a/structure/providers/openai-accounts.md b/structure/providers/openai-accounts.md index 5e716baef5f..f9737bfd144 100644 --- a/structure/providers/openai-accounts.md +++ b/structure/providers/openai-accounts.md @@ -211,7 +211,7 @@ request-scoped: a translated Claude turn resolves through Pool selection like an main keeps its health, quarantine and refresh-and-classify handling. Only a bearer the client itself supplied is caller-owned and exempt from stored state. Both synchronous and asynchronous stored-main substitution in `src/codex/auth-context.ts` remove a caller account header before copying the stored identity; an absent stored account ID leaves no account header. Caller-owned native Direct authentication retains its existing passthrough behavior. -Plain-main HTTP and WebSocket Responses refresh `__main__` only on the canonical OpenAI forward provider when the sent bearer and effective workspace match the main credential already observed under native ownership, the same equality rule used by the hard lock. `src/codex/auth-context.ts` captures a process-local dispatch proof after materialization; `src/server/responses/passthrough-delivery.ts` and `src/server/responses/core-codex-account.ts` recheck its credential/identity generations before publishing. The dispatch also captures the process-wide credential mutation epoch; any OpenCodex-owned credential publication, including native main refresh or same-account reauth before a new quota credential observation, rejects an older dispatch. Publications for other credentials conservatively drop the main update as well. Same-account token rotation and A→B→A changes reject old responses; WebSocket observers retain the original dispatch proof across frames and claim it on every invocation before checking liveness. This process-local claim belongs to one physical attempt and prevents plain-main HTTP publication even when a downstream stream wrapper replaces the Response; response markers remain an additional guard, including separately marked pre-response prelude projections (4xx refusals and 502/504 gateway failures). Prelude headers remain available to Pool replay; real HTTP fallback responses still publish through HTTP delivery. An operator-granted HTTP replacement gets a new unclaimed dispatch object with every captured credential and config fence copied unchanged; the failed WS observer retains the old object. Caller-owned requests acquire no physical-main read or Pool health state. +Plain-main HTTP and WebSocket Responses refresh `__main__` only on the canonical OpenAI forward provider when the sent bearer and effective workspace match the main credential already observed under native ownership, the same equality rule used by the hard lock. `src/codex/auth-context.ts` captures a process-local dispatch proof after materialization; `src/server/responses/passthrough-delivery.ts` and `src/server/responses/core-codex-account.ts` recheck its credential/identity generations before publishing. The dispatch also captures the process-wide credential mutation epoch; any OpenCodex-owned credential publication, including native main refresh or same-account reauth before a new quota credential observation, rejects an older dispatch. Publications for other credentials conservatively drop the main update as well. Same-account token rotation and A→B→A changes reject old responses; Each WebSocket observer renews the live dispatch object with every captured fence unchanged, retains that copy across frames, and claims it on every invocation before checking liveness. A later observer starts unclaimed, so its failed-upgrade HTTP fallback can publish even if the prior WS attempt observed quota. This process-local claim belongs to one physical attempt and prevents plain-main HTTP publication even when a downstream stream wrapper replaces the Response; response markers remain an additional guard, including separately marked pre-response prelude projections (4xx refusals and 502/504 gateway failures). Prelude headers remain available to Pool replay; real HTTP fallback responses still publish through HTTP delivery. An operator-granted HTTP replacement gets a new unclaimed dispatch object with every captured credential and config fence copied unchanged; the failed WS observer retains the old object. Caller-owned requests acquire no physical-main read or Pool health state. `src/providers/openai-sidecar.ts` releases quota-probe ownership on every materialization or usability failure before transferring a resolved context to its caller. Audio reports one terminal upstream outcome after validating the response body; redirects remain diff --git a/structure/providers/openai-tiers.md b/structure/providers/openai-tiers.md index 2a44c52135d..d313fb49af2 100644 --- a/structure/providers/openai-tiers.md +++ b/structure/providers/openai-tiers.md @@ -438,9 +438,11 @@ cannot publish main usage. Pool health/failover handling stays scoped to Pool co The dispatch additionally fences the process-wide credential mutation epoch, so native main refresh and same-account reauth commits reject an older response before quota observation catches up. Publications for other credentials also conservatively drop the main update. -Every plain-main WS observer invocation claims its dispatch before checking liveness, preventing +Each plain-main WS observer renews its live dispatch object with every captured fence unchanged. +Every invocation claims that observer's own copy before checking liveness, preventing HTTP publication from stream-wrapper replacement Responses. A failed-upgrade HTTP fallback with no WS quota -frames remains unclaimed and publishes normally; response markers remain an additional guard. +frames remains unclaimed and publishes normally, including after an earlier attempt observed quota; +response markers remain an additional guard. An operator-granted HTTP replacement renews only the dispatch object identity, copying all captured credential and config fences unchanged. Its unclaimed attempt can publish only while those original fences remain live; the failed WS observer retains its old claimed object. diff --git a/structure/transports/responses-failover.md b/structure/transports/responses-failover.md index 5967a7e6bc1..355b61bcc36 100644 --- a/structure/transports/responses-failover.md +++ b/structure/transports/responses-failover.md @@ -409,7 +409,7 @@ the public server reference already documents. The 504 and a drop after the resp are never replaced. Only the 502 of a socket that closed or errored before any Responses event may be replaced over HTTP, when the provider opted into `retryOnReset` (#4191). That replacement claims from the request's one allowance; if it resets before its head, that is the pre-header row -again and may use a configured second replacement, otherwise it settles as the refusal. For plain-main quota, each physical HTTP replacement renews the dispatch object identity while copying every original credential and config fence unchanged. The failed WS observer retains its claimed object, so only the live replacement attempt may publish fresh HTTP headers. +again and may use a configured second replacement, otherwise it settles as the refusal. For plain-main quota, each physical HTTP replacement renews the dispatch object identity while copying every original credential and config fence unchanged. The failed WS observer retains its claimed object, so only the live replacement attempt may publish fresh HTTP headers. Each rebuilt WS observer also renews its own live dispatch copy; an upgrade failure with no quota frames leaves that attempt unclaimed for HTTP fallback publication. This reclassification is the recorded behaviour change: before it, the pre-header refusal borrowed `upstream_closed_before_response` and its 502, which multiplied the duplicate send diff --git a/structure/transports/responses.md b/structure/transports/responses.md index 81ea11f63b1..db12d050427 100644 --- a/structure/transports/responses.md +++ b/structure/transports/responses.md @@ -23,7 +23,7 @@ Canonical forward auth retains its separate fixed credential/metadata allowlist; Retired Codex Spark has no model-specific tool or Responses Lite override; general Lite handling and namespace scrubbing remain shared compatibility behavior. Codex quota/reset evidence follows the -[shared/Reserve policy](../providers/openai-tiers.md#public-provider-contract), including suppression of retired model-derived evidence before shared recovery. Plain-main HTTP/WS quota headers refresh `__main__` only on the canonical OpenAI forward provider when the materialized bearer and workspace match the owned main observation, using the hard-lock credential-match rule. `src/server/responses/passthrough-delivery.ts` rechecks the captured credential generation after its awaited import; `src/server/responses/core-codex-account.ts` retains one dispatch proof per WS observer and rechecks it for each frame. Rotated credentials and unmatched workspaces publish nothing; dispatch proofs also reject any OpenCodex-owned credential publication epoch change, including native main refresh or same-account reauth before quota re-observation; each plain-main WS observer invocation claims its captured dispatch before liveness checks, and this claim is the authoritative HTTP publication guard across Response replacement; WS upstream responses and separately marked pre-response prelude projections (4xx refusals and 502/504 gateway failures) skip plain-main HTTP quota writes; prelude headers remain available to Pool replay, real HTTP fallbacks still publish; each operator-granted HTTP replacement renews the dispatch object with its original credential and config fences unchanged, and Pool health/failover gates remain unchanged. +[shared/Reserve policy](../providers/openai-tiers.md#public-provider-contract), including suppression of retired model-derived evidence before shared recovery. Plain-main HTTP/WS quota headers refresh `__main__` only on the canonical OpenAI forward provider when the materialized bearer and workspace match the owned main observation, using the hard-lock credential-match rule. `src/server/responses/passthrough-delivery.ts` rechecks the captured credential generation after its awaited import; `src/server/responses/core-codex-account.ts` renews a live dispatch object for each WS observer without recapturing its credential or config fences, retains that copy across frames, and rechecks it for each frame. Rotated credentials and unmatched workspaces publish nothing; dispatch proofs also reject any OpenCodex-owned credential publication epoch change, including native main refresh or same-account reauth before quota re-observation; each plain-main WS observer invocation claims its captured dispatch before liveness checks, and this claim is the authoritative HTTP publication guard across Response replacement; WS upstream responses and separately marked pre-response prelude projections (4xx refusals and 502/504 gateway failures) skip plain-main HTTP quota writes; prelude headers remain available to Pool replay, real HTTP fallbacks still publish; each operator-granted HTTP replacement renews the dispatch object with its original credential and config fences unchanged, and Pool health/failover gates remain unchanged. ### Credential-bearing HTTP redirects diff --git a/tests/responses/responses-main-quota-observation.test.ts b/tests/responses/responses-main-quota-observation.test.ts index 7e75c16da80..0b828c16cc3 100644 --- a/tests/responses/responses-main-quota-observation.test.ts +++ b/tests/responses/responses-main-quota-observation.test.ts @@ -511,8 +511,8 @@ describe("credential-bound plain-main Responses quota", () => { test(`attempt renewal preserves retired fences after ${mutation}`, async () => { observe(); const ctx = materialized(); - const original = ctx.mainQuotaDispatch!; const observer = codexWsQuotaObserver(ctx, provider, "gpt-5.5"); + const original = ctx.mainQuotaDispatch!; expect(observer).toBeDefined(); observer!(quotaHeaders("17")); expect(mainCache.isMainQuotaDispatchWsClaimed(original)).toBe(true); @@ -534,6 +534,83 @@ describe("credential-bound plain-main Responses quota", () => { }); } + test("full handler sanitized recovery with failed WS upgrade publishes fresh HTTP quota", async () => { + observe(); + const originalSocket = globalThis.WebSocket; + const sockets: RecoverySocket[] = []; + const opaqueBytes = Buffer.alloc(73, 1); + opaqueBytes[0] = 0x80; + const opaqueOutput = opaqueBytes.toString("base64").replace(/\+/g, "-").replace(/\//g, "_"); + class RecoverySocket extends EventTarget { + readyState = 0; + sent: string[] = []; + constructor() { + super(); + sockets.push(this); + const attempt = sockets.length; + queueMicrotask(() => { + if (attempt === 1) { this.readyState = 1; this.dispatchEvent(new Event("open")); } + else this.close(); + }); + } + send(text: string): void { + this.sent.push(text); + queueMicrotask(() => { + const emit = (payload: unknown) => this.dispatchEvent(new MessageEvent("message", { data: JSON.stringify(payload) })); + emit({ type: "codex.rate_limits", rate_limits: { primary: { used_percent: 17, window_minutes: 10080 } } }); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(17); + expect(getMainPolicyQuota()?.weeklyPercent).toBe(17); + emit({ type: "error", status_code: 400, error: { type: "invalid_request_error", + code: "invalid_encrypted_content", message: "The encrypted content could not be verified." } }); + }); + } + close(): void { + if (this.readyState === 3) return; + this.readyState = 3; + this.dispatchEvent(new Event("close")); + } + } + globalThis.WebSocket = RecoverySocket as unknown as typeof WebSocket; + let httpCalls = 0; + globalThis.fetch = (async (_input: RequestInfo | URL, init?: RequestInit) => { + httpCalls++; + const body = JSON.parse(String(init?.body)); + const output = body.input.find((item: { type?: string }) => item.type === "function_call_output"); + expect(output.output).toEqual([{ type: "input_text", text: "[encrypted content omitted]" }]); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(17); + const headers = quotaHeaders("99"); + headers.set("content-type", "text/event-stream"); + return new Response(`event: response.completed\ndata: ${JSON.stringify({ + type: "response.completed", response: { id: "fixture-sanitized-http", status: "completed", output: [] }, + })}\n\n`, { status: 200, headers }); + }) as typeof fetch; + try { + const response = await handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", headers: new Headers({ ...Object.fromEntries(caller()), "content-type": "application/json" }), + body: JSON.stringify({ model: "gpt-5.5", stream: true, store: false, input: [ + { type: "function_call", call_id: "fixture-call", name: "fixture_tool", arguments: "{}" }, + { type: "function_call_output", call_id: "fixture-call", + output: [{ type: "encrypted_content", encrypted_content: opaqueOutput }] }, + ] }), + }), { ...config, defaultProvider: "openai", providers: { openai: { ...provider, codexAccountMode: "direct" } } }, + { model: "", provider: "" }, { codexWsRuntimeIdentity: BOUNDED_WS_RUNTIME }); + expect(response.status).toBe(200); + expect(await response.text()).toContain("fixture-sanitized-http"); + expect(sockets).toHaveLength(2); + expect(sockets[0]!.sent).toHaveLength(1); + const firstOutput = JSON.parse(sockets[0]!.sent[0]!).input.find((item: { type?: string }) => item.type === "function_call_output"); + expect(firstOutput.output[0].type).toBe("encrypted_content"); + expect(sockets[1]!.sent).toHaveLength(0); + expect(httpCalls).toBe(1); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(99); + expect(getMainPolicyQuota()?.weeklyPercent).toBe(99); + expect(getMainAccountHardLockStatus({ codexMainAccountHardLock: true }).state).toBe("blocked"); + } finally { + for (const socket of sockets) socket.close(); + globalThis.WebSocket = originalSocket; + } + }); + test("full handler preflight for encrypted function output retains newer WS quota", async () => { observe(); const newerObserver = codexWsQuotaObserver(materialized(), provider, "gpt-5.5"); @@ -619,11 +696,43 @@ describe("credential-bound plain-main Responses quota", () => { expect(getMainAccountHardLockStatus({ codexMainAccountHardLock: true }).state).toBe("blocked"); }); + test("each WS observer owns a fresh dispatch and an old observer cannot claim its successor", async () => { + observe(); + const ctx = materialized(); + const materializedDispatch = ctx.mainQuotaDispatch!; + const firstObserver = codexWsQuotaObserver(ctx, provider, "gpt-5.5"); + const firstDispatch = ctx.mainQuotaDispatch!; + expect(firstObserver).toBeDefined(); + expect(firstDispatch).not.toBe(materializedDispatch); + expect(firstDispatch).toEqual(materializedDispatch); + firstObserver!(quotaHeaders("17")); + expect(mainCache.isMainQuotaDispatchWsClaimed(firstDispatch)).toBe(true); + const nextObserver = codexWsQuotaObserver(ctx, provider, "gpt-5.5"); + const nextDispatch = ctx.mainQuotaDispatch!; + expect(nextObserver).toBeDefined(); + expect(nextDispatch).not.toBe(firstDispatch); + expect(nextDispatch).toEqual(firstDispatch); + expect(nextDispatch.writer).toBe(firstDispatch.writer); + expect(mainCache.isMainQuotaDispatchWsClaimed(nextDispatch)).toBe(false); + // A callback retained by a previous attempt must never claim the later fallback. + firstObserver!(new Headers()); + expect(mainCache.isMainQuotaDispatchWsClaimed(nextDispatch)).toBe(false); + await deliver(ctx, quotaResponse(quotaHeaders("99"))); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(99); + expect(getMainPolicyQuota()?.weeklyPercent).toBe(99); + // If this attempt's own observer runs, its HTTP-shaped projection is suppressed. + nextObserver!(quotaHeaders("99")); + expect(mainCache.isMainQuotaDispatchWsClaimed(nextDispatch)).toBe(true); + await deliver(ctx, quotaResponse(quotaHeaders("17"))); + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(99); + expect(getMainPolicyQuota()?.weeklyPercent).toBe(99); + }); + test("every observer invocation claims its captured dispatch before checking liveness", async () => { observe(); const ctx = materialized(); - const dispatch = ctx.mainQuotaDispatch!; const observer = codexWsQuotaObserver(ctx, provider, "gpt-5.5"); + const dispatch = ctx.mainQuotaDispatch!; expect(observer).toBeDefined(); expect(mainCache.isMainQuotaDispatchWsClaimed(dispatch)).toBe(false); mainCache.observeMainQuotaCredential("fixture-replaced-before-frame", ACCOUNT); From faca1ae2616f37cc20835725d18fd42a244fea54 Mon Sep 17 00:00:00 2001 From: JUN Date: Fri, 9 Oct 2026 22:02:03 +0900 Subject: [PATCH 7/7] fix(quota): publish main quota under the proof of the response it came from Review follow-up for #6800: plaintext V2 delivery awaits a prefix read before the plain-main branch, and a deferred reset replacement can renew the context's proof meanwhile. Delivery now captures the arrival proof before its first await and uses it for the claim check, the liveness re-check and the write. --- src/server/responses/passthrough-delivery.ts | 12 +++-- structure/providers/openai-accounts.md | 2 +- structure/providers/openai-tiers.md | 4 +- structure/transports/responses-failover.md | 2 +- structure/transports/responses.md | 2 +- .../responses-main-quota-observation.test.ts | 53 +++++++++++++++++++ 6 files changed, 66 insertions(+), 9 deletions(-) diff --git a/src/server/responses/passthrough-delivery.ts b/src/server/responses/passthrough-delivery.ts index e17f93da8f2..89cdd6c8533 100644 --- a/src/server/responses/passthrough-delivery.ts +++ b/src/server/responses/passthrough-delivery.ts @@ -409,6 +409,9 @@ export async function deliverPassthroughResponse( | "localUpstream" >, ): Promise { + const { route } = requestState; + // The proof must belong to the response whose headers are published. + const arrivalMainDispatch = liveMainQuotaDispatch(admissionState.authCtx, route.provider); const { logCtx, config, options, req } = requestContext; const { codexSafetyBufferingOptions, @@ -432,7 +435,7 @@ export async function deliverPassthroughResponse( normalizeFunctionCompletionJson, } = nativeExchange; const { commitReasoningReplayServingRoute, recordTerminalOutcomes } = responseEffects; - const { parsed, route, subagentQuotaFailureModel, clientRequestedStream, translatorBudget, inboundWire } = requestState; + const { parsed, subagentQuotaFailureModel, clientRequestedStream, translatorBudget, inboundWire } = requestState; const enforceDeclaredToolNames = inboundWire !== "chat" && inboundWire !== "anthropic"; const { openAiSidecar } = sidecarState; const { requestBindings } = transportState; @@ -539,18 +542,17 @@ export async function deliverPassthroughResponse( }); } } else { - const mainDispatch = liveMainQuotaDispatch(admissionState.authCtx, route.provider); // The WS observer is the only plain-main publisher for a WebSocket exchange; // a prelude projection carries the prelude snapshot, not fresh evidence. // The dispatch claim is authoritative because downstream wrappers can replace the Response. - if (mainDispatch && !isMainQuotaDispatchWsClaimed(mainDispatch) + if (arrivalMainDispatch && !isMainQuotaDispatchWsClaimed(arrivalMainDispatch) && !(isCodexWsUpstreamResponse(upstreamResponse) || isCodexWsPreludeProjection(upstreamResponse))) { const { applyAccountQuotaFromUpstreamHeaders } = await import("../../codex/auth-api"); // Import yields; same-account token replacement leaves the identity writer live. // Re-check the credential fence with no await before publication. - if (isMainQuotaDispatchLive(mainDispatch)) { + if (isMainQuotaDispatchLive(arrivalMainDispatch)) { applyAccountQuotaFromUpstreamHeaders(MAIN_CODEX_ACCOUNT_ID, upstreamResponse.headers, - mainDispatch.configGeneration, mainDispatch.writer, { modelId: route.modelId }); + arrivalMainDispatch.configGeneration, arrivalMainDispatch.writer, { modelId: route.modelId }); } } } diff --git a/structure/providers/openai-accounts.md b/structure/providers/openai-accounts.md index f9737bfd144..ac511652d6f 100644 --- a/structure/providers/openai-accounts.md +++ b/structure/providers/openai-accounts.md @@ -211,7 +211,7 @@ request-scoped: a translated Claude turn resolves through Pool selection like an main keeps its health, quarantine and refresh-and-classify handling. Only a bearer the client itself supplied is caller-owned and exempt from stored state. Both synchronous and asynchronous stored-main substitution in `src/codex/auth-context.ts` remove a caller account header before copying the stored identity; an absent stored account ID leaves no account header. Caller-owned native Direct authentication retains its existing passthrough behavior. -Plain-main HTTP and WebSocket Responses refresh `__main__` only on the canonical OpenAI forward provider when the sent bearer and effective workspace match the main credential already observed under native ownership, the same equality rule used by the hard lock. `src/codex/auth-context.ts` captures a process-local dispatch proof after materialization; `src/server/responses/passthrough-delivery.ts` and `src/server/responses/core-codex-account.ts` recheck its credential/identity generations before publishing. The dispatch also captures the process-wide credential mutation epoch; any OpenCodex-owned credential publication, including native main refresh or same-account reauth before a new quota credential observation, rejects an older dispatch. Publications for other credentials conservatively drop the main update as well. Same-account token rotation and A→B→A changes reject old responses; Each WebSocket observer renews the live dispatch object with every captured fence unchanged, retains that copy across frames, and claims it on every invocation before checking liveness. A later observer starts unclaimed, so its failed-upgrade HTTP fallback can publish even if the prior WS attempt observed quota. This process-local claim belongs to one physical attempt and prevents plain-main HTTP publication even when a downstream stream wrapper replaces the Response; response markers remain an additional guard, including separately marked pre-response prelude projections (4xx refusals and 502/504 gateway failures). Prelude headers remain available to Pool replay; real HTTP fallback responses still publish through HTTP delivery. An operator-granted HTTP replacement gets a new unclaimed dispatch object with every captured credential and config fence copied unchanged; the failed WS observer retains the old object. Caller-owned requests acquire no physical-main read or Pool health state. +Plain-main HTTP and WebSocket Responses refresh `__main__` only on the canonical OpenAI forward provider when the sent bearer and effective workspace match the main credential already observed under native ownership, the same equality rule used by the hard lock. `src/codex/auth-context.ts` captures a process-local dispatch proof after materialization; `src/server/responses/passthrough-delivery.ts` captures the response-arrival proof before any awaited body classification and retains it for that response's header publication; it and `src/server/responses/core-codex-account.ts` recheck credential/identity generations before publishing. The dispatch also captures the process-wide credential mutation epoch; any OpenCodex-owned credential publication, including native main refresh or same-account reauth before a new quota credential observation, rejects an older dispatch. Publications for other credentials conservatively drop the main update as well. Same-account token rotation and A→B→A changes reject old responses; Each WebSocket observer renews the live dispatch object with every captured fence unchanged, retains that copy across frames, and claims it on every invocation before checking liveness. A later observer starts unclaimed, so its failed-upgrade HTTP fallback can publish even if the prior WS attempt observed quota. This process-local claim belongs to one physical attempt and prevents plain-main HTTP publication even when a downstream stream wrapper replaces the Response; response markers remain an additional guard, including separately marked pre-response prelude projections (4xx refusals and 502/504 gateway failures). Prelude headers remain available to Pool replay; real HTTP fallback responses still publish through HTTP delivery. An operator-granted HTTP replacement gets a new unclaimed dispatch object with every captured credential and config fence copied unchanged; the failed WS observer retains the old object. Caller-owned requests acquire no physical-main read or Pool health state. `src/providers/openai-sidecar.ts` releases quota-probe ownership on every materialization or usability failure before transferring a resolved context to its caller. Audio reports one terminal upstream outcome after validating the response body; redirects remain diff --git a/structure/providers/openai-tiers.md b/structure/providers/openai-tiers.md index d313fb49af2..3e54593d065 100644 --- a/structure/providers/openai-tiers.md +++ b/structure/providers/openai-tiers.md @@ -432,7 +432,9 @@ remain process-local and never enter disk, logs, or management DTOs. Plain-main HTTP and WebSocket Responses on the canonical OpenAI forward provider refresh cached main usage under that same credential/workspace match, including stored-main substitution and an identical caller-owned credential. Materialization captures a process-local dispatch proof; -publication rechecks identity and credential generations, including after an awaited HTTP import +HTTP delivery captures the response-arrival proof before any body await and retains it for those +headers even if deferred recovery renews the context. Publication rechecks identity and credential +generations, including after an awaited HTTP import and for every WebSocket frame. A replaced credential, unmatched workspace, or custom destination cannot publish main usage. Pool health/failover handling stays scoped to Pool contexts. The dispatch additionally fences the process-wide credential mutation epoch, so native main diff --git a/structure/transports/responses-failover.md b/structure/transports/responses-failover.md index 355b61bcc36..7fe0ff326b2 100644 --- a/structure/transports/responses-failover.md +++ b/structure/transports/responses-failover.md @@ -409,7 +409,7 @@ the public server reference already documents. The 504 and a drop after the resp are never replaced. Only the 502 of a socket that closed or errored before any Responses event may be replaced over HTTP, when the provider opted into `retryOnReset` (#4191). That replacement claims from the request's one allowance; if it resets before its head, that is the pre-header row -again and may use a configured second replacement, otherwise it settles as the refusal. For plain-main quota, each physical HTTP replacement renews the dispatch object identity while copying every original credential and config fence unchanged. The failed WS observer retains its claimed object, so only the live replacement attempt may publish fresh HTTP headers. Each rebuilt WS observer also renews its own live dispatch copy; an upgrade failure with no quota frames leaves that attempt unclaimed for HTTP fallback publication. +again and may use a configured second replacement, otherwise it settles as the refusal. For plain-main quota, each physical HTTP replacement renews the dispatch object identity while copying every original credential and config fence unchanged. The failed WS observer retains its claimed object, so only the live replacement attempt may publish fresh HTTP headers. Each rebuilt WS observer also renews its own live dispatch copy; an upgrade failure with no quota frames leaves that attempt unclaimed for HTTP fallback publication. HTTP delivery retains the arrival dispatch for the arrival headers across deferred body recovery, even when that recovery renews the auth context's dispatch. This reclassification is the recorded behaviour change: before it, the pre-header refusal borrowed `upstream_closed_before_response` and its 502, which multiplied the duplicate send diff --git a/structure/transports/responses.md b/structure/transports/responses.md index db12d050427..e6fcd0b6981 100644 --- a/structure/transports/responses.md +++ b/structure/transports/responses.md @@ -23,7 +23,7 @@ Canonical forward auth retains its separate fixed credential/metadata allowlist; Retired Codex Spark has no model-specific tool or Responses Lite override; general Lite handling and namespace scrubbing remain shared compatibility behavior. Codex quota/reset evidence follows the -[shared/Reserve policy](../providers/openai-tiers.md#public-provider-contract), including suppression of retired model-derived evidence before shared recovery. Plain-main HTTP/WS quota headers refresh `__main__` only on the canonical OpenAI forward provider when the materialized bearer and workspace match the owned main observation, using the hard-lock credential-match rule. `src/server/responses/passthrough-delivery.ts` rechecks the captured credential generation after its awaited import; `src/server/responses/core-codex-account.ts` renews a live dispatch object for each WS observer without recapturing its credential or config fences, retains that copy across frames, and rechecks it for each frame. Rotated credentials and unmatched workspaces publish nothing; dispatch proofs also reject any OpenCodex-owned credential publication epoch change, including native main refresh or same-account reauth before quota re-observation; each plain-main WS observer invocation claims its captured dispatch before liveness checks, and this claim is the authoritative HTTP publication guard across Response replacement; WS upstream responses and separately marked pre-response prelude projections (4xx refusals and 502/504 gateway failures) skip plain-main HTTP quota writes; prelude headers remain available to Pool replay, real HTTP fallbacks still publish; each operator-granted HTTP replacement renews the dispatch object with its original credential and config fences unchanged, and Pool health/failover gates remain unchanged. +[shared/Reserve policy](../providers/openai-tiers.md#public-provider-contract), including suppression of retired model-derived evidence before shared recovery. Plain-main HTTP/WS quota headers refresh `__main__` only on the canonical OpenAI forward provider when the materialized bearer and workspace match the owned main observation, using the hard-lock credential-match rule. `src/server/responses/passthrough-delivery.ts` captures the response-arrival dispatch before any await and keeps that proof for the arrival headers across body classification and deferred replacement, rechecking its credential generation after the awaited import; `src/server/responses/core-codex-account.ts` renews a live dispatch object for each WS observer without recapturing its credential or config fences, retains that copy across frames, and rechecks it for each frame. Rotated credentials and unmatched workspaces publish nothing; dispatch proofs also reject any OpenCodex-owned credential publication epoch change, including native main refresh or same-account reauth before quota re-observation; each plain-main WS observer invocation claims its captured dispatch before liveness checks, and this claim is the authoritative HTTP publication guard across Response replacement; WS upstream responses and separately marked pre-response prelude projections (4xx refusals and 502/504 gateway failures) skip plain-main HTTP quota writes; prelude headers remain available to Pool replay, real HTTP fallbacks still publish; each operator-granted HTTP replacement renews the dispatch object with its original credential and config fences unchanged, and Pool health/failover gates remain unchanged. ### Credential-bearing HTTP redirects diff --git a/tests/responses/responses-main-quota-observation.test.ts b/tests/responses/responses-main-quota-observation.test.ts index 0b828c16cc3..45d0bf6e3d9 100644 --- a/tests/responses/responses-main-quota-observation.test.ts +++ b/tests/responses/responses-main-quota-observation.test.ts @@ -534,6 +534,59 @@ describe("credential-bound plain-main Responses quota", () => { }); } + test("plaintext V2 deferred reset keeps original headers bound to their arrival dispatch", async () => { + observe(); + const renewalSpy = spyOn(mainCache, "renewMainQuotaDispatchForAttempt"); + const claimSpy = spyOn(mainCache, "isMainQuotaDispatchWsClaimed"); + let httpCalls = 0; + globalThis.fetch = (async (_input: RequestInfo | URL, init?: RequestInit) => { + httpCalls++; + expect(JSON.parse(String(init?.body))).toMatchObject({ stream: true, store: false }); + expect(getAccountQuota(MAIN)).toBeNull(); + expect(getMainPolicyQuota()).toBeNull(); + if (httpCalls === 1) { + return new Response(new ReadableStream({ + async pull(controller) { + await new Promise(resolve => setTimeout(resolve, 0)); + controller.error(Object.assign(new Error("fixture pre-output reset"), { code: "ECONNRESET" })); + }, + }), { status: 200, headers: quotaHeaders("17") }); + } + return new Response(new TextEncoder().encode(`event: response.completed\ndata: ${JSON.stringify({ + type: "response.completed", response: { id: "fixture-plaintext-reset", status: "completed", output: [] }, + })}\n\n`), { status: 200, headers: quotaHeaders("99") }); + }) as typeof fetch; + try { + const response = await handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", headers: new Headers({ ...Object.fromEntries(caller()), "content-type": "application/json" }), + body: JSON.stringify({ model: "gpt-5.5", stream: true, store: false, input: "fixture delegate request", + tools: [{ type: "namespace", name: "collaboration", tools: [{ type: "function", name: "spawn_agent", + parameters: { type: "object", properties: { message: { type: "string", encrypted: true } } } }] }], + }), + }), { ...config, defaultProvider: "openai", plaintextV2AgentMessages: true, + providers: { openai: { ...provider, codexAccountMode: "direct", upstreamWebsocket: false, retryOnReset: {} } } }, + { model: "", provider: "" }, { codexWsRuntimeIdentity: BOUNDED_WS_RUNTIME }); + expect(response.status).toBe(200); + expect(await response.text()).toContain("fixture-plaintext-reset"); + expect(httpCalls).toBe(2); + expect(renewalSpy).toHaveBeenCalledTimes(2); + const arrival = renewalSpy.mock.results[0]!.value as mainCache.MainQuotaDispatch; + const successor = renewalSpy.mock.results[1]!.value as mainCache.MainQuotaDispatch; + expect(successor).not.toBe(arrival); + expect(successor).toEqual(arrival); + // The prefix probe awaits the replacement before publication; the original proof stays live. + // Its 17% may publish only under that proof. Replacement-header publication is a follow-up. + expect(getAccountQuota(MAIN)?.weeklyPercent).toBe(17); + expect(getMainPolicyQuota()?.weeklyPercent).toBe(17); + expect(claimSpy).toHaveBeenCalledTimes(1); + expect(claimSpy.mock.calls[0]![0]).toBe(arrival); + expect(claimSpy.mock.calls[0]![0]).not.toBe(successor); + } finally { + claimSpy.mockRestore(); + renewalSpy.mockRestore(); + } + }); + test("full handler sanitized recovery with failed WS upgrade publishes fresh HTTP quota", async () => { observe(); const originalSocket = globalThis.WebSocket;