From e12dcc6d0d3af38c6f0118a6efae912d82f5b6ab Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Thu, 17 Sep 2026 14:58:38 +0900 Subject: [PATCH 1/3] fix(adapters): charge adapter-owned retry sends to the request send budget mimo-free's 401 JWT replay, command-code's reasoning-effort repair, and the google-http transient loop each issued bare fetches that never touched ctx.sendBudget, so a request holding only its final recovery permit still dispatched and a refused retry still paid the backoff sleep. Route each physical send through createAdapterPhysicalSend: admission precedes pacing, backoff and superseded-response cancellation, a credential hop's pending permit pays for the first send exactly once, and a refused retry returns the real upstream response instead of a synthetic error. Follow-up to #4621. --- src/adapters/command-code.ts | 17 ++- src/adapters/google-http.ts | 51 ++++++--- src/adapters/mimo-free.ts | 46 +++++--- src/adapters/physical-send.ts | 50 ++++++++ .../google/google-vertex-http.test.ts | 42 ++++++- tests/adapters/physical-send.test.ts | 107 ++++++++++++++++++ tests/fixtures/test-layout-expected.json | 1 + tests/helpers/send-budget-owner.ts | 23 ++++ tests/providers/command-code-provider.test.ts | 89 ++++++++++++++- tests/providers/mimo-free-provider.test.ts | 36 +++++- .../responses/responses-core-modules.test.ts | 18 +-- 11 files changed, 417 insertions(+), 63 deletions(-) create mode 100644 src/adapters/physical-send.ts create mode 100644 tests/adapters/physical-send.test.ts create mode 100644 tests/helpers/send-budget-owner.ts diff --git a/src/adapters/command-code.ts b/src/adapters/command-code.ts index 3c466829e2..111319236a 100644 --- a/src/adapters/command-code.ts +++ b/src/adapters/command-code.ts @@ -13,6 +13,8 @@ import { commandCodeReasoningEfforts, refreshCommandCodeReasoningEfforts } from import { identifyRoutedModel } from "./identity"; import { buildNonOpenAIToolCatalogNudgeForTools } from "./tool-catalog-nudge"; import { parseDataUrl } from "./image"; +import { createAdapterPhysicalSend } from "./physical-send"; +import { SendBudgetExhaustedError } from "../lib/upstream-retry"; // Retain the short ids emitted by the first local integration. New requests use the live catalog's // provider-native IDs directly; this map is compatibility-only and is not a model fallback list. @@ -469,7 +471,7 @@ async function fetchCommandCode(request: AdapterRequest, ctx: AdapterFetchContex const timer = setTimeout(() => timeout.abort(new DOMException("Timeout elapsed", "TimeoutError")), ctx?.timeoutMs ?? 200_000); const callerSignal = ctx?.abortSignal ?? new AbortController().signal; try { - return await (ctx?.executor ?? executor)(request.url, { + return await executor(request.url, { method: request.method, headers: request.headers, body: request.body, @@ -556,7 +558,8 @@ export function createCommandCodeAdapter(provider: OcxProviderConfig): ProviderA }; }, async fetchResponse(request: AdapterRequest, ctx?: AdapterFetchContext): Promise { - const response = await fetchCommandCode(request, ctx, executor); + const send = createAdapterPhysicalSend(ctx, executor); + const response = await send({ url: request.url, dispatch: physical => fetchCommandCode(request, ctx, physical) }); if (response.ok) return response; const currentEffort = (() => { try { return (JSON.parse(request.body) as { params?: { reasoning_effort?: unknown } }).params?.reasoning_effort; } catch { return undefined; } @@ -577,8 +580,14 @@ export function createCommandCodeAdapter(provider: OcxProviderConfig): ProviderA if (!refreshed || refreshed.includes(currentEffort)) return response; const retry = requestWithoutReasoningEffort(request); if (!retry) return response; - try { void response.body?.cancel(); } catch { /* already closed */ } - return fetchCommandCode(retry, ctx, executor); + try { + return await send({ url: retry.url, sendClass: "repair", recovery: "reasoning-effort-downgrade", + beforeDispatch: () => { try { void response.body?.cancel().catch(() => {}); } catch { /* already closed */ } }, + dispatch: physical => fetchCommandCode(retry, ctx, physical) }); + } catch (error) { + if (error instanceof SendBudgetExhaustedError) return response; + throw error; + } }, async *parseStream(response: Response, budget: TranslatorBudget): AsyncGenerator { let sawFinish = false; diff --git a/src/adapters/google-http.ts b/src/adapters/google-http.ts index f7b90de87e..de1332dae7 100644 --- a/src/adapters/google-http.ts +++ b/src/adapters/google-http.ts @@ -1,4 +1,7 @@ import type { AdapterFetchContext, AdapterRequest } from "./base"; +import { createAdapterPhysicalSend } from "./physical-send"; +import type { SendClass } from "../lib/request-execution-budget"; +import type { AttemptRecoveryKind } from "../usage/log"; import { isQuotaExhaustedBody, retryableGoogleStatus, safeGoogleHttpErrorMessage } from "./google-errors"; import { repairGoogleInvalidRequestBody } from "./google-wire-compiler"; import { normalizeUpstreamHttpErrorResponse, readDisplaySafeErrorPayloadText } from "./upstream-http-error"; @@ -8,6 +11,8 @@ import { fetchWithAttemptDeadline, retryBackoffDelayMs, sleepWithAbort, + SendBudgetExhaustedError, + isConnectionResetError, } from "../lib/upstream-retry"; const GOOGLE_RETRY_ATTEMPTS = 3; @@ -41,18 +46,30 @@ export async function fetchGoogleWithRetry( ): Promise { const repairInvalid400 = opts.repairInvalid400 ?? true; const timeoutMs = ctx.timeoutMs ?? 200_000; - const executor = ctx.executor ?? globalThis.fetch; + const send = createAdapterPhysicalSend(ctx); let lastError: unknown; let activeRequest = request; let compatibilityReplayUsed = false; + let pendingResponse: Response | undefined; + let retryDelayMs = 0; + let sendClass: SendClass = "transient"; + let recovery: AttemptRecoveryKind | undefined; for (let attempt = 0; attempt < GOOGLE_RETRY_ATTEMPTS; attempt++) { if (ctx.abortSignal?.aborted) throw abortError(ctx.abortSignal); try { - const res = await fetchWithAttemptDeadline(activeRequest.url, { - method: activeRequest.method, - headers: activeRequest.headers, - body: activeRequest.body, - }, timeoutMs, ctx.abortSignal, ctx.stream, executor); + const res = await send({ url: activeRequest.url, sendClass, recovery, + beforeDispatch: async () => { + if (retryDelayMs > 0) await sleepWithAbort(retryDelayMs, ctx.abortSignal); + if (pendingResponse) cancelResponseBodyBestEffort(pendingResponse); + pendingResponse = undefined; + }, + dispatch: executor => fetchWithAttemptDeadline(activeRequest.url, { + method: activeRequest.method, headers: activeRequest.headers, body: activeRequest.body, + }, timeoutMs, ctx.abortSignal, ctx.stream, executor), + }); + retryDelayMs = 0; + sendClass = "transient"; + recovery = undefined; if (res.status === 400 && repairInvalid400 && !compatibilityReplayUsed) { let payloadText = ""; try { @@ -64,7 +81,8 @@ export async function fetchGoogleWithRetry( if (repairedBody !== undefined) { compatibilityReplayUsed = true; activeRequest = { ...activeRequest, body: repairedBody }; - cancelResponseBodyBestEffort(res); + pendingResponse = res; + sendClass = "repair"; attempt--; // The changed-request replay is separate from transient retry accounting. continue; } @@ -75,7 +93,7 @@ export async function fetchGoogleWithRetry( // A 429 may be a transient rate limit (retry) or hard quota exhaustion (do NOT retry — // it won't recover for hours and burns retries). Peek the body to tell them apart. if (res.status === 429) { - const peekTarget = ctx.returnRawErrors ? res.clone() : res; + const peekTarget = res.clone(); const peek = await readDisplaySafeErrorPayloadText(peekTarget, ctx.abortSignal); if (isQuotaExhaustedBody(peek)) { return ctx.returnRawErrors ? res : normalizeUpstreamHttpErrorResponse(res, { @@ -84,20 +102,27 @@ export async function fetchGoogleWithRetry( }); } } - cancelResponseBodyBestEffort(res); - await sleepWithAbort(retryBackoffDelayMs(attempt, { + pendingResponse = res; + recovery = res.status === 429 ? "rate-limit-429" : "transient-5xx"; + retryDelayMs = retryBackoffDelayMs(attempt, { baseDelayMs: GOOGLE_RETRY_BASE_MS, maxDelayMs: GOOGLE_RETRY_MAX_MS, headers: res.headers, - }), ctx.abortSignal); + }); } catch (err) { if (ctx.abortSignal?.aborted) throw err; + if (err instanceof SendBudgetExhaustedError) { + if (pendingResponse) return ctx.returnRawErrors ? pendingResponse : normalizeFinalGoogleError(label, pendingResponse, ctx.abortSignal); + throw err; + } lastError = err; if (attempt === GOOGLE_RETRY_ATTEMPTS - 1) throw err; - await sleepWithAbort(retryBackoffDelayMs(attempt, { + sendClass = "transient"; + recovery = isConnectionResetError(err) ? "connection-reset" : undefined; + retryDelayMs = retryBackoffDelayMs(attempt, { baseDelayMs: GOOGLE_RETRY_BASE_MS, maxDelayMs: GOOGLE_RETRY_MAX_MS, - }), ctx.abortSignal); + }); } } throw lastError ?? new Error(`${label} fetch failed`); diff --git a/src/adapters/mimo-free.ts b/src/adapters/mimo-free.ts index 6500e023b5..668691b2a1 100644 --- a/src/adapters/mimo-free.ts +++ b/src/adapters/mimo-free.ts @@ -6,6 +6,8 @@ import { recordOwnedConfigPath } from "../lib/config-ownership"; import type { OcxProviderConfig, OcxParsedRequest } from "../types"; import { createOpenAIChatAdapter } from "./openai-chat"; import type { ProviderAdapter, AdapterRequest, IncomingMeta } from "./base"; +import { createAdapterPhysicalSend } from "./physical-send"; +import { SendBudgetExhaustedError } from "../lib/upstream-retry"; const BOOTSTRAP_URL = "https://api.xiaomimimo.com/api/free-ai/bootstrap"; export const MIMO_CHAT_URL = "https://api.xiaomimimo.com/api/free-ai/openai/chat"; @@ -248,33 +250,43 @@ export function createMimoFreeAdapter(provider: OcxProviderConfig): ProviderAdap }, async fetchResponse(request: AdapterRequest, ctx): Promise { - const response = await fetch(request.url, { + const send = createAdapterPhysicalSend(ctx); + const response = await send({ url: request.url, dispatch: executor => executor(request.url, { method: request.method, redirect: "manual", headers: request.headers as Record, body: request.body, signal: ctx?.abortSignal, - }); + }) }); // Retry predicate: 401 (expired/invalid JWT) retries ONCE with a fresh token. // 403 is NOT retried — Xiaomi uses it for anti-abuse "Illegal access" and there is // no documented token-expiry signature that would mark a 403 as retryable. if (response.status === 401) { - // Drain the first response body before issuing the retry. - try { await response.body?.cancel(); } catch { /* already consumed */ } - resetMimoJwtCache(); - const freshJwt = await getMimoJwt(ctx?.abortSignal); - const retryHeaders = { - ...(request.headers as Record), - "Authorization": `Bearer ${freshJwt}`, - }; - return fetch(request.url, { - method: request.method, - redirect: "manual", - headers: retryHeaders, - body: request.body, - signal: ctx?.abortSignal, - }); + let retryHeaders = request.headers; + try { + return await send({ url: request.url, sendClass: "auth-recovery", recovery: "oauth-401", + beforeDispatch: async () => { + // Drain the first response body and refresh the JWT only after admission. + resetMimoJwtCache(); + const freshJwt = await getMimoJwt(ctx?.abortSignal); + retryHeaders = { + ...(request.headers as Record), + "Authorization": `Bearer ${freshJwt}`, + }; + try { void response.body?.cancel().catch(() => {}); } catch { /* already consumed */ } + }, + dispatch: executor => executor(request.url, { + method: request.method, + redirect: "manual", + headers: retryHeaders, + body: request.body, + signal: ctx?.abortSignal, + }) }); + } catch (error) { + if (error instanceof SendBudgetExhaustedError) return response; + throw error; + } } return response; diff --git a/src/adapters/physical-send.ts b/src/adapters/physical-send.ts new file mode 100644 index 0000000000..9f7c2a9f93 --- /dev/null +++ b/src/adapters/physical-send.ts @@ -0,0 +1,50 @@ +import type { AdapterFetchContext } from "./base"; +import type { SendClass } from "../lib/request-execution-budget"; +import type { AttemptRecoveryKind } from "../usage/log"; +import { abortError, SendBudgetExhaustedError } from "../lib/upstream-retry"; + +type PacedFetch = typeof globalThis.fetch & { + waitForPacing?: (signal?: AbortSignal) => Promise; + unpacedFetch?: typeof globalThis.fetch; +}; + +/** One ordinal sequence per adapter fetchResponse call, across all of its inference retries. + * Consumption starts at underlying executor invocation; its own later preflight may still fail. */ +export function createAdapterPhysicalSend(ctx: AdapterFetchContext = {}, fallback = globalThis.fetch) { + const executor = (ctx.executor ?? fallback) as PacedFetch; + let ordinal = 0; + return async (options: { + url: string; + sendClass?: SendClass; + recovery?: AttemptRecoveryKind; + /** Runs only after admission, e.g. backoff and cancellation of a superseded response. */ + beforeDispatch?: () => void | Promise; + dispatch: (executor: typeof globalThis.fetch) => Promise; + }): Promise => { + if (ctx.abortSignal?.aborted) throw abortError(ctx.abortSignal); + const decision = ctx.sendBudget?.reserveDispatch({ + sendClass: options.sendClass ?? "transient", targetKey: options.url, + }); + if (decision && !decision.allowed) throw new SendBudgetExhaustedError(options.url); + const permit = decision?.allowed ? decision.permit : undefined; + let dispatched = false; + const physicalExecutor = (async (input, init) => { + if (ctx.abortSignal?.aborted) throw abortError(ctx.abortSignal); + if (init?.signal?.aborted) throw abortError(init.signal); + if (dispatched || (permit && !permit.use())) throw new SendBudgetExhaustedError(options.url); + dispatched = true; + ordinal += 1; + ctx.onPhysicalSend?.({ ordinal, ...(options.recovery ? { recovery: options.recovery } : {}) }); + return (executor.unpacedFetch ?? executor)(input, init); + }) as typeof globalThis.fetch; + try { + await executor.waitForPacing?.(ctx.abortSignal); + if (ctx.abortSignal?.aborted) throw abortError(ctx.abortSignal); + await options.beforeDispatch?.(); + if (ctx.abortSignal?.aborted) throw abortError(ctx.abortSignal); + return await options.dispatch(physicalExecutor); + } finally { + permit?.release(); + } + }; +} diff --git a/tests/adapters/google/google-vertex-http.test.ts b/tests/adapters/google/google-vertex-http.test.ts index 7d91793226..e2cb4aef7b 100644 --- a/tests/adapters/google/google-vertex-http.test.ts +++ b/tests/adapters/google/google-vertex-http.test.ts @@ -1,4 +1,7 @@ -import { afterEach, describe, expect, test } from "bun:test"; +import { afterEach, describe, expect, spyOn, test } from "bun:test"; +import * as retry from "../../../src/lib/upstream-retry"; +import { createRequestExecutionBudget } from "../../../src/lib/request-execution-budget"; +import { budgetOwner } from "../../helpers/send-budget-owner"; import type { AdapterRequest } from "../../../src/adapters/base"; import { fetchAntigravityWithRetry, fetchDirectGeminiWithRetry, fetchVertexWithRetry } from "../../../src/adapters/google-http"; import { safeVertexHttpErrorMessage, retryableGoogleStatus } from "../../../src/adapters/google-errors"; @@ -30,6 +33,43 @@ function vertexError(code: number, status: string, message: string): string { } describe("vertex retry fetch", () => { + for (const [name, fetchResponse] of [["Vertex", fetchVertexWithRetry], ["Antigravity", fetchAntigravityWithRetry]] as const) { + test.each([400, 429, 503, "reset"] as const)(`${name} prepaid final send prevents another inference or backoff (%s)`, async status => { + const parent = createRequestExecutionBudget(); + parent.used = 3; + const { owner, dispose } = budgetOwner(parent); + const raw = status === 400 ? vertexError(400, "INVALID_ARGUMENT", "tools.0.custom.input_schema: JSON schema is invalid") : `fixture ${status}`; + const first = status === "reset" ? Object.assign(new Error("fixture reset"), { code: "ECONNRESET" }) + : new Response(raw, { status, headers: { "Retry-After": "60" } }); + const fixture = mockFetch([first, new Response("unexpected replay")]); + const waits = spyOn(retry, "sleepWithAbort").mockImplementation(async () => {}); + const ordinals: number[] = []; + try { + const hop = owner.reserveCredentialHop("auth-recovery", request.url, true); + if (!hop.allowed || !hop.permit) throw new Error("Expected final prepaid send"); + owner.pendingHopPermit = hop.permit; + const scope = owner.adapterDispatchBudget; + if (!scope) throw new Error("Expected an adapter dispatch budget"); + const result = fetchResponse({ ...request, body: JSON.stringify({ request: { + contents: [{ role: "user", parts: [{ text: "hi" }] }], + tools: [{ functionDeclarations: [{ name: "replace_in_files", parameters: { + type: "object", properties: { occurrence_ids: { type: "array", items: { type: "string" } } }, + } }] }], + } }) }, { sendBudget: scope, returnRawErrors: true, onPhysicalSend: send => ordinals.push(send.ordinal) }); + if (status === "reset") await expect(result).rejects.toBeInstanceOf(retry.SendBudgetExhaustedError); + else { + const response = await result; + expect(response).toBe(first); + expect(await response.text()).toBe(raw); + } + expect(fixture.calls).toHaveLength(1); + expect(waits).not.toHaveBeenCalled(); + expect(parent.used).toBe(4); + expect(ordinals).toEqual([1]); + } finally { waits.mockRestore(); dispose(); } + }); + } + test("successful response bodies survive beyond the response-header timeout", async () => { globalThis.fetch = (async () => new Response(new ReadableStream({ async start(controller) { diff --git a/tests/adapters/physical-send.test.ts b/tests/adapters/physical-send.test.ts new file mode 100644 index 0000000000..1b76ce0e4a --- /dev/null +++ b/tests/adapters/physical-send.test.ts @@ -0,0 +1,107 @@ +import { describe, expect, test } from "bun:test"; +import { createAdapterPhysicalSend } from "../../src/adapters/physical-send"; +import { createRequestExecutionBudget } from "../../src/lib/request-execution-budget"; +import { SendBudgetExhaustedError } from "../../src/lib/upstream-retry"; +import { budgetOwner } from "../helpers/send-budget-owner"; + +const url = "https://adapter-fixture.invalid/inference"; + +/** + * A credential hop has already reserved the replay it hands to the adapter, so the adapter's + * first send spends that permit through the dispatch view instead of reserving again. + */ +function prepaid() { + const parent = createRequestExecutionBudget(); + parent.used = 3; + const { owner, dispose } = budgetOwner(parent); + const hop = owner.reserveCredentialHop("auth-recovery", url, true); + if (!hop.allowed || !hop.permit) throw new Error("Expected prepaid final send"); + owner.pendingHopPermit = hop.permit; + const scope = owner.adapterDispatchBudget; + if (!scope) throw new Error("Expected an adapter dispatch budget"); + return { parent, scope, dispose }; +} + +describe("adapter physical inference admission", () => { + test("a prepaid scope admits exactly one physical send and rejects replay before backoff", async () => { + const { parent, scope, dispose } = prepaid(); + let sends = 0, waits = 0, pacingSlots = 0; + const ordinals: number[] = []; + try { + const send = createAdapterPhysicalSend({ sendBudget: scope, onPhysicalSend: event => ordinals.push(event.ordinal) }, + Object.assign(async () => { sends += 1; return new Response("ok"); }, { + waitForPacing: async () => { pacingSlots += 1; }, + }) as typeof fetch); + await send({ url, dispatch: executor => executor(url) }); + await expect(send({ url, sendClass: "repair", beforeDispatch: () => { waits += 1; }, + dispatch: executor => executor(url) })).rejects.toBeInstanceOf(SendBudgetExhaustedError); + expect(sends).toBe(1); + expect(waits).toBe(0); + expect(pacingSlots).toBe(1); + expect(ordinals).toEqual([1]); + expect(parent.used).toBe(4); + } finally { dispose(); } + }); + + test.each(["pacing", "backoff", "abort", "adapter"] as const)("unused reservation refunds after %s refusal", async phase => { + const parent = createRequestExecutionBudget(); + parent.used = 3; + let sends = 0; + const controller = new AbortController(); + const failure = new Error(`fixture ${phase} refusal`); + const executor = Object.assign(async () => { sends += 1; return new Response("unexpected"); }, { + waitForPacing: async () => { if (phase === "pacing") throw failure; }, + }) as typeof fetch; + const send = createAdapterPhysicalSend({ sendBudget: parent, abortSignal: controller.signal }, executor); + // A reserve-funded class still gets a real permit once the base allowance is spent; the + // refusal paths below never reach its dispatch, so the reservation must be handed back. + await expect(send({ url, sendClass: "repair", beforeDispatch: () => { + if (phase === "backoff") throw failure; + if (phase === "abort") controller.abort(failure); + }, dispatch: physical => { + if (phase === "adapter") throw failure; + return physical(url); + } })).rejects.toBe(failure); + expect(parent.used).toBe(3); + expect(parent.reserveSpent).toBe(false); + expect(sends).toBe(0); + }); + + test.each(["pacing", "backoff", "abort", "adapter"] as const)("a settled hop charge stays charged when the %s leg never dispatches", async phase => { + const { parent, scope, dispose } = prepaid(); + let sends = 0; + const controller = new AbortController(); + const failure = new Error(`fixture ${phase} refusal`); + const executor = Object.assign(async () => { sends += 1; return new Response("unexpected"); }, { + waitForPacing: async () => { if (phase === "pacing") throw failure; }, + }) as typeof fetch; + try { + const send = createAdapterPhysicalSend({ sendBudget: scope, abortSignal: controller.signal }, executor); + await expect(send({ url, beforeDispatch: () => { + if (phase === "backoff") throw failure; + if (phase === "abort") controller.abort(failure); + }, dispatch: physical => { + if (phase === "adapter") throw failure; + return physical(url); + } })).rejects.toBe(failure); + // The hop's reservation was the charge and the dispatch view settled it at admission; + // the adapter's release has nothing left to refund. + expect(parent.used).toBe(4); + expect(parent.reserveSpent).toBe(true); + expect(sends).toBe(0); + } finally { dispose(); } + }); + + test("an exhausted initial send performs no inference or retry preparation", async () => { + const budget = createRequestExecutionBudget(); + budget.used = 4; + let prepared = false, sends = 0; + const send = createAdapterPhysicalSend({ sendBudget: budget }, + (async () => { sends += 1; return new Response("unexpected"); }) as typeof fetch); + await expect(send({ url, beforeDispatch: () => { prepared = true; }, + dispatch: physical => physical(url) })).rejects.toBeInstanceOf(SendBudgetExhaustedError); + expect(prepared).toBe(false); + expect(sends).toBe(0); + expect(budget.used).toBe(4); + }); +}); diff --git a/tests/fixtures/test-layout-expected.json b/tests/fixtures/test-layout-expected.json index 36a7c442d2..4d7999425a 100644 --- a/tests/fixtures/test-layout-expected.json +++ b/tests/fixtures/test-layout-expected.json @@ -20,6 +20,7 @@ "adapter-event-oauth-failover.test.ts": "oauth", "adapter-inner-send-budget-wiring.test.ts": "adapters", "adapter-inner-send-budget.test.ts": "adapters", + "physical-send.test.ts": "adapters", "adapter-registry-authority.test.ts": "adapters", "adapter-resolve.test.ts": "server", "adapter-tool-conformance.test.ts": "adapters", diff --git a/tests/helpers/send-budget-owner.ts b/tests/helpers/send-budget-owner.ts new file mode 100644 index 0000000000..15a1b96cbe --- /dev/null +++ b/tests/helpers/send-budget-owner.ts @@ -0,0 +1,23 @@ +import { createResponsesSendBudget } from "../../src/server/responses/request-send-budget"; +import { createTranslatorBudget } from "../../src/lib/translator-budget"; +import type { TransientSendBudget } from "../../src/lib/upstream-retry"; + +/** + * The production send-budget owner wired the way the responses stack wires it: the adapter's + * dispatch view, credential-hop reservations, and the pending hop permit all come from + * `createResponsesSendBudget`, so a test that needs a prepaid hop exercises the same path a + * real failover takes instead of reimplementing the view. + */ +export function budgetOwner(sendBudget: TransientSendBudget) { + const translatorBudget = createTranslatorBudget(); + const result = createResponsesSendBudget({ + req: new Request("http://localhost/v1/responses"), + logCtx: { model: "test", provider: "test" }, + options: { translatorBudget, sendBudget }, + }); + if (result instanceof Response) { + translatorBudget.dispose(); + throw new Error("Unexpected workflow refusal without a workflow root"); + } + return { owner: result, dispose: () => translatorBudget.dispose() }; +} diff --git a/tests/providers/command-code-provider.test.ts b/tests/providers/command-code-provider.test.ts index d588c550ee..91980f5fd5 100644 --- a/tests/providers/command-code-provider.test.ts +++ b/tests/providers/command-code-provider.test.ts @@ -1,4 +1,14 @@ import { afterEach, describe, expect, test } from "bun:test"; +import { mkdtempSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { saveCredential, getAccountSet, setActiveAccount } from "../../src/oauth/store"; +import { clearGenericFailoverHealth } from "../../src/oauth/generic-account-failover"; +import { createRequestExecutionBudget, CODEX_TEXT_GUARDED_BUDGET_POLICY } from "../../src/lib/request-execution-budget"; +import { handleResponses } from "../../src/server/responses"; +import { removeTreeWithRetry } from "../helpers/remove-tree"; +import { budgetOwner } from "../helpers/send-budget-owner"; +import type { OcxConfig } from "../../src/types"; import { commandCodeSessionId, createCommandCodeAdapter } from "../../src/adapters/command-code"; import { loginCommandCode, parseCommandCodeCallback, shouldImportLocalCommandCodeAuth } from "../../src/oauth/command-code"; import { buildModelsRequest, OAUTH_PROVIDERS } from "../../src/oauth"; @@ -39,6 +49,50 @@ async function builtRequest(...args: Parameters resetCommandCodeReasoningEffortsForTest()); describe("Command Code provider", () => { + test("empty-completion OAuth continuation counts initial sends and keeps its prepaid hop charged", async () => { + const previousHome = process.env.OPENCODEX_HOME; + const fixtureHome = mkdtempSync(join(tmpdir(), "ocx-command-hop-")); + process.env.OPENCODEX_HOME = fixtureHome; + const originalFetch = globalThis.fetch; + clearGenericFailoverHealth(); + try { + for (let index = 0; index < 4; index++) await saveCredential("command-code", { + access: `synthetic-command-${index}`, refresh: `synthetic-refresh-${index}`, + expires: Date.now() + 3_600_000, accountId: `fixture-${index}`, source: "oauth", + }, { addAccount: true }); + await setActiveAccount("command-code", getAccountSet("command-code")!.accounts[0]!.id); + // Every physical inference send, including the initial and continuation, shares this cap. + const budget = createRequestExecutionBudget({ ...CODEX_TEXT_GUARDED_BUDGET_POLICY, + maxTotalModelSends: 3, baseSendAllowance: 3, finalRecoveryAllowance: 0 }); + const authorizations: string[] = []; + globalThis.fetch = (async (input, init) => { + const url = input instanceof Request ? input.url : String(input); + if (url !== "https://api.commandcode.ai/alpha/generate") throw new Error(`Unexpected fixture request: ${url}`); + authorizations.push(new Headers(init?.headers).get("authorization") ?? ""); + if (authorizations.length === 1) return new Response('{"type":"finish","finishReason":"stop"}\n'); + return Response.json({ error: { message: "rate limited" } }, { status: 429 }); + }) as typeof fetch; + const cfg = { defaultProvider: "command-code", emptyCompletionRetry: true, providers: { + "command-code": { adapter: "command-code", baseUrl: "https://api.commandcode.ai", authMode: "oauth", + models: ["deepseek/deepseek-v4-flash"] }, + } } as OcxConfig; + const response = await handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", headers: { "content-type": "application/json" }, + body: JSON.stringify({ model: "command-code/deepseek/deepseek-v4-flash", input: "hello", stream: false }), + }), cfg, { model: "", provider: "" }, { sendBudget: budget }); + await response.text(); + expect(authorizations).toEqual(["Bearer synthetic-command-0", "Bearer synthetic-command-0", "Bearer synthetic-command-1"]); + expect(budget.used).toBe(3); + expect(getAccountSet("command-code")!.activeAccountId).toBe(getAccountSet("command-code")!.accounts[1]!.id); + } finally { + globalThis.fetch = originalFetch; + clearGenericFailoverHealth(); + if (previousHome === undefined) delete process.env.OPENCODEX_HOME; + else process.env.OPENCODEX_HOME = previousHome; + removeTreeWithRetry(fixtureHome); + } + }, 20_000); + test("registry and OAuth surfaces stay in parity", () => { const registry = PROVIDER_REGISTRY.find(row => row.id === "command-code"); expect(registry).toMatchObject({ @@ -650,7 +704,8 @@ describe("Command Code provider", () => { expect(JSON.parse(bareBuilt.body).params.tools).toEqual(tools); }); - test("refreshes a stale official effort record only after a reasoning rejection and retries without it", async () => { + test.each(["fallback", "supplied", "prepaid"] as const)("refreshes stale effort metadata separately from inference executor (%s)", async mode => { + const supplied = mode !== "fallback"; const requests: Array<{ url: string; body?: string }> = []; const fetch = (async (url: string | URL | Request, init?: RequestInit) => { const href = String(url); @@ -664,11 +719,35 @@ describe("Command Code provider", () => { }) as typeof globalThis.fetch; const adapter = createCommandCodeAdapter({ ...provider, fetch } as OcxProviderConfig); const request = await adapter.buildRequest({ ...parsed(), options: { reasoning: "max" } }); - const response = await adapter.fetchResponse!(request); - expect(response.ok).toBe(true); + let suppliedCalls = 0; + const executor = (async (input, init) => { + expect(String(input).endsWith("/alpha/generate")).toBe(true); + suppliedCalls += 1; + return fetch(input, init); + }) as typeof globalThis.fetch; + const budget = createRequestExecutionBudget(); + const { owner, dispose } = budgetOwner(budget); + try { + if (mode === "prepaid") { + budget.used = 3; + const hop = owner.reserveCredentialHop("auth-recovery", request.url, true); + if (!hop.allowed || !hop.permit) throw new Error("Expected final prepaid send"); + owner.pendingHopPermit = hop.permit; + } + const scope = mode === "prepaid" ? owner.adapterDispatchBudget : budget; + const observed: number[] = []; + const response = await adapter.fetchResponse!(request, { ...(supplied ? { executor } : {}), sendBudget: scope, + onPhysicalSend: send => observed.push(send.ordinal) }); + expect(suppliedCalls).toBe(supplied ? mode === "prepaid" ? 1 : 2 : 0); + expect(response.ok).toBe(mode !== "prepaid"); + expect(budget.used).toBe(mode === "prepaid" ? 4 : 2); + expect(observed).toEqual(mode === "prepaid" ? [1] : [1, 2]); + const generated = requests.filter(request => request.url.endsWith("/alpha/generate")); + expect(generated).toHaveLength(mode === "prepaid" ? 1 : 2); + if (mode === "prepaid") expect(await response.text()).toContain("unsupported reasoning_effort"); + else expect(JSON.parse(generated[1]!.body!).params).not.toHaveProperty("reasoning_effort"); + } finally { dispose(); } expect(commandCodeReasoningEfforts("deepseek/deepseek-v4-flash")).toEqual(["high"]); - const generated = requests.filter(request => request.url.endsWith("/alpha/generate")); - expect(JSON.parse(generated[1]!.body!).params).not.toHaveProperty("reasoning_effort"); }); // Pins the profileUrl of each id added for #2647 — nothing more. diff --git a/tests/providers/mimo-free-provider.test.ts b/tests/providers/mimo-free-provider.test.ts index 00176d2589..8425ada95b 100644 --- a/tests/providers/mimo-free-provider.test.ts +++ b/tests/providers/mimo-free-provider.test.ts @@ -16,6 +16,8 @@ import { mkdtempSync } from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; import { removeTreeWithRetry } from "../helpers/remove-tree"; +import { createRequestExecutionBudget } from "../../src/lib/request-execution-budget"; +import { budgetOwner } from "../helpers/send-budget-owner"; for (const phase of ["bootstrap", "chat", "401-replay"] as const) test.each([307, 308])(`MiMo ${phase} never follows %i`, async status => { const nativeFetch = globalThis.fetch; @@ -319,7 +321,8 @@ describe("mimo-free auth retry predicate", () => { return createMimoFreeAdapter(provider); } - test("401 retries exactly once with a fresh JWT after draining the first body", async () => { + test.each(["fallback", "supplied", "prepaid"] as const)("401 retry preserves inference admission and separate bootstrap (%s)", async mode => { + const supplied = mode !== "fallback"; const fakeJwt = "h." + Buffer.from(JSON.stringify({ exp: Math.floor(Date.now() / 1000) + 3600 })).toString("base64") + ".s"; const calls: string[] = []; const originalFetch = globalThis.fetch; @@ -336,21 +339,42 @@ describe("mimo-free auth retry predicate", () => { } return new Response(JSON.stringify({ ok: true }), { status: 200 }); }) as unknown as typeof fetch; + const budget = createRequestExecutionBudget(); + const { owner, dispose } = budgetOwner(budget); try { const adapter = adapterForRetry(); + let suppliedCalls = 0; + const executor = (async (input, init) => { + expect(String(input)).toBe(MIMO_CHAT_URL); + suppliedCalls += 1; + return globalThis.fetch(input, init); + }) as typeof fetch; + if (mode === "prepaid") { + budget.used = 3; + const hop = owner.reserveCredentialHop("auth-recovery", MIMO_CHAT_URL, true); + if (!hop.allowed || !hop.permit) throw new Error("Expected final prepaid send"); + owner.pendingHopPermit = hop.permit; + } + const scope = mode === "prepaid" ? owner.adapterDispatchBudget : budget; + const observed: number[] = []; const res = await adapter.fetchResponse!( { url: MIMO_CHAT_URL, method: "POST", headers: { "Authorization": "Bearer stale" }, body: "{}" }, - {} as never, + { ...(supplied ? { executor } : {}), sendBudget: scope, onPhysicalSend: send => observed.push(send.ordinal) }, ); - expect(res.status).toBe(200); + expect(suppliedCalls).toBe(supplied ? mode === "prepaid" ? 1 : 2 : 0); + expect(res.status).toBe(mode === "prepaid" ? 401 : 200); + expect(budget.used).toBe(mode === "prepaid" ? 4 : 2); + expect(observed).toEqual(mode === "prepaid" ? [1] : [1, 2]); // Sequence: first chat with stale token -> 401 -> bootstrap -> retry with fresh JWT. expect(calls[0]).toBe("chat:Bearer stale"); - expect(calls[1]).toBe("bootstrap"); - expect(calls[2]).toBe(`chat:Bearer ${fakeJwt}`); - expect(calls.length).toBe(3); + if (mode !== "prepaid") expect(calls[1]).toBe("bootstrap"); + if (mode === "prepaid") expect(await res.text()).toBe("expired"); + else expect(calls[2]).toBe(`chat:Bearer ${fakeJwt}`); + expect(calls.length).toBe(mode === "prepaid" ? 1 : 3); } finally { globalThis.fetch = originalFetch; resetMimoJwtCache(); + dispose(); } }); diff --git a/tests/responses/responses-core-modules.test.ts b/tests/responses/responses-core-modules.test.ts index 1d9ff1f7aa..589a6176cd 100644 --- a/tests/responses/responses-core-modules.test.ts +++ b/tests/responses/responses-core-modules.test.ts @@ -5,10 +5,8 @@ import { RESPONSES_CORE_MODULES, readResponsesCoreModule, } from "../helpers/responses-core-source"; -import { createResponsesSendBudget } from "../../src/server/responses/request-send-budget"; import { createRequestExecutionBudget } from "../../src/lib/request-execution-budget"; -import { createTranslatorBudget } from "../../src/lib/translator-budget"; -import type { TransientSendBudget } from "../../src/lib/upstream-retry"; +import { budgetOwner } from "../helpers/send-budget-owner"; // Existing, separately owned siblings at the extraction boundary. A new owner // cannot silently disappear from source-oracle coverage by being absent from the inventory. @@ -112,20 +110,6 @@ describe("Responses core module boundaries", () => { }); }); -function budgetOwner(sendBudget: TransientSendBudget) { - const translatorBudget = createTranslatorBudget(); - const result = createResponsesSendBudget({ - req: new Request("http://localhost/v1/responses"), - logCtx: { model: "test", provider: "test" }, - options: { translatorBudget, sendBudget }, - }); - if (result instanceof Response) { - translatorBudget.dispose(); - throw new Error("Unexpected workflow refusal without a workflow root"); - } - return { owner: result, dispose: () => translatorBudget.dispose() }; -} - describe("Responses request-owned send budget after extraction", () => { test("legacy holders retain identity and an exhausted remainder stays zero", () => { const holder = { used: 2 }; From 919e9dc4842700c03a3a678c14e12b58d84c0d7f Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Thu, 17 Sep 2026 15:37:51 +0900 Subject: [PATCH 2/3] test(layout): map physical-send.test.ts to the adapters domain --- scripts/test-layout/layout.json | 1 + 1 file changed, 1 insertion(+) diff --git a/scripts/test-layout/layout.json b/scripts/test-layout/layout.json index c56dfe6117..8d08c84538 100644 --- a/scripts/test-layout/layout.json +++ b/scripts/test-layout/layout.json @@ -188,6 +188,7 @@ "adapter-event-oauth-failover.test.ts": "oauth", "adapter-inner-send-budget-wiring.test.ts": "adapters", "adapter-inner-send-budget.test.ts": "adapters", + "physical-send.test.ts": "adapters", "adapter-registry-authority.test.ts": "adapters", "adapter-resolve.test.ts": "server", "adapter-tool-conformance.test.ts": "adapters", From 31229b26849a8d21524ad090d4cabf97d8f57939 Mon Sep 17 00:00:00 2001 From: JUN Date: Thu, 17 Sep 2026 17:51:17 +0900 Subject: [PATCH 3/3] fix(mimo-free): drain the 401 body before the JWT refresh can reject The 401 replay moved its drain behind admission so that a budget-refused replay can still return that same response with a readable body. Inside beforeDispatch it ran last, after resetMimoJwtCache and getMimoJwt. getMimoJwt issues its own bootstrap request and rejects on a failed or oversized response. When it did, fetchResponse threw and the 401 body was never released - a leak the pre-change code did not have, because it cancelled first and refreshed second. Draining first WITHIN beforeDispatch keeps both properties: it is still after admission, so a refusal returns the untouched response, and it no longer depends on the refresh succeeding. --- src/adapters/mimo-free.ts | 7 ++++-- tests/providers/mimo-free-provider.test.ts | 29 ++++++++++++++++++++++ 2 files changed, 34 insertions(+), 2 deletions(-) diff --git a/src/adapters/mimo-free.ts b/src/adapters/mimo-free.ts index 668691b2a1..55185019aa 100644 --- a/src/adapters/mimo-free.ts +++ b/src/adapters/mimo-free.ts @@ -267,14 +267,17 @@ export function createMimoFreeAdapter(provider: OcxProviderConfig): ProviderAdap try { return await send({ url: request.url, sendClass: "auth-recovery", recovery: "oauth-401", beforeDispatch: async () => { - // Drain the first response body and refresh the JWT only after admission. + // Drain the first response body and refresh the JWT only after admission: a + // refused replay still returns THIS response to the caller, body intact. + // Draining comes first within the block because getMimoJwt issues its own + // network call and may throw, and the 401 body would then never be released. + try { void response.body?.cancel().catch(() => {}); } catch { /* already consumed */ } resetMimoJwtCache(); const freshJwt = await getMimoJwt(ctx?.abortSignal); retryHeaders = { ...(request.headers as Record), "Authorization": `Bearer ${freshJwt}`, }; - try { void response.body?.cancel().catch(() => {}); } catch { /* already consumed */ } }, dispatch: executor => executor(request.url, { method: request.method, diff --git a/tests/providers/mimo-free-provider.test.ts b/tests/providers/mimo-free-provider.test.ts index 8425ada95b..e8993a92fb 100644 --- a/tests/providers/mimo-free-provider.test.ts +++ b/tests/providers/mimo-free-provider.test.ts @@ -394,6 +394,35 @@ describe("mimo-free auth retry predicate", () => { resetMimoJwtCache(); } }); + + test("a rejected JWT refresh still releases the first 401 body", async () => { + // The drain belongs inside `beforeDispatch` so that a budget-refused replay can still hand + // the 401 back with a readable body. Within that block it has to come FIRST, because + // `getMimoJwt` issues its own bootstrap request and can reject -- and the 401 body would + // then never be released. + const originalFetch = globalThis.fetch; + let cancelled = false; + globalThis.fetch = mock(async (url: string | URL | Request) => { + if (String(url).includes("/bootstrap")) { + return new Response(JSON.stringify({ jwt: "x".repeat(64 * 1024 + 1) }), { status: 200 }); + } + return new Response(new ReadableStream({ + start(controller) { controller.enqueue(new TextEncoder().encode("expired")); }, + cancel() { cancelled = true; }, + }), { status: 401 }); + }) as unknown as typeof fetch; + try { + const adapter = adapterForRetry(); + await expect(adapter.fetchResponse!( + { url: MIMO_CHAT_URL, method: "POST", headers: { "Authorization": "Bearer stale" }, body: "{}" }, + {} as never, + )).rejects.toThrow("MiMo bootstrap response too large"); + expect(cancelled).toBe(true); + } finally { + globalThis.fetch = originalFetch; + resetMimoJwtCache(); + } + }); }); describe("mimo-free adapter request building", () => {