diff --git a/src/server/responses/core.ts b/src/server/responses/core.ts index 635d5e37c17..aea67f586d2 100644 --- a/src/server/responses/core.ts +++ b/src/server/responses/core.ts @@ -130,6 +130,7 @@ async function handleResponsesInner( responseEffects, sendBudgetState, ); + if (sidecarPlans instanceof Response && !sidecarPlans.ok) sidecarState.openAiSidecar?.releaseProbeLease?.(); if (sidecarPlans instanceof Response) return sidecarPlans; const completionPolicy = createResponsesCompletionPolicy(requestContext, sidecarState); if (transportState.adapter.runTurn) return await executeResponsesRunTurn( diff --git a/src/server/responses/sidecar-execution.ts b/src/server/responses/sidecar-execution.ts index 24ffc9f126c..3bea7c078de 100644 --- a/src/server/responses/sidecar-execution.ts +++ b/src/server/responses/sidecar-execution.ts @@ -104,6 +104,17 @@ export async function executeResponsesSidecars( cancelResponseCompletion, } = responseEffects; + // Resolving the OpenAI search credential may hold the account's sole + // cooldown-recovery probe lease. A streamed sidecar result keeps it until the + // stream settles — completion or client cancel — so a later in-stream search + // outcome can still clear the cooldown; a response with no live body is + // terminal, so the lease is handed back before returning it. A recorded + // search outcome already settled the lease, making each release a + // generation-bound no-op. + const releaseSearchProbeLease = (): void => { + openAiSidecar?.releaseProbeLease?.(); + }; + // Tool results are PAIRED by call_id. parseRequest writes it into OcxToolResultMessage.toolCallId // (parser.ts:738/752) without validating it, because inputItemSchema's permissive catch-all @@ -402,11 +413,12 @@ export async function executeResponsesSidecars( if (imgResponse.body) { const imgTurnAc = new AbortController(); imgTurnAc.signal.addEventListener("abort", cancelResponseCompletion, { once: true }); - return new Response(trackStreamLifetime(imgResponse.body, imgTurnAc, undefined, options.turnAdmissionLease), { + return new Response(trackStreamLifetime(imgResponse.body, imgTurnAc, releaseSearchProbeLease, options.turnAdmissionLease), { status: imgResponse.status, headers: imgResponse.headers, }); } + releaseSearchProbeLease(); return imgResponse; } // end else (streaming bridge) } @@ -478,11 +490,12 @@ export async function executeResponsesSidecars( if (wsResponse.body) { const wsTurnAc = new AbortController(); wsTurnAc.signal.addEventListener("abort", cancelResponseCompletion, { once: true }); - return new Response(trackStreamLifetime(wsResponse.body, wsTurnAc, undefined, options.turnAdmissionLease), { + return new Response(trackStreamLifetime(wsResponse.body, wsTurnAc, releaseSearchProbeLease, options.turnAdmissionLease), { status: wsResponse.status, headers: wsResponse.headers, }); } + releaseSearchProbeLease(); return wsResponse; } diff --git a/tests/responses/responses-run-turn-web-search.test.ts b/tests/responses/responses-run-turn-web-search.test.ts index 1477290f7ab..afbf1d50c49 100644 --- a/tests/responses/responses-run-turn-web-search.test.ts +++ b/tests/responses/responses-run-turn-web-search.test.ts @@ -28,9 +28,23 @@ function fixture(provider: OcxProviderConfig): ProviderAdapter { }, }; } +function fetchFixture(provider: OcxProviderConfig): ProviderAdapter { + return { + name: "fetchonly", + buildRequest: () => ({ url: provider.baseUrl, method: "POST", headers: {}, body: "{}" }), + fetchResponse: async () => new Response("{}", { status: 200 }), + async *parseStream() { + yield { type: "text_delta", text: "search-enabled answer" } as AdapterEvent; + yield { type: "done" } as AdapterEvent; + }, + async parseResponse() { return [{ type: "done" }] as AdapterEvent[]; }, + }; +} mock.module("../../src/server/adapter-resolve", () => ({ ...resolver, resolveAdapter: (provider: OcxProviderConfig, cache?: "none" | "short" | "long") => - provider.adapter === "cursor" ? fixture(provider) : resolveAdapter(provider, cache), + provider.adapter === "cursor" ? fixture(provider) + : provider.adapter === "fetchonly" ? fetchFixture(provider) + : resolveAdapter(provider, cache), })); const pacing = await import("../../src/providers/request-pacing"); const originalWaitForSlot = pacing.waitForProviderRequestSlot; @@ -40,6 +54,20 @@ mock.module("../../src/providers/request-pacing", () => ({ ...pacing, return originalWaitForSlot(...args); }, })); +const sidecarAuth = await import("../../src/server/responses/request-sidecar-auth"); +const prepareResponsesSidecarAuth = sidecarAuth.prepareResponsesSidecarAuth; +let releasedFixtureProbe = false; +mock.module("../../src/server/responses/request-sidecar-auth", () => ({ ...sidecarAuth, + prepareResponsesSidecarAuth: async (...args: Parameters) => { + if (args[0].req.headers.get("x-fixture-probe") !== "held") { + return prepareResponsesSidecarAuth(...args); + } + return { + routedCompaction: false, + openAiSidecar: { releaseProbeLease: () => { releasedFixtureProbe = true; } }, + } as Awaited>; + }, +})); const { handleResponses } = await import("../../src/server/responses"); const originalHome = process.env.OPENCODEX_HOME; let home = ""; @@ -131,6 +159,120 @@ test.each(["image", "video"] as const)("media-only %s bridge still injects its t expect(attempts[0].context.tools?.some(t => t.webSearch)).toBe(false); }); +test("releases a search probe when pre-dispatch validation rejects the request", async () => { + releasedFixtureProbe = false; + const config = { + port: 0, defaultProvider: "cursor", + webSearchSidecar: { backend: "exa", exaApiKey: "fixture-search-key" }, + providers: { + cursor: { adapter: "cursor", baseUrl: "https://api2.cursor.sh", authMode: "oauth", models: ["model"] }, + }, + } as OcxConfig; + const response = await handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", headers: { "content-type": "application/json", "x-fixture-probe": "held" }, + body: JSON.stringify({ model: "cursor/model", input: [{ + type: "function_call_output", output: "fixture result", + }], tools: [{ type: "web_search" }] }), + }), config, { model: "", provider: "" }); + + expect(response.status).toBe(400); + expect(releasedFixtureProbe).toBe(true); + expect(attempts).toHaveLength(0); +}); + +test("a streamed sidecar response keeps the search probe until the stream settles", async () => { + releasedFixtureProbe = false; + const realFetch = globalThis.fetch; + globalThis.fetch = (async (input: RequestInfo | URL) => { + const url = String(typeof input === "object" && "url" in input ? input.url : input); + if (url.includes("exa")) { + return new Response(JSON.stringify({ results: [{ title: "fixture", url: "https://fixture.test" }] }), + { status: 200, headers: { "content-type": "application/json" } }); + } + return new Response( + 'event: response.completed\ndata: {"type":"response.completed","response":{"output":[]}}\n\n', + { status: 200, headers: { "content-type": "text/event-stream" } }); + }) as typeof fetch; + try { + const config = { + port: 0, defaultProvider: "fetchonly", + webSearchSidecar: { backend: "exa", exaApiKey: "fixture-search-key" }, + providers: { + fetchonly: { adapter: "fetchonly", baseUrl: "https://fetchonly.test/v1", apiKey: "fixture-key", models: ["model"] }, + }, + } as OcxConfig; + const response = await handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", headers: { "content-type": "application/json", "x-fixture-probe": "held" }, + body: JSON.stringify({ model: "fetchonly/model", input: "search this", stream: true, + tools: [{ type: "web_search" }] }), + }), config, { model: "", provider: "" }); + + expect(response.status).toBe(200); + expect(response.headers.get("content-type")).toContain("event-stream"); + expect(releasedFixtureProbe).toBe(false); + await response.text(); + // The routed model answered without a web_search call, so no sidecar outcome + // settled the lease — the stream's own completion hands the probe back. + expect(releasedFixtureProbe).toBe(true); + } finally { + globalThis.fetch = realFetch; + } +}); + +test("a cancelled streamed sidecar response releases the search probe", async () => { + releasedFixtureProbe = false; + const realFetch = globalThis.fetch; + globalThis.fetch = (async () => new Response( + 'event: response.completed\ndata: {"type":"response.completed","response":{"output":[]}}\n\n', + { status: 200, headers: { "content-type": "text/event-stream" } })) as typeof fetch; + try { + const config = { + port: 0, defaultProvider: "fetchonly", + webSearchSidecar: { backend: "exa", exaApiKey: "fixture-search-key" }, + providers: { + fetchonly: { adapter: "fetchonly", baseUrl: "https://fetchonly.test/v1", apiKey: "fixture-key", models: ["model"] }, + }, + } as OcxConfig; + const response = await handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", headers: { "content-type": "application/json", "x-fixture-probe": "held" }, + body: JSON.stringify({ model: "fetchonly/model", input: "search this", stream: true, + tools: [{ type: "web_search" }] }), + }), config, { model: "", provider: "" }); + + expect(response.status).toBe(200); + expect(releasedFixtureProbe).toBe(false); + // Client disconnect: the tracked stream's cancel path must settle the lease the same + // way a completed stream does, or the probe stays held until process exit. + await response.body!.cancel(); + expect(releasedFixtureProbe).toBe(true); + } finally { + globalThis.fetch = realFetch; + } +}); + +test("a media-bridge stream releases the search probe when it settles", async () => { + releasedFixtureProbe = false; + events = [[{ type: "text_delta", text: "media answer" }, { type: "done" }]]; + const config = { + port: 0, defaultProvider: "cursor", + images: { bridgeEnabled: true }, + providers: { + cursor: { adapter: "cursor", baseUrl: "https://api2.cursor.sh", authMode: "oauth", models: ["model"] }, + }, + } as OcxConfig; + const response = await handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", headers: { "content-type": "application/json", "x-fixture-probe": "held" }, + body: JSON.stringify({ model: "cursor/model", input: "draw a fixture", stream: true, + tools: [{ type: "image_generation" }] }), + }), config, { model: "", provider: "" }); + + expect(response.status).toBe(200); + expect(response.headers.get("content-type")).toContain("event-stream"); + expect(releasedFixtureProbe).toBe(false); + await response.text(); + expect(releasedFixtureProbe).toBe(true); +}); + // Streaming only: a first-event 429 replays the turn while the superseded // attempt is still in-flight (the buffered path awaits it before collecting // events, so the race cannot exist there). When that attempt finally returns,