Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 6 additions & 1 deletion src/server/responses/core.ts
Original file line number Diff line number Diff line change
Expand Up @@ -130,7 +130,12 @@ async function handleResponsesInner(
responseEffects,
sendBudgetState,
);
if (sidecarPlans instanceof Response) return sidecarPlans;
if (sidecarPlans instanceof Response) {
// A streamed sidecar result (web search / media bridge) still owns the probe: the loop
// releases it when the stream settles. Only a rejection is terminal here.
if (!sidecarPlans.ok) sidecarState.openAiSidecar?.releaseProbeLease?.();
return sidecarPlans;
Comment thread
devin-ai-integration[bot] marked this conversation as resolved.
}
const completionPolicy = createResponsesCompletionPolicy(requestContext, sidecarState);
if (transportState.adapter.runTurn) return await executeResponsesRunTurn(
requestContext,
Expand Down
17 changes: 15 additions & 2 deletions src/server/responses/sidecar-execution.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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), {
Comment thread
devin-ai-integration[bot] marked this conversation as resolved.
status: imgResponse.status,
headers: imgResponse.headers,
});
}
releaseSearchProbeLease();
return imgResponse;
} // end else (streaming bridge)
}
Expand Down Expand Up @@ -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;
}

Expand Down
144 changes: 143 additions & 1 deletion tests/responses/responses-run-turn-web-search.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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<typeof sidecarAuth.prepareResponsesSidecarAuth>) => {
if (args[0].req.headers.get("x-fixture-probe") !== "held") {
return prepareResponsesSidecarAuth(...args);
}
return {
routedCompaction: false,
openAiSidecar: { releaseProbeLease: () => { releasedFixtureProbe = true; } },
} as Awaited<ReturnType<typeof sidecarAuth.prepareResponsesSidecarAuth>>;
},
}));
const { handleResponses } = await import("../../src/server/responses");
const originalHome = process.env.OPENCODEX_HOME;
let home = "";
Expand Down Expand Up @@ -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,
Expand Down
Loading