diff --git a/tests/lib/execution-budget-permits.test.ts b/tests/lib/execution-budget-permits.test.ts index f9163860bcc..3f784e3c145 100644 --- a/tests/lib/execution-budget-permits.test.ts +++ b/tests/lib/execution-budget-permits.test.ts @@ -1,5 +1,4 @@ import { describe, expect, test } from "bun:test"; -import { readFileSync } from "node:fs"; import { CODEX_TEXT_GUARDED_BUDGET_POLICY, createRequestExecutionBudget, @@ -200,76 +199,73 @@ describe("layer caps intersect the shared budget", () => { }); /** - * The refund property above is only worth something if every caller actually uses it. + * A credential hop reserves before it knows whether account resolution, request rebuilding, or + * admission will reach the wire. The reservation is a real charge immediately, so every exit + * before dispatch must release it. Once bytes leave, the same permit must become non-refundable. * - * The generic-OAuth 429 ladder reserves a hop before it knows whether a rotation is possible. - * Two of its three exits released correctly and the `catch` did not, so a throw from the - * snapshot fetch or from credential application charged the request for a send that never left - * the process — and a later recovery in the same request was then refused on an allowance - * nothing had spent. The passthrough and runTurn ladders already had it right; these two did not. - * - * This is a source oracle because the defect lives in the caller's control flow, not in the - * budget: a unit test of the budget cannot see a caller that forgets to hand the permit back. + * These cases assert the permit state and spend observer directly. They fail on the historical + * accounting defect without depending on a particular server function name or catch-block shape. */ -describe("generic-OAuth hop reservations are handed back when no send happens", () => { - // Bounded to each ladder's own span and matched on the catch that opens it. An earlier version - // of this test searched from the first following "catch {" and found the inline body-cancel - // catch instead, so it passed while the defect was still present. - const ladder = (relativePath: string, fromMarker: string, toMarker: string): string => { - const source = readFileSync(new URL("../../" + relativePath, import.meta.url), "utf8"); - const from = source.indexOf(fromMarker); - const to = source.indexOf(toMarker, from); - expect(from).toBeGreaterThan(-1); - expect(to).toBeGreaterThan(from); - return source.slice(from, to); +describe("dispatch permits distinguish pre-send failures from physical sends", () => { + const recordingObserver = () => { + const events: string[] = []; + return { + events, + observer: { + charge: () => { events.push("charge"); return true; }, + refund: () => { events.push("refund"); }, + }, + }; }; - const refundsOnThrow = /catch \{[^}]*hop\.permit\?\.release\(\)/; - - test("the adapter dispatch ladder confirms at the dispatch boundary and refunds otherwise", () => { - const source = readFileSync(new URL("../../src/server/responses/adapter-dispatch.ts", import.meta.url), "utf8"); - // Confirming before the rebuild is not enough: buildRequest failures return { failed } - // without reaching the wire, so the hop is confirmed by the callback the rebuild invokes at - // its dispatch boundary, and the { failed } arm refunds whatever that callback did not spend. - expect(source).toContain("onDispatch?.()"); - // Confirmed at the wire, not before the pacer: waitForProviderRequestSlot can reject for an - // abort, a saturated queue, an expired slot or a removed provider without ever calling the - // adapter, and release() is a no-op once used, so an early confirm could never be refunded. - const slotWait = source.indexOf("await waitForProviderRequestSlot("); - const confirmAfterWait = source.indexOf("onDispatch?.()", slotWait); - const adapterSend = source.indexOf("transportState.activeAdapter.fetchResponse(retryRequest", confirmAfterWait); - expect(slotWait).toBeGreaterThan(-1); - expect(confirmAfterWait).toBeGreaterThan(slotWait); - expect(adapterSend).toBeGreaterThan(confirmAfterWait); - // The helper path has the same boundary inside the thunk that reaches the wire. - const thunkConfirm = source.indexOf("onDispatch?.()", adapterSend); - const headerTimeout = source.indexOf("fetchWithHeaderTimeout(retryRequest.url", thunkConfirm); - expect(thunkConfirm).toBeGreaterThan(adapterSend); - expect(headerTimeout).toBeGreaterThan(thunkConfirm); - const block = ladder( - "src/server/responses/adapter-dispatch.ts", - "adapter-recovery-oauth-429", - "attemptOpaqueBlobRecovery", - ); - expect(block).toContain('rebuildAndRefetch("oauth-account-429", () => {'); - // ...except on an adapter-owned ladder, which confirms through its own reservation. Settling - // here as well would close the permit before `adapterDispatchBudget` could hand it over, and - // an adapter whose `use()` fails reads the request as exhausted and stops sending (#4709). - expect(block).toContain("if (!adapterOwnsDispatch) hop.permit?.use();"); - expect(block).toContain("sendBudgetState.pendingHopPermit = hop.permit;"); - expect(block).toMatch(/if \("failed" in result\) \{[^}]*hop\.permit\?\.release\(\)/); - expect(block).toMatch(refundsOnThrow); + + test("a reservation released after a pre-dispatch failure books no spend", () => { + const spy = recordingObserver(); + const budget = createRequestExecutionBudget(ONE_SEND_LEFT, "lr-pre-dispatch", spy.observer); + const hop = budget.reserveDispatch({ + sendClass: "auth-recovery", + targetKey: "provider|model", + countedExternally: true, + }); + if (!hop.allowed) throw new Error("unreachable"); + + let physicalSends = 0; + try { + throw new Error("credential application failed"); + } catch { + hop.permit.release(); + } + + expect(physicalSends).toBe(0); + expect(budget.used).toBe(0); + expect(spy.events).toEqual(["charge", "refund"]); + expect(hop.permit.use()).toBe(false); + expect(budget.reserveDispatch({ sendClass: "auth-recovery", targetKey: "provider|model" }).allowed) + .toBe(true); }); - test("the continuation ladder refunds, because its send happens after the loop continues", () => { - const block = ladder( - "src/server/responses/adapter-continuation.ts", - "continuation-oauth-429", - "shouldAttemptImageTierRetry", - ); - // Nothing in that try dispatches: the replay is the next iteration, so a throw must return - // the reservation rather than confirm it. - expect(block).not.toContain("hop.permit?.use()"); - expect(block).toMatch(refundsOnThrow); + test("a reservation confirmed at dispatch stays charged after a later failure", () => { + const spy = recordingObserver(); + const budget = createRequestExecutionBudget(ONE_SEND_LEFT, "lr-post-dispatch", spy.observer); + const hop = budget.reserveDispatch({ + sendClass: "auth-recovery", + targetKey: "provider|model", + }); + if (!hop.allowed) throw new Error("unreachable"); + + let physicalSends = 0; + try { + physicalSends += 1; + expect(hop.permit.use()).toBe(true); + throw new Error("upstream rejected after dispatch"); + } catch { + hop.permit.release(); + } + + expect(physicalSends).toBe(1); + expect(budget.used).toBe(1); + expect(spy.events).toEqual(["charge"]); + expect(budget.reserveDispatch({ sendClass: "transient", targetKey: "provider|model" })) + .toEqual({ allowed: false, reason: "total-exhausted" }); }); }); @@ -327,27 +323,6 @@ describe("a credential hop is settled by whichever layer dispatches its replay", expect(budget.used).toBe(2); }); - test("the three adapter hop sites hand their reservation down instead of double-charging", () => { - const responses = (name: string): string => - readFileSync(new URL("../../src/server/responses/" + name, import.meta.url), "utf8"); - // The adapter recovery loop and the continuation loop both pick their settlement from the - // shape of the dispatcher, so neither promises an external report an adapter would never make. - for (const name of ["adapter-dispatch.ts", "adapter-continuation.ts"]) { - const source = responses(name); - expect(source).toContain("const adapterOwnsDispatch = transportState.activeAdapter.fetchResponse !== undefined;"); - expect(source).toContain("!adapterOwnsDispatch && transientRetryPolicyFor(route.provider) !== null,"); - } - // runTurn has only one shape: the adapter owns the transport, so it never reports and the - // reservation is always handed down rather than confirmed here. - const runTurn = responses("run-turn-execution.ts"); - expect(runTurn).toContain("sendBudgetState.pendingHopPermit = hop.permit;"); - expect(runTurn).not.toContain("hop.permit?.use();"); - // Every adapter-owned transport now reserves against the view, which is what spends the - // handed-down permit. Passing the bare holder is the regression this pins. - for (const name of ["adapter-dispatch.ts", "adapter-continuation.ts", "run-turn-execution.ts"]) { - expect(responses(name)).not.toContain("sendBudget: adapterSendBudget"); - } - }); }); describe("derived policy scopes", () => { diff --git a/tests/lib/transient-budget-scope-source.test.ts b/tests/lib/transient-budget-scope-source.test.ts index a26a18bfa8b..ecf41ca3613 100644 --- a/tests/lib/transient-budget-scope-source.test.ts +++ b/tests/lib/transient-budget-scope-source.test.ts @@ -1,185 +1,180 @@ -import { readResponsesCoreSource } from "../helpers/responses-core-source"; - test("the gated-model 400 ladder is charged, and keeps its own bound", () => { - const core = readResponsesCoreSource(); - // Every rung reserves and charges, so the ladder is visible to later legs instead of - // spending the request's allowance invisibly -- that part was the real defect. - expect(core).toContain("targetKey: ladderTargetKey,"); - expect(core).toContain("if (rung.allowed) rung.permit.use();"); - // A same-account replay must reserve under the SAME target key the other legs use. Folding - // the account id in made every rung read as a target change and spent the one cross-account - // slot a genuine move needs. - expect(core).toContain("const ladderTargetKey = `${route.providerName}|${route.modelId}`;"); - expect(core).not.toContain("|${retryAuthCtx.accountId}`;"); - // The ladder keeps its own bound and a budget refusal does NOT end it. #2097 pins this - // recovery at eight same-account dispatches; clamping it to what the request has left would - // cut a working path to four, which is the flat-ceiling mistake 040 warns about. - expect(core).toContain("const maxRetrySends = retrySameConfirmedAccount ? 7 : 1;"); - expect(core).not.toContain("Math.min(retrySameConfirmedAccount ? 7 : 1, sharedSendsLeft)"); - });import { describe, expect, test } from "bun:test"; -import { readFileSync } from "node:fs"; -import { join } from "node:path"; -import { repoPath } from "../helpers/repo-root"; - -const source = (relative: string): string => - readFileSync(repoPath("src", ...relative.split("/")), "utf8"); +import { describe, expect, test } from "bun:test"; +import { + createRequestExecutionBudget, + deriveRequestExecutionBudget, + type RequestExecutionBudget, + type RequestExecutionBudgetPolicy, + type RequestSendObserver, +} from "../../src/lib/request-execution-budget"; +import { + fetchWithResetRetry, + fetchWithTransientRetry, + SendBudgetExhaustedError, +} from "../../src/lib/upstream-retry"; + +const THREE_SEND_POLICY: RequestExecutionBudgetPolicy = { + maxTotalModelSends: 3, + baseSendAllowance: 3, + finalRecoveryAllowance: 0, + maxAlternateTargetSends: 1, + maxTargetTransitions: 1, +}; + +const recordingObserver = (): RequestSendObserver & { readonly charges: number; readonly refunds: number } => { + let charges = 0; + let refunds = 0; + return { + get charges() { return charges; }, + get refunds() { return refunds; }, + charge() { + charges += 1; + return true; + }, + refund() { + refunds += 1; + }, + }; +}; /** - * `transientRetryOn5xx.attempts` is ONE request-wide total-send budget, not a per-leg - * allowance. A Responses request can reach upstream on several legs — the initial send, a - * 429/account-rotation refetch, and the terminal-guard continuation — and each leg calls - * `fetchWithTransientRetry` separately. The budget only holds if every leg draws from the - * shared request-scoped counter. + * `transientRetryOn5xx.attempts` is one request-wide total-send budget, not a per-leg + * allowance. The historical failure gave a continuation or combo child a fresh allowance, + * so several individually bounded retry helpers multiplied into an unbounded request. * - * The continuation leg shipped on the raw policy value instead, so a request that reached it - * received a fresh full `attempts` allowance: with `attempts: 3` an initial send that had - * already spent its budget could still emit three more upstream sends. Runtime coverage in - * `tests/providers/upstream-transient-retry.test.ts` proves the helper reports and honors a remainder; - * it cannot prove that every call site asks for one, because a site that forgets simply - * passes a larger number. This asserts the wiring at the source, which is the only place the - * omission is visible. + * These cases drive the `src/lib` boundary directly: separate helper invocations report their + * physical sends into one request ledger, and a derived child shares that exact ledger. The + * observer is the durable-spend boundary, so its event count independently proves that one + * physical send produced one charge. */ -describe("transient send budget stays request-scoped", () => { - test("every transient-retry call site draws from the shared counter", () => { - const core = readResponsesCoreSource(); - - // One holder per LOGICAL request, read before any leg can send and inherited by combo - // children through the options spread rather than recreated per child turn. - expect(core.match(/const sendBudget = options\.sendBudget \?\? createRequestExecutionBudget\(\);/g)) - .toHaveLength(1); - // Genuine ingress mints it; a child arrives with the parent's and must not replace it. - expect(core).toContain("sendBudget: options.sendBudget ?? createRequestExecutionBudget("); - // ...and the durable spend observer is installed WITH it, for the same reason: a child that - // inherited the holder must not open a second set of ledger entries for the same sends. - expect(core).toContain("attachRequestSpendTracker(req, logCtx)"); - // The regressed shape: a counter local to one call frame, which a combo child restarts. - expect(core).not.toContain("let transientSendsUsed = 0;"); - expect(core.match(/const remainingTransientSendBudget = \(budget: number\): number =>/g)).toHaveLength(1); - // Zero has to mean zero. The Math.max(1, ...) floor funded one more send on every recovery - // leg, which is most of how a bounded per-leg allowance composed into an unbounded - // per-request count (#4546 REQ-B04). - expect(core).not.toContain("Math.max(1, budget - sendBudget.used)"); - - // Seven legs report into the same counter: the adapter initial send, the 429/rotation - // refetch, the terminal-guard continuation, and the four Codex passthrough sends (initial, - // rebuild refetch, OAuth 401 replay, rate-limit 429 replay). The passthrough four were added - // for #4546: the owner used to be declared BELOW that branch, which put it in the temporal - // dead zone there, so each of those legs silently took the helper's fresh default of 3. - expect(core.match(/onSendsConsumed: noteTransientSends/g)).toHaveLength(7); - - // EVERY leg asks for the remainder now, including the adapter initial send. That one used - // to pass the raw policy on the argument that nothing had been spent yet -- true for a first - // turn, false for a combo child, which inherits the parent's holder and then took a fresh - // full allowance on its own first send. Five sites spell it directly; the two rebuild legs - // go through recoverySendAllowance, which spends the base allowance first and only then - // draws the single shared final-recovery reserve. - expect(core.match(/attempts: remainingTransientSendBudget\(/g)).toHaveLength(5); - expect(core).toContain("attempts: remainingTransientSendBudget(transientPolicy.attempts)"); - expect(core).toContain("attempts: remainingTransientSendBudget(continuationTransientPolicy.attempts)"); - // The reserve path: an account move and a validated rebuild share ONE final send, so a - // request cannot take both and reach five. - expect(core.match(/recoverySendAllowance\(/g)).toHaveLength(2); - expect(core).toContain("countedExternally: true"); - // The passthrough legs have no adapter policy to draw from, so they name the helper's own - // ceiling rather than re-spelling the number. - expect(core).toContain("attempts: remainingTransientSendBudget(TRANSIENT_RETRY_MAX_ATTEMPTS)"); - // The trap that would make the passthrough wiring a silent no-op: transientRetryPolicyFor - // returns null for Codex forward auth, so gating these sites on it would restore a fresh 3. - expect(core).not.toContain("transientPolicy ? { attempts: remainingTransientSendBudget(TRANSIENT_RETRY_MAX_ATTEMPTS)"); - - // The regressed shape: a leg handing itself a fresh full budget. - expect(core).not.toContain("attempts: continuationTransientPolicy.attempts }"); - expect(core).not.toContain("attempts: refetchTransientPolicy.attempts }"); - expect(core).not.toContain("attempts: transientPolicy.attempts,"); +describe("transient send accounting stays request-scoped", () => { + test("separate retry legs and a derived child consume one shared allowance", async () => { + const observer = recordingObserver(); + const parent = createRequestExecutionBudget(THREE_SEND_POLICY, "lr-shared", observer); + const child = deriveRequestExecutionBudget(parent, THREE_SEND_POLICY); + let physicalSends = 0; + + const sendOne = async (budget: RequestExecutionBudget): Promise => + fetchWithTransientRetry(async () => { + physicalSends += 1; + return new Response("ok"); + }, { + attempts: budget.remainingBaseSends(3), + onSendsConsumed: (count) => { budget.used += count; }, + }); + + expect((await sendOne(parent)).status).toBe(200); + expect((await sendOne(child)).status).toBe(200); + expect((await sendOne(parent)).status).toBe(200); + await expect(sendOne(child)).rejects.toBeInstanceOf(SendBudgetExhaustedError); + + expect(physicalSends).toBe(3); + expect(parent.used).toBe(3); + expect(child.used).toBe(3); + expect(observer.charges).toBe(3); + expect(observer.refunds).toBe(0); + }); + + test("an externally counted reservation and its helper report book one send", async () => { + const observer = recordingObserver(); + const budget = createRequestExecutionBudget(THREE_SEND_POLICY, "lr-external", observer); + const reserved = budget.reserveDispatch({ + sendClass: "auth-recovery", + targetKey: "provider|model", + countedExternally: true, + }); + expect(reserved.allowed).toBe(true); + if (!reserved.allowed) throw new Error("unreachable"); + + let physicalSends = 0; + const response = await fetchWithTransientRetry(async () => { + physicalSends += 1; + return new Response("ok"); + }, { + attempts: 1, + onSendsConsumed: (count) => { budget.used += count; }, + }); + + expect(response.status).toBe(200); + expect(physicalSends).toBe(1); + expect(budget.used).toBe(1); + expect(observer.charges).toBe(1); + expect(reserved.permit.use()).toBe(true); + expect(reserved.permit.use()).toBe(false); }); - test("the helper still exposes the seam those call sites depend on", () => { - const retry = source("lib/upstream-retry.ts"); - expect(retry).toContain("onSendsConsumed?: (sends: number) => void;"); - // Reported in `finally` so every exit path — return, throw, abort — feeds the counter. - expect(retry).toMatch(/} finally \{\n\s*opts\.onSendsConsumed\?\.\(sent\);/); - // A spent budget must refuse rather than round itself up to one more send. - expect(retry).not.toContain("Math.max(1, opts.attempts ?? RESET_RETRY_MAX_ATTEMPTS)"); - expect(retry).not.toContain("Math.max(1, opts.attempts ?? TRANSIENT_RETRY_MAX_ATTEMPTS)"); - expect(retry).not.toContain("Math.max(1, budget - sent)"); - expect(retry).toContain("class SendBudgetExhaustedError extends Error"); + test("the transient wrapper reports a rejected physical send exactly once", async () => { + const observer = recordingObserver(); + const budget = createRequestExecutionBudget(THREE_SEND_POLICY, "lr-rejected", observer); + const reports: number[] = []; + let physicalSends = 0; + const rejection = new Error("transport failed before a response"); + + await expect(fetchWithTransientRetry(async () => { + physicalSends += 1; + throw rejection; + }, { + attempts: budget.remainingBaseSends(3), + onSendsConsumed: (count) => { + reports.push(count); + budget.used += count; + }, + })).rejects.toBe(rejection); + + expect(physicalSends).toBe(1); + expect(reports).toEqual([1]); + expect(budget.used).toBe(1); + expect(observer.charges).toBe(1); }); }); /** - * The dispatch paths that were not merely uncounted but UNCOUNTABLE (#4546). - * - * Three holes survived the earlier slices, and each is invisible at runtime until a real account - * pool is hot: `fetchWithResetRetry` had no reporting seam at all, so every leg without a - * transient policy sent off the books; the compact endpoint's routed fallback called - * `handleResponses` with no budget, so a native attempt's spend was forgotten the moment it fell - * through; and the credential hops enforced their own per-roster caps against a counter that knew - * nothing about the rest of the request. The wiring is what these assert -- the arithmetic is - * pinned in `request-execution-budget.test.ts`. + * The reset helper is also a physical-send owner. It must report before awaiting the transport, + * while the transient wrapper must suppress that inner report because its counted fetch already + * owns the same send. Otherwise a rejected send disappears, or a successful send is charged twice. */ -describe("every dispatch path reports into the shared budget", () => { - test("the reset-only helper counts its own physical sends", () => { - const retry = source("lib/upstream-retry.ts"); - // The seam moved onto ResetRetryOptions. On TransientRetryOptions it could not be reached by - // the non-policy adapter send or by any rebuildAndRefetch leg with a null transient policy. - const resetOptions = retry.slice( - retry.indexOf("export interface ResetRetryOptions {"), - retry.indexOf("export interface TransientRetryOptions"), - ); - expect(resetOptions).toContain("onSendsConsumed?: (sends: number) => void;"); - // One report per physical send, before the await, so a rejected send still counts. - expect(retry).toContain("opts.onSendsConsumed?.(1);"); - // ...and the transient layer, which already counts the same sends through countedFetch, - // suppresses the inner reporter. Forwarding it would count every inner send twice. - expect(retry).toContain("onSendsConsumed: undefined,"); - expect(retry).not.toContain("fetchWithResetRetry(countedFetch, { ...opts, attempts: remaining() })"); - }); +describe("retry helpers expose one accounting event per physical send", () => { + test("the reset-only helper reports a rejected send", async () => { + const reports: number[] = []; + let physicalSends = 0; + const rejection = new Error("connection refused"); + + await expect(fetchWithResetRetry(async () => { + physicalSends += 1; + throw rejection; + }, { + attempts: 1, + onSendsConsumed: (count) => reports.push(count), + })).rejects.toBe(rejection); - test("compact holds ONE budget for the native attempt, the handoff child and the routed turn", () => { - const compact = source("server/responses/compact.ts"); - // Declared once, at function scope. Inside the native branch it was out of reach of the - // routed fallback below, which is reached by a 404 native compact and by a quota failure. - expect(compact.match(/const sendBudget: RequestExecutionBudget = options\.sendBudget \?\? createRequestExecutionBudget\(\);/g)) - .toHaveLength(1); - // The routed compaction turn inherits it instead of letting handleResponsesInner mint a - // fresh four. - expect(compact).toContain("turnAdmissionLease, sendBudget,"); - // The handoff child already inherited; both paths must keep doing so. - expect(compact).toContain("{ ...options, sendBudget }"); + expect(physicalSends).toBe(1); + expect(reports).toEqual([1]); }); - test("credential hops keep their roster cap AND reserve from the shared budget", () => { - const core = readResponsesCoreSource(); - // Six hop sites: the native passthrough 429, the shared sidecar hook's generic and - // Anthropic arms, the runTurn preflight 429, the adapter recovery loop, and the - // continuation loop. The last two were the arms that actually iterate the roster, so - // leaving them out meant the claim held everywhere except where it mattered most. - expect(core.match(/reserveCredentialHop\(/g)).toHaveLength(6); - // The per-roster caps are NOT replaced. The effective allowance is the intersection, so - // removing either half is a behaviour change that has to be argued for. - expect(core).toContain("genericFailovers < GENERIC_OAUTH_MAX_FAILOVERS_PER_REQUEST"); - expect(core).toContain("genericFailovers >= GENERIC_OAUTH_MAX_FAILOVERS_PER_REQUEST"); - expect(core).toContain("anthropicPoolFailovers < ANTHROPIC_POOL_MAX_FAILOVERS_PER_REQUEST"); - // A refused hop hands the reservation back rather than spending a send it never made. - expect(core.match(/hop\.permit\?\.release\(\);/g)?.length ?? 0).toBeGreaterThanOrEqual(6); - // The passthrough hop's replay spends the hop's own reservation; a second one would be - // refused as final-recovery-spent and would answer 502 instead of the real 429. - expect(core).toContain("pendingHopPermit = hop.permit;"); + test("the transient wrapper suppresses the nested reset report", async () => { + const reports: number[] = []; + let physicalSends = 0; + + const response = await fetchWithTransientRetry(async () => { + physicalSends += 1; + return new Response("ok"); + }, { + attempts: 1, + onSendsConsumed: (count) => reports.push(count), + }); + + expect(response.status).toBe(200); + expect(physicalSends).toBe(1); + expect(reports).toEqual([1]); }); - test("the gated-model 400 ladder is charged, and keeps its own bound", () => { - const core = readResponsesCoreSource(); - // Every rung reserves and charges, so the ladder is visible to later legs instead of - // spending the request's allowance invisibly -- that was the real defect. - expect(core).toContain("targetKey: ladderTargetKey,"); - expect(core).toContain("if (rung.allowed) rung.permit.use();"); - // A same-account replay reserves under the SAME target key the other legs use. Folding the - // account id in made every rung read as a target change and spent the one cross-account slot - // a genuine move needs. - expect(core).toContain("const ladderTargetKey = `${route.providerName}|${route.modelId}`;"); - // The ladder keeps its own bound and a budget refusal does NOT end it. #2097 pins this - // recovery at eight same-account dispatches; clamping it to what the request has left cut a - // working path to four, which is the flat-ceiling mistake 040_send_budget.md warns about. - expect(core).toContain("const maxRetrySends = retrySameConfirmedAccount ? 7 : 1;"); - expect(core).not.toContain("Math.min(retrySameConfirmedAccount ? 7 : 1, sharedSendsLeft)"); + test("zero remaining sends refuses both helpers before dispatch", async () => { + for (const send of [fetchWithResetRetry, fetchWithTransientRetry]) { + let physicalSends = 0; + await expect(send(async () => { + physicalSends += 1; + return new Response("must not send"); + }, { attempts: 0 })).rejects.toBeInstanceOf(SendBudgetExhaustedError); + expect(physicalSends).toBe(0); + } }); }); diff --git a/tests/responses/responses-preview-main-read-fence.test.ts b/tests/responses/responses-preview-main-read-fence.test.ts index 5242c0dbc22..cee53c1dd6a 100644 --- a/tests/responses/responses-preview-main-read-fence.test.ts +++ b/tests/responses/responses-preview-main-read-fence.test.ts @@ -1,97 +1,435 @@ -import { describe, expect, test } from "bun:test"; -import { repoPath } from "../helpers/repo-root"; +import { afterEach, beforeEach, describe, expect, spyOn, test } from "bun:test"; +import * as fs from "node:fs"; +import { mkdtempSync, unlinkSync, writeFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { saveCodexAccountCredential } from "../../src/codex/account-store"; +import { + resolveCodexAuthContext, + type CodexAuthContext, +} from "../../src/codex/auth-context"; +import { getMainAccountToken, MAIN_CODEX_ACCOUNT_ID } from "../../src/codex/main-account"; +import { + resetCodexModelEntitlementCacheForTests, + seedCodexModelEntitlementsForTests, +} from "../../src/codex/model-entitlements"; +import { + blockNativeMainRecovery, + completeNativeMainRecovery, + nativeMainStartupGateSnapshot, +} from "../../src/codex/native-profile-startup"; +import { clearAccountQuota } from "../../src/codex/quota"; +import { clearCodexUpstreamHealth, clearThreadAccountMap } from "../../src/codex/routing"; +import { + noteSubagentModelFailure, + resetSubagentModelFallbackStateForTests, +} from "../../src/codex/subagent-model-fallback"; +import { clearComboTargetCooldowns } from "../../src/combos/failover"; +import { clearComboSelectionState } from "../../src/combos/resolve"; +import { + clearResponseStateForTests, + clearResponseStateMemoryForTests, +} from "../../src/responses/state"; +import { handleResponses } from "../../src/server/responses"; +import { resetAgentTaskRecoveryState } from "../../src/server/responses/agent-task-recovery"; +import type { ActiveTurnLease } from "../../src/server/lifecycle"; +import type { RequestLogContext } from "../../src/server/request-log"; +import type { OcxConfig } from "../../src/types"; +import { + codexHeaders, + encryptedInput, + recoverySse, +} from "../helpers/agent-task-recovery"; +import { removeTreeWithRetry } from "../helpers/remove-tree"; /** - * Request preview exists to predict what final authentication will decide, so the two must apply - * the same native-main read fence. Final auth forbids those reads for three reasons and the first - * of them is ownership: a request that authenticates with the CALLER's own credential may not - * read, reconcile or score the physical main token (`resolveCodexAuthContext`). Preview computed - * the same-named constant from recovery and drain state only, so a `thread_spawn` carrying a - * forwardable caller bearer previewed with main included -- a fence violation and a - * preview/final disagreement at once. + * Request preview predicts final authentication, so both decisions must fence the physical native + * main credential on the same three facts: caller ownership, retained recovery, and a draining + * selector. The historical ownership omission was especially dangerous: a `thread_spawn` with a + * forwardable caller bearer let preview open `auth.json` while final authentication correctly + * treated that file as belonging to a different credential domain. * - * Asserted on the source, like the sibling preview-site contract in - * `tests/routing/subagent-fallback-preview-sites.test.ts`. Driving it end to end needs a - * thread_spawn whose caller bearer is forwardable, an account-gated candidate model, and a - * populated denial cache whose only entry is main; the fixture that arrangement demands is more - * fragile than the divergence it would catch. What this does catch is the regression that - * actually threatens the fix -- one of the two preview fences being reconstructed from drain - * state alone again, which is how the recovery path came to repeat the omission. + * The denial cache is the behavioral oracle here. A cached native-main denial validates the + * physical token before it can influence account scoring; excluding main skips that validation + * before the file is opened. The fence cases therefore reach the real preview/final path and + * observe `auth.json` reads rather than the spelling of the fence expression. */ -describe("preview and final agree on the native-main read fence (source contract)", () => { - const requestPrepareSource = async (): Promise => - Bun.file(repoPath("src", "server", "responses", "request-prepare.ts")).text(); - const authContextSource = async (): Promise => - Bun.file(repoPath("src", "codex", "auth-context.ts")).text(); - - const fenceExpression = (source: string): string => { - const match = source.match(/const nativeMainReadsForbidden =([\s\S]*?);\n/); - if (!match) throw new Error("no nativeMainReadsForbidden declaration found"); - return match[1]!; - }; - - test("final authentication still ORs request-owned ownership into its fence", async () => { - // The thing preview is copying. If final auth ever stops fencing on ownership, the copy below - // is no longer parity and this file should be revisited rather than quietly kept. - const source = await authContextSource(); - - expect(fenceExpression(source)).toContain("requestScopedMainCredential"); - // And it validates the caller's option against the header it will actually send. - expect(source).toMatch(/options\.requestScopedMainCredential === true\s*\n?\s*&& hasCallerCodexBearer\(headers\)/); + +const NOW = 1_800_000_000_000; +const PREFERRED_MODEL = "gpt-5.6-sol"; +const FALLBACK_MODEL = "xai/grok-4.5"; +const originalFetch = globalThis.fetch; +const originalNow = Date.now; + +let testDir = ""; +let previousOpenCodexHome: string | undefined; +let previousCodexHome: string | undefined; +let authJsonReads = 0; +let authJsonReadStacks: string[] = []; +let readSpy: ReturnType | undefined; +let blockedHomeId: string | null = null; + +function providerConfig(overrides: Partial = {}): OcxConfig { + return { + port: 0, + defaultProvider: "openai", + activeCodexAccountId: "pool-a", + autoSwitchThreshold: 0, + subagentModelFallback: [FALLBACK_MODEL], + providers: { + openai: { + adapter: "openai-responses", + baseUrl: "https://chatgpt.com/backend-api/codex", + authMode: "forward", + codexAccountMode: "pool", + }, + xai: { + adapter: "openai-chat", + baseUrl: "https://api.x.ai/v1", + authMode: "key", + apiKey: "xai-test", + }, + }, + codexAccounts: [ + { id: MAIN_CODEX_ACCOUNT_ID, email: "main@example.test", isMain: true }, + { id: "pool-a", email: "pool@example.test", isMain: false, chatgptAccountId: "pool-account" }, + ], + ...overrides, + } as OcxConfig; +} + +function installCredentials(): void { + writeFileSync(join(testDir, "auth.json"), JSON.stringify({ + tokens: { + access_token: "physical-main-token", + refresh_token: "physical-main-refresh", + account_id: "physical-main-account", + }, + })); + saveCodexAccountCredential("pool-a", { + accessToken: "pool-access-token", + refreshToken: "pool-refresh-token", + expiresAt: NOW + 24 * 60 * 60_000, + chatgptAccountId: "pool-account", }); +} + +function seedMainDenial(): void { + seedCodexModelEntitlementsForTests( + MAIN_CODEX_ACCOUNT_ID, + [], + NOW, + "0.146.0", + "main:physical-main-account", + ); +} + +function calibrateMainReadCounter(): void { + expect(getMainAccountToken()).toEqual({ + accessToken: "physical-main-token", + chatgptAccountId: "physical-main-account", + }); + expect(authJsonReads).toBeGreaterThan(0); + resetMainReadObservations(); +} + +function resetMainReadObservations(): void { + authJsonReads = 0; + authJsonReadStacks = []; +} + +/** + * The request also asks ordinary pool selection whether native main is live: once for the direct + * preview and once when fallback invokes the preview callback. Those reads travel through + * `isMainAccountCredentialUsable`, not the denial-cache credential validator this conversion + * protects. Keep every stack for diagnostics, then select the exact observable whose exclusion + * would regress if either request-prepare fence dropped ownership. + */ +function denialCacheMainReadStacks(): string[] { + return authJsonReadStacks.filter(stack => + stack.replaceAll("\\", "/").includes("/src/codex/model-entitlements.ts")); +} + +function completedResponses(model = PREFERRED_MODEL): Response { + return Response.json({ + id: "resp_main_read_fence", + object: "response", + status: "completed", + model, + output: [], + usage: { input_tokens: 1, output_tokens: 1, total_tokens: 2 }, + }); +} + +function readableInput(): unknown[] { + return [{ + type: "message", + role: "user", + content: [{ type: "input_text", text: "keep the main credential fenced" }], + }]; +} + +async function postSpawn( + config: OcxConfig, + options: Parameters[3] = {}, + headers: HeadersInit = codexHeaders("caller-account"), + input: unknown[] = readableInput(), + model = PREFERRED_MODEL, + logCtx: RequestLogContext = { model: "", provider: "" }, +): Promise { + const requestHeaders = new Headers(headers); + requestHeaders.set("content-type", "application/json"); + requestHeaders.set("x-openai-subagent", "collab_spawn"); + const response = await handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", + headers: requestHeaders, + body: JSON.stringify({ model, input, stream: false }), + }), config, logCtx, options); + // handleResponses owns its translator budget through the returned body lifecycle. Draining the + // body also lets completed Responses schedule their state write before afterEach cancels it. + await response.arrayBuffer(); + return response; +} - test("the preview fence carries the same ownership term", async () => { - const source = await requestPrepareSource(); - const fence = fenceExpression(source); +beforeEach(() => { + testDir = mkdtempSync(join(tmpdir(), "ocx-preview-main-fence-")); + previousOpenCodexHome = process.env.OPENCODEX_HOME; + previousCodexHome = process.env.CODEX_HOME; + process.env.OPENCODEX_HOME = testDir; + process.env.CODEX_HOME = testDir; + Date.now = () => NOW; + clearThreadAccountMap(); + clearCodexUpstreamHealth(); + clearAccountQuota(); + clearComboSelectionState(); + clearComboTargetCooldowns(); + clearResponseStateMemoryForTests(); + resetSubagentModelFallbackStateForTests(); + resetCodexModelEntitlementCacheForTests(); + resetAgentTaskRecoveryState(); + installCredentials(); - expect(fence).toContain("previewRequestScopedMainCredential"); - // Still the other two inputs as well -- adding ownership must not have replaced them. - expect(fence).toContain("nativeMainRecoveryBlocked"); - expect(fence).toContain("mainProfileDraining"); + const originalReadFileSync = fs.readFileSync as (...args: unknown[]) => unknown; + readSpy = spyOn(fs, "readFileSync"); + readSpy.mockImplementation(((...args: unknown[]) => { + const target = args[0]; + if (typeof target === "string" && target.endsWith("auth.json")) { + authJsonReads += 1; + authJsonReadStacks.push(new Error("auth.json read").stack ?? "stack unavailable"); + } + return originalReadFileSync(...args); + }) as unknown as typeof fs.readFileSync); + resetMainReadObservations(); + blockedHomeId = null; +}); + +afterEach(() => { + globalThis.fetch = originalFetch; + Date.now = originalNow; + readSpy?.mockRestore(); + readSpy = undefined; + if (blockedHomeId !== null) completeNativeMainRecovery(blockedHomeId); + blockedHomeId = null; + clearThreadAccountMap(); + clearCodexUpstreamHealth(); + clearAccountQuota(); + clearComboSelectionState(); + clearComboTargetCooldowns(); + clearResponseStateForTests(); + resetSubagentModelFallbackStateForTests(); + resetCodexModelEntitlementCacheForTests(); + resetAgentTaskRecoveryState(); + if (previousOpenCodexHome === undefined) delete process.env.OPENCODEX_HOME; + else process.env.OPENCODEX_HOME = previousOpenCodexHome; + if (previousCodexHome === undefined) delete process.env.CODEX_HOME; + else process.env.CODEX_HOME = previousCodexHome; + removeTreeWithRetry(testDir); + testDir = ""; + resetMainReadObservations(); +}); + +describe("preview and final authentication agree on the native-main read fence", () => { + test("final authentication validates request ownership against the bearer before fencing main", async () => { + seedMainDenial(); + calibrateMainReadCounter(); + const config = providerConfig(); + + const owned = await resolveCodexAuthContext(codexHeaders("caller-account"), config, "pool", { + requestScopedMainCredential: true, + modelId: PREFERRED_MODEL, + }); + expect(owned).toMatchObject({ kind: "pool", accountId: "pool-a" }); + expect(authJsonReads).toBe(0); + + // The option is only a claim by the caller. Without the bearer final auth must reject that + // claim, leave the ownership fence open, and perform the ordinary physical-main reads. + await resolveCodexAuthContext(new Headers(), config, "pool", { + requestScopedMainCredential: true, + modelId: PREFERRED_MODEL, + }); + expect(authJsonReads).toBeGreaterThan(0); }); - test("both preview sites derive ownership the way final auth validates it", async () => { - const source = await requestPrepareSource(); + test("the initial preview excludes physical main from denial-cache credential validation", async () => { + seedMainDenial(); + calibrateMainReadCounter(); + const upstreamAuth: Array = []; + globalThis.fetch = (async (_input, init) => { + upstreamAuth.push(new Headers(init?.headers).get("authorization")); + return completedResponses(); + }) as typeof fetch; - // The initial preview and the encrypted-recovery re-preview. Recovery recomputes rather than - // reusing, because a subagent fallback above it may have re-routed and ownership is a - // function of the route as well as the headers. - const validated = [...source.matchAll( - /\)\.requestScopedMainCredential\s*&&\s*hasCallerCodexBearer\(/g, - )]; + const response = await postSpawn(providerConfig()); - expect(validated).toHaveLength(2); + expect(response.status).toBe(200); + expect(upstreamAuth).toEqual(["Bearer pool-access-token"]); + expect(denialCacheMainReadStacks()).toEqual([]); }); - test("no main exclusion is guarded by drain state alone", async () => { - const source = await requestPrepareSource(); + test("the initial preview also fences main for recovery blocking and selector drain", async () => { + seedMainDenial(); + calibrateMainReadCounter(); + globalThis.fetch = (async () => completedResponses()) as typeof fetch; - // Every place preview withholds main from a credential-validating read. Each must be guarded - // either by the shared fence above -- which the previous case pins to ownership -- or by its - // own ownership term. The recovery site reconstructed this condition inline and lost the - // ownership half; that is the regression this asserts against. - const guards = [...source.matchAll( - /excludeAccountIds:\s*([\s\S]*?)\?\s*new Set\(\[MAIN_CODEX_ACCOUNT_ID\]\)/g, - )].map(match => match[1]!); + const snapshot = nativeMainStartupGateSnapshot(); + blockedHomeId = snapshot.homeId ?? testDir; + expect(blockNativeMainRecovery(blockedHomeId)).toBe(true); + const blockedResponse = await postSpawn(providerConfig(), {}, new Headers()); + expect(blockedResponse.status).toBe(200); + expect(authJsonReads).toBe(0); + expect(completeNativeMainRecovery(blockedHomeId)).toBe(true); + blockedHomeId = null; - expect(guards.length).toBeGreaterThanOrEqual(2); - const unfenced = guards.filter( - guard => !/nativeMainReadsForbidden|RequestScopedMainCredential/.test(guard), + let selectionStarts = 0; + const turnAdmissionLease = { + release() {}, + beginCodexAccountSelection() { + selectionStarts += 1; + return { + mainProfileDraining: true, + claimMainProfile: () => false, + release() {}, + }; + }, + } satisfies Pick; + resetMainReadObservations(); + await postSpawn(providerConfig(), { turnAdmissionLease }, new Headers()); + expect(selectionStarts).toBeGreaterThan(0); + expect(authJsonReads).toBe(0); + }); + + /** + * A bare preferred native model can reach the request-prepare recovery re-preview only after + * routing has already chosen a noncanonical provider, but bare native ids are reserved to the + * canonical OpenAI provider. Combo recovery is the reachable response-driven boundary: the + * canonical target rejects, recovery decrypts once, and only then may the routed target run. + */ + test("response-triggered encrypted recovery replays plaintext without post-recovery main reads", async () => { + calibrateMainReadCounter(); + const config = providerConfig({ + defaultProvider: "openai", + agentTaskRecovery: { enabled: true }, + providers: { + openai: { + adapter: "openai-responses", + baseUrl: "https://chatgpt.com/backend-api/codex", + authMode: "forward", + codexAccountMode: "direct", + }, + backup: { + adapter: "openai-responses", + baseUrl: "https://backup.example/v1", + authMode: "key", + apiKey: "backup-test-key", + }, + }, + combos: { + recovery: { + strategy: "failover", + targets: [ + { provider: "openai", model: PREFERRED_MODEL }, + { provider: "backup", model: "m2" }, + ], + }, + }, + }); + let recoveryCalls = 0; + const backupBodies: string[] = []; + globalThis.fetch = (async (input, init) => { + const body = typeof init?.body === "string" ? init.body : ""; + if (body.includes("capture_assignment")) { + recoveryCalls += 1; + // The canonical target has already failed. Reads after this response belong only to the + // recovered routed replay, so the observation cannot be satisfied by pre-recovery work. + resetMainReadObservations(); + return new Response(recoverySse("Use the recovered assignment."), { + status: 200, + headers: { "content-type": "text/event-stream" }, + }); + } + const url = new URL(input instanceof Request ? input.url : String(input)); + if (url.hostname === "backup.example") { + backupBodies.push(body); + return completedResponses("m2"); + } + // The canonical target is the only target allowed to receive unreadable ciphertext. Its + // response forces the combo owner to recover the assignment before backup becomes eligible. + return Response.json({ error: { message: "caller credential rejected" } }, { status: 401 }); + }) as typeof fetch; + + const response = await postSpawn( + config, + {}, + codexHeaders("caller-account"), + encryptedInput(), + "combo/recovery", ); - expect(unfenced).toEqual([]); + + expect(response.status).toBe(200); + expect(recoveryCalls).toBe(1); + expect(backupBodies).toHaveLength(1); + expect(backupBodies[0]).toContain("Use the recovered assignment."); + expect(authJsonReads).toBe(0); }); - test("selection-only stays derived from the drain alone, in both files", async () => { - // The asymmetry is deliberate: final auth derives `nativeMainSelectionOnly` from the drain - // without ownership, so adding an ownership term to the preview copy would diverge from it in - // the other direction. Pinned so the symmetry above is not "fixed" onto this one too. - for (const source of [await requestPrepareSource(), await authContextSource()]) { - const derivations = [...source.matchAll( - /nativeMainSelectionOnly\s*[:=]([\s\S]*?)mainProfileDraining === true/g, - )].map(match => match[1]!); - - expect(derivations.length).toBeGreaterThanOrEqual(1); - expect(derivations.filter(d => /equestScopedMainCredential/.test(d))).toEqual([]); - } + test("ownership alone leaves selection-only off in preview and final authentication", async () => { + seedMainDenial(); + calibrateMainReadCounter(); + // Make the two selection modes observably different. Ordinary selection finds physical main + // unreadable and uses pool-a; selection-only would retain main as a synthetic candidate + // without opening the missing file. + unlinkSync(join(testDir, "auth.json")); + const config = providerConfig({ activeCodexAccountId: MAIN_CODEX_ACCOUNT_ID }); + // If preview incorrectly treated ownership as selection-only, it would score the request as + // native main, observe this account-scoped failure, and route to the XAI fallback. + noteSubagentModelFailure(PREFERRED_MODEL, "429", config, MAIN_CODEX_ACCOUNT_ID, NOW); + const upstreamUrls: string[] = []; + const upstreamBodies: string[] = []; + const upstreamAuth: Array = []; + globalThis.fetch = (async (input, init) => { + upstreamUrls.push(String(input)); + upstreamBodies.push(typeof init?.body === "string" ? init.body : ""); + upstreamAuth.push(new Headers(init?.headers).get("authorization")); + return completedResponses(); + }) as typeof fetch; + let finalAuth: CodexAuthContext | undefined; + const logCtx: RequestLogContext = { model: "", provider: "" }; + + const response = await postSpawn( + config, + { onCodexAuthContextResolved: context => { finalAuth = context; } }, + codexHeaders("caller-account"), + readableInput(), + PREFERRED_MODEL, + logCtx, + ); + + expect(response.status).toBe(200); + expect(finalAuth).toMatchObject({ kind: "pool", accountId: "pool-a" }); + expect(upstreamUrls).toHaveLength(1); + expect(upstreamUrls[0]).toContain("chatgpt.com/backend-api/codex"); + expect(upstreamBodies[0]).toContain(`"model":"${PREFERRED_MODEL}"`); + expect(upstreamAuth).toEqual(["Bearer pool-access-token"]); + expect((logCtx as unknown as Record).subagentModelFallbackTo).toBeUndefined(); }); }); diff --git a/tests/routing/probe-lease-dispatch-wiring.test.ts b/tests/routing/probe-lease-dispatch-wiring.test.ts index 9ac2157106e..00c4a6b942e 100644 --- a/tests/routing/probe-lease-dispatch-wiring.test.ts +++ b/tests/routing/probe-lease-dispatch-wiring.test.ts @@ -1,11 +1,12 @@ import { afterEach, beforeEach, describe, expect, test } from "bun:test"; -import { existsSync, mkdirSync, readFileSync } from "node:fs"; +import { existsSync, mkdirSync } from "node:fs"; import { join } from "node:path"; import { canAcquireTransientProbe, classifyPoolRecoveryDispatch, clearPoolRecoveryState, createPoolBackpressureLimiter, + sharedPoolBackpressure, transientProbeDiagnostics, tryAcquireTransientProbe, TRANSIENT_PROBE_INTERVAL_MS, @@ -27,9 +28,9 @@ import { import { clearPoolRotationState } from "../../src/codex/pool-rotation"; import { saveCodexAccountCredential } from "../../src/codex/account-store"; import { clearAccountQuota, updateAccountQuota } from "../../src/codex/auth-api"; -import { repoPath } from "../helpers/repo-root"; -import { removeTreeWithRetry } from "../helpers/remove-tree"; +import { handleResponses } from "../../src/server/responses"; import type { OcxConfig } from "../../src/types"; +import { removeTreeWithRetry } from "../helpers/remove-tree"; /** * The pool-wide recovery limiter, wired to the dispatch that actually sends (#4701). @@ -38,7 +39,8 @@ import type { OcxConfig } from "../../src/types"; * bounded nothing: no file under `src/` imported the module, so every hit for * `resolveHeldAccountDispatch` was its own definition or a direct unit test. An implementation * nothing calls is indistinguishable from an absent one at runtime, which is the whole of the - * issue -- and it is why the first case here is a source oracle rather than a behaviour. + * issue. The first case therefore drives the public Responses handler and observes the shared + * limiter's demand counter at the physical-send boundary. * * The defect that reached production lived at the end of both transient-hold branches of * `resolveCodexAccountForThreadDetailed`: when no sibling could take the request they returned @@ -80,33 +82,41 @@ function streakTransientFailures(config: OcxConfig, accountId: string, now: numb } describe("recovery limiter wiring is reachable from production (#4701)", () => { - test("the transient-hold module is imported by the selector and the dispatch boundary", () => { - // Not a style assertion. Before this change the module had complete unit coverage and zero - // production callers, so the suite was green while nothing in a running proxy was bounded. - // If a refactor ever detaches it again, that is the symptom to catch -- the behaviour tests - // below would keep passing against primitives nobody calls. - const holdDispatch = readFileSync( - repoPath("src", "codex", "routing", "transient-hold-dispatch.ts"), "utf8", - ); - expect(holdDispatch).toContain('from "../../routing/probe-lease"'); - expect(holdDispatch).toContain("resolveHeldAccountDispatch"); - - // The selector reaches the bound through that seam, on the production path. - const routing = readFileSync(repoPath("src", "codex", "routing.ts"), "utf8"); - expect(routing).toContain('from "./routing/transient-hold-dispatch"'); - expect(routing).toContain("resolveTransientHoldDispatch"); - - // The physical-send boundary itself owns no routing policy -- `responses-fetch-helpers- - // boundary.test.ts` pins its runtime imports to three transport modules -- so the dispatch - // call sites name their own class instead. - const passthrough = readFileSync(repoPath("src", "server", "responses", "passthrough-dispatch.ts"), "utf8"); - expect(passthrough).toContain('from "../../routing/probe-lease"'); - expect(passthrough).toContain('classifyPoolRecoveryDispatch("initial")'); - - // Two modules are named probe-lease, one directory apart, and they are different domains. - // The selector keeps importing the QUOTA one; merging them would make one settle the - // other's probe. - expect(routing).toContain('from "./routing/probe-lease"'); + test("a production Responses dispatch records demand in the shared limiter", async () => { + const originalFetch = globalThis.fetch; + globalThis.fetch = (async () => Response.json({ + id: "resp-probe-lease-wiring", + object: "response", + status: "completed", + model: "fixture-model", + output: [], + usage: { input_tokens: 1, output_tokens: 1, total_tokens: 2 }, + })) as typeof fetch; + + const config = { + defaultProvider: "fixture", + providers: { + fixture: { + adapter: "openai-responses", + baseUrl: "https://fixture.example.test/v1", + apiKey: "sk-test", + }, + }, + } as OcxConfig; + + try { + const response = await handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ model: "fixture/fixture-model", input: "hello", stream: false }), + }), config, { model: "", provider: "" }); + await response.text(); + + expect(response.status).toBe(200); + expect(sharedPoolBackpressure().state().initialSends).toBe(1); + } finally { + globalThis.fetch = originalFetch; + } }); }); diff --git a/tests/server/cancel-body-on-abort.test.ts b/tests/server/cancel-body-on-abort.test.ts index decf0ebcd21..075fb80c084 100644 --- a/tests/server/cancel-body-on-abort.test.ts +++ b/tests/server/cancel-body-on-abort.test.ts @@ -1,7 +1,8 @@ -import { readResponsesCoreSource } from "../helpers/responses-core-source"; import { describe, expect, test } from "bun:test"; import { cancelBodyOnAbort } from "../../src/lib/abort"; -import { readBodyCapped } from "../../src/server/live"; +import { handleLive, readBodyCapped } from "../../src/server/live"; +import { handleResponses } from "../../src/server/responses"; +import type { OcxConfig } from "../../src/types"; function bodyWithCancelSpy(): { body: ReadableStream; cancelled: () => boolean } { let cancelled = false; @@ -65,56 +66,220 @@ describe("readBodyCapped settles the stream when a read throws", () => { expect(cancelled).toBe(true); }); - // Wiring guard. The unit tests above exercise readBodyCapped and cancelBodyOnAbort - // directly, which means they ALL still pass when the /v1/live relay forgets to call the - // guard — an earlier revision of this change imported the helper and never invoked it, and - // no test noticed. Asserting the call site is crude but it is the thing that was actually - // missing. - test("the live relay attaches the body guard before consuming the upstream body", async () => { - const source = await Bun.file(new URL("../../src/server/live.ts", import.meta.url)).text(); - - const guardAt = source.indexOf("cancelBodyOnAbort(upstreamResponse.body"); - const readAt = source.indexOf("payload = await readBodyCapped("); - expect(guardAt).toBeGreaterThan(-1); - expect(readAt).toBeGreaterThan(-1); - // Guard first, read second. - expect(guardAt).toBeLessThan(readAt); - // And detached on the normal path. - expect(source).toContain("detachBodyGuard()"); + test("an aborted live relay reaches fetch before settling its locked upstream body", async () => { + const originalFetch = globalThis.fetch; + const requestAbort = new AbortController(); + const events: string[] = []; + let rejectRead!: (reason: unknown) => void; + let markReadStarted!: () => void; + const readStarted = new Promise(resolve => { markReadStarted = resolve; }); + + const reader = { + read(): Promise> { + events.push("reader.read"); + markReadStarted(); + return new Promise((_resolve, reject) => { rejectRead = reject; }); + }, + cancel(): Promise { + events.push("reader.cancel"); + return Promise.resolve(); + }, + releaseLock(): void { + events.push("reader.releaseLock"); + }, + } as ReadableStreamDefaultReader; + const body = { + getReader(): ReadableStreamDefaultReader { + events.push("body.getReader"); + return reader; + }, + cancel(): Promise { + events.push("body.cancel"); + return Promise.reject(new TypeError("body is locked")); + }, + } as ReadableStream; + + globalThis.fetch = (async (_input, init) => { + const signal = init?.signal; + if (!(signal instanceof AbortSignal)) throw new Error("live relay omitted its upstream abort signal"); + signal.addEventListener("abort", () => { + events.push("fetch.abort"); + rejectRead(signal.reason); + }, { once: true }); + return { + status: 201, + headers: new Headers({ "content-type": "application/sdp" }), + body, + } as Response; + }) as typeof fetch; + + const config = { + defaultProvider: "openai-apikey", + providers: { + "openai-apikey": { + adapter: "openai-responses", + baseUrl: "https://api.openai.com/v1", + apiKey: "sk-test-live", + }, + }, + } as OcxConfig; + + try { + const pending = handleLive(new Request("http://localhost/v1/live", { + method: "POST", + headers: { "content-type": "application/sdp" }, + body: "offer", + signal: requestAbort.signal, + }), config, { model: "", provider: "" }); + await readStarted; + requestAbort.abort(new DOMException("client closed request", "AbortError")); + + expect((await pending).status).toBe(499); + // Fetch observes the client abort first; the pre-reader guard then attempts body-level + // settlement, and the reader owns the locked-stream fallback before releasing its lock. + expect(events).toEqual([ + "body.getReader", + "reader.read", + "fetch.abort", + "body.cancel", + "reader.cancel", + "reader.releaseLock", + ]); + } finally { + globalThis.fetch = originalFetch; + } }); - test("the bounded reader exclusively owns all non-combo Responses error bodies", async () => { - const source = readResponsesCoreSource(); + test.each([ + { label: "passthrough", adapter: "openai-responses", model: "fixture/model", combos: undefined }, + { label: "translated adapter", adapter: "openai-chat", model: "fixture/model", combos: undefined }, + { + label: "combo", + adapter: "openai-responses", + model: "combo/fallback", + combos: { fallback: { strategy: "failover" as const, targets: [{ provider: "fixture", model: "model" }] } }, + }, + ] as const)("$label Responses failure consumes each original body exactly once", async ({ adapter, model, combos }) => { + const originalFetch = globalThis.fetch; + const bodyReads: number[] = []; + globalThis.fetch = (async () => { + const body = new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode(JSON.stringify({ error: { message: "upstream failed" } }))); + controller.close(); + }, + }); + const response = new Response(body, { + status: 503, + headers: { "content-type": "application/json" }, + }); + const index = bodyReads.push(0) - 1; + Object.defineProperty(response, "body", { + configurable: true, + get() { + bodyReads[index] += 1; + return body; + }, + }); + return response; + }) as typeof fetch; - expect(source.match(/\breadDisplaySafeErrorText\(/g)).toHaveLength(4); - expect(source).not.toContain("detachPassthroughErrorGuard"); - expect(source).not.toContain("detachErrorBodyGuard"); - expect(source).not.toContain("detachContinuationErrorGuard"); - expect(source).not.toContain("upstreamResponse.text().catch(() => \"\")"); - expect(source).not.toContain("upstreamResponse.text().catch(() => \"unknown error\")"); - expect(source).not.toContain("response.text().catch(() => \"unknown error\")"); + const config = { + defaultProvider: "fixture", + providers: { + fixture: { + adapter, + baseUrl: "https://fixture.example.test/v1", + apiKey: "sk-test", + }, + }, + ...(combos ? { combos } : {}), + } as OcxConfig; + + try { + const response = await handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ model, input: "hello", stream: false }), + }), config, { model: "", provider: "" }); + await response.text(); + expect(bodyReads.length).toBeGreaterThan(0); + expect(bodyReads.every(reads => reads === 1)).toBe(true); + } finally { + globalThis.fetch = originalFetch; + } }); - // The combo branches are deliberately NOT guarded: consumeComboFailure -> - // readBoundedResponseBody reads `response.body` itself with the abort signal threaded - // through, and the combo contract is that the getter is touched exactly once (pinned by - // "captures passthrough failed usage from its original bounded body exactly once" in - // tests/server/server-combo-failover-e2e.test.ts). An earlier revision guarded them anyway and - // broke that test by adding a second `.body` read. - test("the combo failure branches do not add a second body read", async () => { - const source = readResponsesCoreSource(); - - for (const marker of ["const failure = await consumeComboFailure("]) { - let from = 0; - for (;;) { - const at = source.indexOf(marker, from); - if (at === -1) break; - // Look back a short window: no body guard may be attached immediately before a - // combo consumption. - const preceding = source.slice(Math.max(0, at - 400), at); - expect(preceding).not.toContain("cancelBodyOnAbort(upstreamResponse.body"); - from = at + marker.length; + test("terminal continuation failure consumes its original body exactly once", async () => { + const originalFetch = globalThis.fetch; + const continuationBodyReads: number[] = []; + let sends = 0; + globalThis.fetch = (async () => { + sends += 1; + if (sends === 1) { + return new Response([ + 'data: {"choices":[{"delta":{"content":"我接下来会修改相关文件。"}}]}\n\n', + 'data: {"choices":[{"delta":{},"finish_reason":"stop"}]}\n\n', + "data: [DONE]\n\n", + ].join(""), { headers: { "content-type": "text/event-stream" } }); } + + const body = new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode(JSON.stringify({ error: { message: "continuation failed" } }))); + controller.close(); + }, + }); + const response = new Response(body, { + status: 503, + headers: { "content-type": "application/json" }, + }); + const index = continuationBodyReads.push(0) - 1; + Object.defineProperty(response, "body", { + configurable: true, + get() { + continuationBodyReads[index] += 1; + return body; + }, + }); + return response; + }) as typeof fetch; + + const config = { + defaultProvider: "fixture", + providers: { + fixture: { + adapter: "openai-chat", + baseUrl: "https://fixture.example.test/v1", + apiKey: "sk-test", + terminalContinuationGuard: true, + }, + }, + } as OcxConfig; + + try { + const response = await handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ + model: "fixture/model", + input: "请检查这个问题并修复代码", + stream: true, + tools: [{ + type: "function", + name: "exec_command", + description: "run a command", + parameters: { type: "object" }, + }], + }), + }), config, { model: "", provider: "" }); + await response.text(); + + expect(response.status).toBe(200); + expect(continuationBodyReads.length).toBeGreaterThan(0); + expect(continuationBodyReads.every(reads => reads === 1)).toBe(true); + } finally { + globalThis.fetch = originalFetch; } });