diff --git a/src/images/loop.ts b/src/images/loop.ts index c09a07336b5..dd12fbb9d48 100644 --- a/src/images/loop.ts +++ b/src/images/loop.ts @@ -325,12 +325,19 @@ export interface ImageBridgeDeps { * `retryParsed` is the exact iteration-local request the retry will be built from. The loop * sends a shallow copy of the outer parsed request, so a rotation that rebinds only the outer * object never reaches the wire. Optional so existing callers keep compiling. + * + * A rotation may cross ACCOUNTS, not just keys, and the attempt row is where an operator reads + * which one happened. A rotator that knows which kind it performed returns it alongside the + * adapter -- the same `{ adapter, recoveryKind }` shape `onCredentialError` already uses below. */ on429?: ( retryAfterHeader: string | null, responseHeaders?: Headers, retryParsed?: OcxParsedRequest, - ) => ProviderAdapter | null | Promise; + ) => + | { adapter: ProviderAdapter; recoveryKind: AttemptRecoveryKind } + | null + | Promise<{ adapter: ProviderAdapter; recoveryKind: AttemptRecoveryKind } | null>; /** Opt-in same-target 429 policy (key-auth providers). When present, 429 replays on the SAME key before on429 rotation. */ retryOn429Policy?: Required | null; /** Called when the bridged Responses stream completes (parity with runTurn / routed paths). */ @@ -668,9 +675,9 @@ export async function runWithImageBridge(deps: ImageBridgeDeps): Promise {}); } catch { /* already closed */ } - adapter = rotated; + adapter = rotated.adapter; yield { type: "heartbeat" }; - prepared = await fetchOnce(adapter, "key-429"); + prepared = await fetchOnce(adapter, rotated.recoveryKind); } // Final headers have arrived. Clear only the deadline timer before ANY body read. diff --git a/src/server/responses/sidecar-execution.ts b/src/server/responses/sidecar-execution.ts index bab63758574..fd927032c8b 100644 --- a/src/server/responses/sidecar-execution.ts +++ b/src/server/responses/sidecar-execution.ts @@ -157,7 +157,12 @@ export async function executeResponsesSidecars( retryAfter: string | null, responseHeaders?: Headers, retryParsed?: OcxParsedRequest, - ): Promise => { + ): Promise<{ adapter: ProviderAdapter; recoveryKind: AttemptRecoveryKind } | null> => { + // Which credential axis actually moved. The main routed path already reports these three + // separately (`adapter-dispatch`: key-429 / anthropic-oauth-429 / oauth-account-429); the + // sidecar loops used to flatten all three to `key-429`, so an account rotation read as a key + // rotation in the attempt row and in the Logs UI. + let recoveryKind: AttemptRecoveryKind = "key-429"; const rotated = rotateProviderTransportOn429(config, route.providerName, route.provider, { retryAfter, now: Date.now(), @@ -206,6 +211,7 @@ export async function executeResponsesSidecars( hop.permit?.release(); return null; } + recoveryKind = "oauth-account-429"; hop.permit?.use(); } else if ( // Anthropic's pool is excluded from generic failover, so without this arm a 429 inside a @@ -249,6 +255,7 @@ export async function executeResponsesSidecars( hop.permit?.release(); return null; } + recoveryKind = "anthropic-oauth-429"; hop.permit?.use(); } else { // No key pool, no generic OAuth roster, no Anthropic pool could produce a replacement @@ -278,7 +285,7 @@ export async function executeResponsesSidecars( provider: route.provider, adapterName: rotatedAdapter.name, }); - return rotatedAdapter; + return { adapter: rotatedAdapter, recoveryKind }; }; if ((imgPlan || vidPlan) && canRunWebSearch) { // Web search takes priority when both are active — the media bridge cannot run diff --git a/src/web-search/loop.ts b/src/web-search/loop.ts index 76f5748130a..0f3784cb8ea 100644 --- a/src/web-search/loop.ts +++ b/src/web-search/loop.ts @@ -337,12 +337,19 @@ export interface WebSearchLoopDeps { * `retryParsed` is the exact iteration-local request the retry will be built from. The loop * sends a shallow copy of the outer parsed request, so a rotation that rebinds only the outer * object never reaches the wire. Optional so existing callers keep compiling. + * + * A rotation may cross ACCOUNTS, not just keys, and the attempt row is where an operator reads + * which one happened. A rotator that knows which kind it performed returns it alongside the + * adapter -- the same `{ adapter, recoveryKind }` shape `onCredentialError` already uses below. */ on429?: ( retryAfterHeader: string | null, responseHeaders?: Headers, retryParsed?: OcxParsedRequest, - ) => ProviderAdapter | null | Promise; + ) => + | { adapter: ProviderAdapter; recoveryKind: AttemptRecoveryKind } + | null + | Promise<{ adapter: ProviderAdapter; recoveryKind: AttemptRecoveryKind } | null>; /** Opt-in same-target 429 policy (key-auth providers). When present, 429 replays on the SAME key before on429 rotation. */ retryOn429Policy?: Required | null; /** Called only when the final bridged Responses stream reaches completed or incomplete. */ @@ -568,10 +575,10 @@ export async function runWithWebSearch(deps: WebSearchLoopDeps): Promise {}); } catch { /* already closed */ } - adapter = rotated; + adapter = rotated.adapter; // Stall-watchdog seam between bounded retry fetches (audit 011 B3). yield { type: "heartbeat" }; - prepared = await fetchOnce(adapter, "key-429"); + prepared = await fetchOnce(adapter, rotated.recoveryKind); } // Final headers have arrived. Clear only the deadline timer before ANY body read. diff --git a/tests/adapters/anthropic/anthropic-sidecar-account-failover.test.ts b/tests/adapters/anthropic/anthropic-sidecar-account-failover.test.ts index 7328d081f42..b4686b8ad5c 100644 --- a/tests/adapters/anthropic/anthropic-sidecar-account-failover.test.ts +++ b/tests/adapters/anthropic/anthropic-sidecar-account-failover.test.ts @@ -9,6 +9,7 @@ import { mkdtempSync } from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; import type { AdapterRequest, IncomingMeta, ProviderAdapter } from "../../../src/adapters/base"; +import type { AttemptRecoveryKind } from "../../../src/usage/log"; import { clearAnthropicAccountPoolState } from "../../../src/oauth/anthropic-routing"; import { clearGenericFailoverHealth } from "../../../src/oauth/generic-account-failover"; import { getAccountSet, saveCredential, setActiveAccount } from "../../../src/oauth/store"; @@ -71,7 +72,9 @@ beforeAll(async () => { adapter: ProviderAdapter; incomingMeta: IncomingMeta; fetchForRequest: (request: AdapterRequest, parsed: OcxParsedRequest) => typeof fetch; - on429?: (retryAfter: string | null) => Promise; + on429?: (retryAfter: string | null) => Promise< + { adapter: ProviderAdapter; recoveryKind: AttemptRecoveryKind } | null + >; }) => { // This is a dispatch seam test. The real loop is covered in anthropic-quota-dispatch. const first = await args.adapter.buildRequest(args.parsed, args.incomingMeta); @@ -83,7 +86,11 @@ beforeAll(async () => { await refused.body?.cancel(); const rotated = await args.on429?.(retryAfter); if (!rotated) throw new Error("Anthropic sidecar did not rotate after 429"); - const second = await rotated.buildRequest(args.parsed, args.incomingMeta); + // Unwrapped exactly as the real loop does. This seam drives the PRODUCTION rotator + // (`rotateSidecarProviderOn429`), so it is the one place the Anthropic arm's kind is + // proven end to end rather than against a hand-written stub. + expect(rotated.recoveryKind).toBe("anthropic-oauth-429"); + const second = await rotated.adapter.buildRequest(args.parsed, args.incomingMeta); return args.fetchForRequest(second, args.parsed)(second.url, { method: second.method, headers: second.headers, body: second.body, }); diff --git a/tests/images/loop.test.ts b/tests/images/loop.test.ts index 84dd5f13257..f5a06d867ae 100644 --- a/tests/images/loop.test.ts +++ b/tests/images/loop.test.ts @@ -4,6 +4,7 @@ import { join } from "node:path"; import { randomUUID } from "node:crypto"; import type { ProviderAdapter, IncomingMeta } from "../../src/adapters/base"; import type { AdapterEvent, OcxParsedRequest } from "../../src/types"; +import type { AttemptRecoveryKind } from "../../src/usage/log"; import type { ImageBridgePlan, ImageCallResult } from "../../src/images/types"; import type { ImageBridgeDeps } from "../../src/images/loop"; import { createTestTranslatorBudget } from "../helpers/translator-budget"; @@ -664,13 +665,16 @@ describe("runWithImageBridge", () => { // First rotation returns a new adapter that also 429s; the exhausted budget must not // re-arm for it. Second call returns null to terminate the pool. return rotations === 1 - ? ({ - ...mockAdapter, - fetchResponse: async () => { - sends += 1; - return new Response("{}", { status: 429 }); - }, - } as ProviderAdapter) + ? { + adapter: { + ...mockAdapter, + fetchResponse: async () => { + sends += 1; + return new Response("{}", { status: 429 }); + }, + } as ProviderAdapter, + recoveryKind: "key-429" as const, + } : null; }, onAttemptSend: recovery => { @@ -927,7 +931,7 @@ describe("runWithImageBridge", () => { retryParsed._kiroAuthContext = { apiRegion: "ap-southeast-2", profileArn: "account-b" }; delete retryParsed._providerContinuation; activeAdapter = secondAdapter; - return secondAdapter; + return { adapter: secondAdapter, recoveryKind: "key-429" }; }, }); const sse = await response.text(); @@ -938,6 +942,50 @@ describe("runWithImageBridge", () => { expect(retryState?._providerContinuation).toBeUndefined(); expect(sse).toContain("after rotate"); }); + + // An account rotation and a key rotation are different operator-facing events, and the + // rotated fetch's recovery kind is the only place the attempt row records which happened. + // The loop used to hardcode `key-429` for both. + test("429 rotation reports the rotator's recovery kind", async () => { + const recoveryKindsFor = async ( + rotation: (next: ProviderAdapter) => { adapter: ProviderAdapter; recoveryKind: AttemptRecoveryKind }, + ): Promise<(AttemptRecoveryKind | undefined)[]> => { + let fetchCalls = 0; + const sends: (AttemptRecoveryKind | undefined)[] = []; + const makeAdapter = (label: string): ProviderAdapter => ({ + name: label, + buildRequest: async () => ({ url: "https://test/v1/chat", method: "POST", headers: {}, body: "{}" }), + fetchResponse: async () => { + fetchCalls++; + if (fetchCalls === 1) return new Response("rate limited", { status: 429, headers: { "retry-after": "1" } }); + streamQueue = [[{ type: "text_delta", text: "after rotate" }, { type: "done" }]]; + return new Response("{}", { status: 200 }); + }, + parseStream: async function* (): AsyncGenerator { + const events = streamQueue.shift(); + if (events) for (const e of events) yield e; + }, + }); + const secondAdapter = makeAdapter("after-rotate"); + const response = await runWithImageBridge({ + parsed: makeParsed(), + adapter: makeAdapter("before-rotate"), + plan, + onAttemptSend: recovery => { sends.push(recovery); }, + on429: () => rotation(secondAdapter), + }); + expect(await response.text()).toContain("after rotate"); + return sends; + }; + + // A rotator that crossed accounts says so, and the rotated send carries that kind. + expect(await recoveryKindsFor(next => ({ adapter: next, recoveryKind: "oauth-account-429" }))) + .toEqual([undefined, "oauth-account-429"]); + expect(await recoveryKindsFor(next => ({ adapter: next, recoveryKind: "anthropic-oauth-429" }))) + .toEqual([undefined, "anthropic-oauth-429"]); + expect(await recoveryKindsFor(next => ({ adapter: next, recoveryKind: "key-429" }))) + .toEqual([undefined, "key-429"]); + }); }); // --------------------------------------------------------------------------- diff --git a/tests/web-search/web-search-timeout-contract.test.ts b/tests/web-search/web-search-timeout-contract.test.ts index 83d1c63f5fa..39e152d57ed 100644 --- a/tests/web-search/web-search-timeout-contract.test.ts +++ b/tests/web-search/web-search-timeout-contract.test.ts @@ -452,7 +452,7 @@ describe("web-search timeout runtime contracts", () => { connectTimeoutMs, on429: () => { rotations++; - return rotatedAdapter; + return { adapter: rotatedAdapter, recoveryKind: "key-429" as const }; }, })); diff --git a/tests/web-search/web-search.test.ts b/tests/web-search/web-search.test.ts index a027a980759..de7f4f336d3 100644 --- a/tests/web-search/web-search.test.ts +++ b/tests/web-search/web-search.test.ts @@ -10,7 +10,8 @@ import { listOpenAiForwardSidecarCandidates, resolveFirstUsableOpenAiSidecar } f import { handleResponses } from "../../src/server/responses/core"; import { providerFetch } from "../../src/server/responses/fetch-helpers"; import type { AdapterEvent, OcxConfig, OcxProviderConfig } from "../../src/types"; -import type { AdapterFetchContext, ProviderAdapter } from "../../src/adapters/base"; +import type { AdapterFetchContext, AdapterRequest, ProviderAdapter } from "../../src/adapters/base"; +import type { AttemptRecoveryKind } from "../../src/usage/log"; import type { OcxMessage, OcxParsedRequest } from "../../src/types"; import { fakeChatGptJwt } from "../helpers/fake-chatgpt-jwt"; import { createTestTranslatorBudget } from "../helpers/translator-budget"; @@ -1148,7 +1149,7 @@ describe("web-search sidecar native web_search_call emission", () => { if (!retryParsed) throw new Error("the loop must pass the iteration request to on429"); retryParsed._kiroAuthContext = { apiRegion: "ap-southeast-2", profileArn: "account-b" }; delete retryParsed._providerContinuation; - return rotatedAdapter; + return { adapter: rotatedAdapter, recoveryKind: "key-429" }; }, }); expect(response.status).toBe(200); @@ -1173,6 +1174,66 @@ describe("web-search sidecar native web_search_call emission", () => { ]); }); + // An account rotation and a key rotation are different operator-facing events, and the + // rotated fetch's recovery kind is the only place the attempt row records which happened. + // The loop used to hardcode `key-429` for both. + test("429 rotation reports the rotator's recovery kind", async () => { + globalThis.fetch = (() => Promise.resolve(new Response( + 'event: response.completed\ndata: {"type":"response.completed"}\n\n', + { headers: { "Content-Type": "text/event-stream" } }, + ))) as typeof fetch; + + const recoveryKindsFor = async ( + rotation: (next: ProviderAdapter) => { adapter: ProviderAdapter; recoveryKind: AttemptRecoveryKind }, + ): Promise<(AttemptRecoveryKind | undefined)[]> => { + const sends: (AttemptRecoveryKind | undefined)[] = []; + const buildRequest = (): AdapterRequest => + ({ url: "https://routed.test/v1", method: "POST", headers: {}, body: "{}" }); + const firstAdapter: ProviderAdapter = { + name: "mock-429", + buildRequest, + fetchResponse: async () => new Response("rate limited", { status: 429, headers: { "retry-after": "30" } }), + async *parseStream() { /* unused */ }, + async parseResponse() { return [{ type: "done" }] as AdapterEvent[]; }, + }; + const rotatedAdapter: ProviderAdapter = { + name: "mock-rotated", + buildRequest, + fetchResponse: async () => new Response("{}", { status: 200 }), + async *parseStream() { + yield { type: "text_delta", text: "answer from rotated account" }; + yield { type: "done" }; + }, + async parseResponse() { throw new Error("parseResponse must be unreachable"); }, + }; + const response = await runWithWebSearch({ + parsed: parseRequest({ model: "routed/model", input: "hi", stream: true, tools: [{ type: "web_search" }] }), + adapter: firstAdapter, + forwardProvider, + hostedTool: { type: "web_search" }, + selectedForwardHeaders: new Headers({ authorization: "Bearer token" }), + settings: { model: "gpt-5.6-luna", reasoning: "low", timeoutMs: 30_000 }, + maxSearches: 1, + onAttemptSend: recovery => { sends.push(recovery); }, + on429: () => rotation(rotatedAdapter), + }); + expect(response.status).toBe(200); + const frames = await collectSse(response.body!); + const completed = frames.find(f => f.event === "response.completed")?.data.response as Record; + const output = completed.output as { type: string; content?: { text?: string }[] }[]; + expect(output.find(o => o.type === "message")?.content?.[0]?.text).toBe("answer from rotated account"); + return sends; + }; + + // A rotator that crossed accounts says so, and the rotated send carries that kind. + expect(await recoveryKindsFor(next => ({ adapter: next, recoveryKind: "oauth-account-429" }))) + .toEqual([undefined, "oauth-account-429"]); + expect(await recoveryKindsFor(next => ({ adapter: next, recoveryKind: "anthropic-oauth-429" }))) + .toEqual([undefined, "anthropic-oauth-429"]); + expect(await recoveryKindsFor(next => ({ adapter: next, recoveryKind: "key-429" }))) + .toEqual([undefined, "key-429"]); + }); + test("retryOn429 replays on the same key before on429 rotation", async () => { globalThis.fetch = (() => Promise.resolve(new Response( 'event: response.completed\ndata: {"type":"response.completed"}\n\n', @@ -1421,7 +1482,7 @@ describe("web-search sidecar native web_search_call emission", () => { settings: { model: "gpt-5.6-luna", reasoning: "low", timeoutMs: 30_000 }, maxSearches: 1, connectTimeoutMs: 100, - on429: () => rotatedAdapter, + on429: () => ({ adapter: rotatedAdapter, recoveryKind: "key-429" }), }); expect(response.status).toBe(504); const body = await response.json() as { error?: { message?: string } };