From 76aa665e645dfe4f1a1f0d340886d3c28607d9cd Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Tue, 22 Sep 2026 05:24:30 +0900 Subject: [PATCH 01/12] fix(web-search): combine replay-cache isolation with deadline-safe quota evidence Bind replay-cache restoration to request and serving identity, and preserve the original quota response when Retry-After would outlive the sidecar deadline. Keep the current developer-role folding contract in the destination fixture. (cherry picked from commit 498498d3c814e066eaca066faacf6a9b004697c2) --- src/adapters/openai-responses/passthrough.ts | 8 +- src/responses/bridge-search-replay-cache.ts | 30 +++-- src/server/responses/passthrough-delivery.ts | 7 +- src/server/responses/request-prepare.ts | 6 + src/server/responses/request-transport.ts | 7 +- src/types/request.ts | 2 + src/web-search/executor.ts | 28 +++-- structure/providers-and-adapters.md | 15 ++- structure/runtime.md | 2 + .../chat-inline-document-bytes.test.ts | 15 ++- tests/server/server-key-failover-e2e.test.ts | 117 +++++++++++++++++- .../web-search-bridge-replay.test.ts | 57 +++++++-- .../web-search/web-search-sidecar-429.test.ts | 48 ++++++- 13 files changed, 291 insertions(+), 51 deletions(-) diff --git a/src/adapters/openai-responses/passthrough.ts b/src/adapters/openai-responses/passthrough.ts index cadcb2a38c2..0bef9c1cf06 100644 --- a/src/adapters/openai-responses/passthrough.ts +++ b/src/adapters/openai-responses/passthrough.ts @@ -329,11 +329,11 @@ export function createResponsesPassthroughAdapter(provider: OcxProviderConfig): outBody = stripUnsupportedReasoningSummaryDelivery(outBody, parsed.modelId); // #4587: on a bridged provider, hand the destination back the search call and result the // proxy executed on its behalf, in place of the hosted cell the caller replays. Scoped to - // this destination and recorded by the bridge itself, so a provider without the opt-in - // computes no identity and keeps the body reference it already had. This runs before the - // query backfill below because a restored cell is no longer a web_search_call to repair. + // its exact conversation and serving identity and recorded by the bridge itself, so a + // provider without the opt-in computes no identity and keeps the body reference it already + // had. This runs before query backfill because a restored cell is no longer one to repair. if (provider.webSearchBridge?.enabled === true) { - outBody = restoreBridgedWebSearchCalls(outBody, bridgeSearchReplayScope(provider.baseUrl)); + outBody = restoreBridgedWebSearchCalls(outBody, bridgeSearchReplayScope(parsed._reasoningReplayScope)); } // Repair stored history from before the bridge emitted both keys, in either // direction: a conversation that already recorded a web_search_call replays it diff --git a/src/responses/bridge-search-replay-cache.ts b/src/responses/bridge-search-replay-cache.ts index 6de24cdce39..07213682c33 100644 --- a/src/responses/bridge-search-replay-cache.ts +++ b/src/responses/bridge-search-replay-cache.ts @@ -14,9 +14,9 @@ * what `appendBridgeSearchTurn` would have written onto a continuation leg, so a replayed turn * and a continued turn show the destination the same conversation. * - * Scope. Entries are keyed by the upstream destination in addition to the cell id. The cell id is - * a v4 UUID minted here, so it cannot collide across conversations, but an unscoped key would let - * a history replayed against a DIFFERENT provider resurrect a call that provider never made. + * Scope. Entries are keyed by the exact conversation and serving identity in addition to the cell + * id. The cell id is a v4 UUID minted here, but possession of a client-visible id is not authority + * to recover result text under another provider, model, destination, or credential. * * Bounds and privacy. Result text is web content the caller already received, but it is still * request-derived data: it lives in memory only, is never logged, serialized, or exported, and is @@ -26,7 +26,7 @@ * alone. Neither re-running the search nor inventing a result is an acceptable recovery. */ -import { reasoningReplayDestinationIdentity } from "./reasoning-replay-cache"; +import type { OcxReasoningReplayScopeRef } from "../types"; const MAX_ENTRIES = 64; const MAX_TOTAL_BYTES = 512 * 1024; @@ -58,14 +58,24 @@ let clockForTests: (() => number) | null = null; const now = (): number => clockForTests?.() ?? Date.now(); /** - * Identify the upstream destination a bridged search belongs to. + * Identify the exact conversation and upstream binding a bridged search belongs to. * - * Reuses the salted process-local destination digest the reasoning replay cache already defines, - * so both stores agree on what "the same upstream" means and neither invents a second notion of - * destination identity. + * The serving route binds this holder only after provider, model, and physical credential + * selection. A missing conversation or binding fails closed: a cell id is client-visible and is + * not itself authority to recover another request's retained result. */ -export function bridgeSearchReplayScope(baseUrl: string | undefined): string | undefined { - return reasoningReplayDestinationIdentity(baseUrl); +export function bridgeSearchReplayScope(scope: OcxReasoningReplayScopeRef | undefined): string | undefined { + const identity = scope?.current; + if (!scope?.clientPrincipalId || !scope.clientThreadId || !identity) return undefined; + return JSON.stringify([ + scope.clientPrincipalId, + scope.clientThreadId, + identity.providerName, + identity.providerDestinationIdentity, + identity.adapterName, + identity.modelId, + identity.credentialIdentity, + ]); } function keyFor(scope: string, cellItemId: string): string { diff --git a/src/server/responses/passthrough-delivery.ts b/src/server/responses/passthrough-delivery.ts index 969cd84e589..b4941fbd6dd 100644 --- a/src/server/responses/passthrough-delivery.ts +++ b/src/server/responses/passthrough-delivery.ts @@ -417,10 +417,9 @@ export async function deliverPassthroughResponse( describeImages: requiresVisionPreprocessing(config, route.provider, route.modelId, route.providerName), sidecar: config.webSearchSidecar, }), - // Scope the executed-search memo to this exact upstream (#4587). The Responses adapter - // derives the same scope from the same base URL before the NEXT turn is dispatched, so - // a replayed hosted cell can be turned back into the destination's own call and result. - destinationScope: bridgeSearchReplayScope(route.provider.baseUrl), + // Snapshot the bound conversation, provider, model, destination, and credential. The + // next turn must match every dimension before its hosted cell can recover this result. + destinationScope: bridgeSearchReplayScope(parsed._reasoningReplayScope), // Appending a search result can push the continuation past the ceiling the first leg // was admitted under, so the same limit is re-applied before every later send. checkOutboundBody: (continuationBody: string) => { diff --git a/src/server/responses/request-prepare.ts b/src/server/responses/request-prepare.ts index 867f816414d..c9f04d73d73 100644 --- a/src/server/responses/request-prepare.ts +++ b/src/server/responses/request-prepare.ts @@ -20,6 +20,7 @@ import { sessionIdHeaderFromRequest, reasoningReplayConversationIdFromResponsesRequest, } from "../request-log-conversation"; +import { contextPrincipalIdOf } from "../auth-cors"; import { isShadowSourceModel, shadowSourceModelPrefix, @@ -407,6 +408,11 @@ export async function prepareResponsesRequest( parsed._reasoningReplayScope = { clientThreadId: reasoningReplayConversationId }; } } + if (parsed._reasoningReplayScope) { + const clientPrincipalId = contextPrincipalIdOf(options.admission) + ?? (options.admission?.kind === "loopback" ? "loopback" : undefined); + parsed._reasoningReplayScope = { ...parsed._reasoningReplayScope, clientPrincipalId }; + } // Prefer a pre-populated id (routed Claude) over Responses headers that may be // absent or synthetically injected (session_id from prompt_cache_key). if (!logCtx.conversationId) { diff --git a/src/server/responses/request-transport.ts b/src/server/responses/request-transport.ts index e6834c7e941..2b302402f38 100644 --- a/src/server/responses/request-transport.ts +++ b/src/server/responses/request-transport.ts @@ -445,6 +445,11 @@ export async function prepareResponsesTransport( return response; } const nextAdapter = await refreshDispatchAdapter(requestParsed); + // Rebind before rebuilding: the rebuild's bridged-search restore and continuation + // restore key on the serving identity, which must be the refreshed route's, not the + // credential whose selection just lapsed. + bindRouteReasoningReplayScope({ parsed: requestParsed, providerName: route.providerName, provider: route.provider, + adapterName: nextAdapter.name, oauthCredentialSnapshot: replayOAuthCredentialSnapshot }); const rebuilt = await nextAdapter.buildRequest(requestParsed, { headers: requestState.selectedForwardHeaders, translatorBudget, ...(imageTierBias > 0 ? { imageTierBias } : {}), @@ -467,8 +472,6 @@ export async function prepareResponsesTransport( sameTargetToken = transportToken; destination = rebuilt.url; dispatchInit = { ...dispatchInit, method: rebuilt.method, headers, body: rebuilt.body }; - bindRouteReasoningReplayScope({ parsed: requestParsed, providerName: route.providerName, provider: route.provider, - adapterName: nextAdapter.name, oauthCredentialSnapshot: replayOAuthCredentialSnapshot }); // The next iteration validates synchronously and calls fetch in that same turn. } throw new Error("OAuth account selection changed repeatedly before dispatch"); diff --git a/src/types/request.ts b/src/types/request.ts index 7a73eaa5987..975db448179 100644 --- a/src/types/request.ts +++ b/src/types/request.ts @@ -30,6 +30,8 @@ export interface OcxReasoningReplayIdentity { * the holder, so late tool-call cache writes see the active physical identity. */ export interface OcxReasoningReplayScopeRef { + /** Process-local caller principal; `loopback` denotes the trusted local-only admission lane. */ + readonly clientPrincipalId?: string; /** * Conversation namespace for replay state. Historically this was always the Codex parent-thread * id; headerless Responses callers use a raw sanitized thread/Cursor/session fallback, never the diff --git a/src/web-search/executor.ts b/src/web-search/executor.ts index 33b6a1eb6cc..e14288ebafa 100644 --- a/src/web-search/executor.ts +++ b/src/web-search/executor.ts @@ -49,10 +49,12 @@ export type SidecarOutcome = WebSearchResult & { error?: string }; * The forward backend throttles burst sidecar traffic, and without a replay the 429 becomes a * failed tool result that poisons the query for the whole turn (see failedQueries in loop.ts). * 1 initial send + 2 replays; Retry-After is honored as a lower bound and capped by - * RETRY_AFTER_CEILING_MS (an instruction past the ceiling ends with the 429 instead of - * parking the search). Each wait releases the unread 429 body first so sockets do not - * accumulate under a rate-limit storm. Abort or timeout ends the wait through the existing - * catch, exactly like an abort during the SSE parse. + * RETRY_AFTER_CEILING_MS and the remaining sidecar deadline (an instruction past either + * ends with the 429 instead of parking the search). Each wait releases the unread 429 body first so sockets do not + * accumulate under a rate-limit storm. The release itself may take up to a second, so a + * deadline landing during release or backoff ends with the 429 already in hand rather than + * a timeout; a caller abort still ends the wait through the shared catch, exactly like an + * abort during the SSE parse. */ const SIDECAR_429_MAX_ATTEMPTS = 3; const SIDECAR_429_BASE_DELAY_MS = 1_000; @@ -98,9 +100,10 @@ export async function runWebSearch( stream: true, }; const url = `${forwardProvider.baseUrl}/responses`; + // t0 precedes the deadline timer's start so the remaining-time check stays conservative. + const t0 = Date.now(); const linkedSignal = signalWithTimeout(settings.timeoutMs, abortSignal); const sidecarExit = sidecarEnter("web-search"); - const t0 = Date.now(); try { const sendOnce = () => fetchWithResetRetry( // Recovery nests INSIDE the version helper: applyUpstreamRecoveryInit then always receives a @@ -129,10 +132,19 @@ export async function runWebSearch( }); // A deadline, not a clamp: an instruction past the ceiling ends the search with the // 429 instead of parking it at a provider that already said it would refuse. - if (delay > RETRY_AFTER_CEILING_MS) break; + if (delay > RETRY_AFTER_CEILING_MS || delay >= settings.timeoutMs - (Date.now() - t0)) break; console.warn(`[web-search] sidecar HTTP 429 — retrying (${attempt + 2}/${SIDECAR_429_MAX_ATTEMPTS}) after ${delay}ms`); - await releaseResponseBodyBestEffort(res.body, linkedSignal.signal); - await sleepWithAbort(delay, linkedSignal.signal); + try { + await releaseResponseBodyBestEffort(res.body, linkedSignal.signal); + await sleepWithAbort(delay, linkedSignal.signal); + } catch (e) { + // The release above may consume up to 1s, so the sidecar deadline can land during + // cleanup or mid-backoff — before the replay is dispatched. The observed 429 is + // already in hand: end with it rather than laundering it into a timeout. A caller + // abort (or a non-deadline throw) still propagates to the shared catch below. + if (!linkedSignal.signal.aborted || linkedSignal.signal.reason === abortSignal?.reason) throw e; + break; + } res = await sendOnce(); } // Attach the body guard before ANY branch reads it. The success path guarded itself below, diff --git a/structure/providers-and-adapters.md b/structure/providers-and-adapters.md index 132265964a9..2256b0ce1ba 100644 --- a/structure/providers-and-adapters.md +++ b/structure/providers-and-adapters.md @@ -177,17 +177,22 @@ searches run, their hosted cells complete, the held client calls are released fo execute, and the leg's own terminal closes the turn with no continuation sent upstream. The destination therefore does not receive that search result during the turn. It gets it on the next one: every search the bridge executes is recorded in `src/responses/bridge-search-replay-cache.ts` -under the hosted cell's proxy-minted id, scoped to the upstream destination and bounded by entry -count, total bytes, and a one-hour TTL. When the caller replays that cell, +under the hosted cell's proxy-minted id, scoped to the admitted caller principal, client +conversation, and exact provider, adapter, model, destination, and physical credential binding, and bounded by entry count, total +bytes, and a one-hour TTL. An unavailable scope fails closed. When the caller replays that cell, `restoreBridgedWebSearchCalls` in `src/adapters/openai-responses/tool-output-recovery.ts` puts the destination's own `function_call` and the executed `function_call_output` back in the cell's position before the next turn's first leg is dispatched, recording exactly the text `appendBridgeSearchTurn` would have sent on a continuation leg so a replayed turn and a continued turn show the destination one consistent conversation. The rewrite runs only for a provider with -`webSearchBridge.enabled`, and a miss — unknown id, expired entry, a different destination, or a -`call_id` the body already carries — leaves the replayed item untouched. Re-running the search or -synthesizing result text is not a permitted recovery. The bridge finalizes request-scoped OpenAI sidecar authority on completion, failure, and client cancellation — cancellation releases immediately rather than waiting on an abandoned upstream read — so a recovery probe lease no search consumed is always returned. +`webSearchBridge.enabled`, and a miss — unknown id, expired entry, a different conversation or +serving binding, or a `call_id` the body already carries — leaves the replayed item untouched. +Re-running the search or synthesizing result text is not a permitted recovery. The bridge finalizes +request-scoped OpenAI sidecar authority on completion, failure, and client cancellation — +cancellation releases immediately rather than waiting on an abandoned upstream read — so a +recovery probe lease no search consumed is always returned. `tests/web-search/web-search-bridge-replay.test.ts` pins the restore and each of those refusals. +A forward OpenAI search sidecar retries a 429 only when the requested delay fits both its retry ceiling and the remaining overall sidecar deadline. A delay that cannot fit returns and records the original 429 so pool routing retains quota evidence. A leg whose upstream terminal is `response.failed` or `response.incomplete` runs no search at all and closes any cell it opened rather than leaving it in progress. Assistant text is not treated as a search diff --git a/structure/runtime.md b/structure/runtime.md index 26f5b2a43f8..cc89632691e 100644 --- a/structure/runtime.md +++ b/structure/runtime.md @@ -395,6 +395,8 @@ Automatic Codex pool selection and account status share the [plan exclusion cont ### Empty forced search answers `src/web-search/loop.ts` makes at most one extra answer attempt after a clean forced-answer terminal with no visible output or tool call. The recovery has no tools and reuses gathered search results. Malformed calls fail before refusal/truncation passthrough, and well-formed recognized refusal/truncation terminals pass through unchanged, including empty or partial answers. The extra generation may incur provider usage. + +OpenAI sidecar 429 replays run only when their backoff fits the remaining sidecar deadline; otherwise the original 429 remains the routing-health outcome rather than becoming a timeout. ## Scoped provider quota for Combo selection `src/providers/quota/report-cache.ts` publishes routing evidence only when a producer explicitly supplies its diff --git a/tests/responses/chat-inline-document-bytes.test.ts b/tests/responses/chat-inline-document-bytes.test.ts index 4221b174921..a2813b574c2 100644 --- a/tests/responses/chat-inline-document-bytes.test.ts +++ b/tests/responses/chat-inline-document-bytes.test.ts @@ -198,12 +198,19 @@ describe("inline document bytes reach a wire that can hold them", () => { stream: false, options: {}, } as unknown as OcxParsedRequest; - const outbound = JSON.parse(createOpenAIChatAdapter(chatProvider).buildRequest(parsed).body) as { - messages: Array<{ role: string; content: unknown }>; - }; - expect(outbound.messages).toEqual([{ + const buildBody = (provider: OcxProviderConfig) => JSON.parse( + createOpenAIChatAdapter(provider).buildRequest(parsed).body, + ) as { messages: Array<{ role: string; content: unknown }> }; + // Carrying a document must not demote the turn to `user`; which role the slot shows on the + // wire is the destination's recorded answer, so the accepting destination keeps `developer` + // and the unrecorded one folds to `system` in place. + expect(buildBody({ ...chatProvider, foldDeveloperRoleToSystem: false }).messages).toEqual([{ role: "developer", content: [{ type: "file", file: { file_data: PDF_DATA_URL, filename: "spec" } }], }]); + expect(buildBody({ ...chatProvider, foldDeveloperRoleToSystem: true }).messages).toEqual([{ + role: "system", + content: [{ type: "file", file: { file_data: PDF_DATA_URL, filename: "spec" } }], + }]); }); }); diff --git a/tests/server/server-key-failover-e2e.test.ts b/tests/server/server-key-failover-e2e.test.ts index a9753c6f27e..5279288c90c 100644 --- a/tests/server/server-key-failover-e2e.test.ts +++ b/tests/server/server-key-failover-e2e.test.ts @@ -7,7 +7,16 @@ import { readUsageEntries, resetUsageReadCacheForTests } from "../../src/usage/l import { loadConfig, saveConfig } from "../../src/config"; import { clearKeyCooldowns, getKeyCooldownUntil, rotateKeyOn429 } from "../../src/providers/key-failover"; import { deriveXaiConvId } from "../../src/providers/xai-transport"; -import { clearReasoningReplayCacheForTests } from "../../src/responses/reasoning-replay-cache"; +import { + clearReasoningReplayCacheForTests, + reasoningReplayDestinationIdentity, + reasoningReplayKeyCredentialIdentity, +} from "../../src/responses/reasoning-replay-cache"; +import { + bridgeSearchReplayScope, + clearBridgeSearchReplayCacheForTests, + rememberBridgeSearchReplay, +} from "../../src/responses/bridge-search-replay-cache"; import { startServer } from "../../src/server"; import { handleResponses } from "../../src/server/responses"; import type { OcxConfig } from "../../src/types"; @@ -31,6 +40,7 @@ beforeEach(() => { process.env.OPENCODEX_HOME = testDir; clearKeyCooldowns(); clearReasoningReplayCacheForTests(); + clearBridgeSearchReplayCacheForTests(); }); afterEach(() => { @@ -43,6 +53,7 @@ afterEach(() => { if (testDir) removeTreeWithRetry(testDir); clearKeyCooldowns(); clearReasoningReplayCacheForTests(); + clearBridgeSearchReplayCacheForTests(); }); describe("server 429 key failover (end-to-end)", () => { @@ -1182,3 +1193,107 @@ test.each([false, true])("key refetch retains transient recovery metadata (strea usage: { inputTokens: 12, outputTokens: 2 } }); } finally { await server.stop(true); } }); + +test("a dispatch-time key switch rebuilds the bridged-search restore under the new credential", async () => { + // Regression for the oauthDispatch rebuild order: the Responses adapter restores a replayed + // web_search_call from the memo keyed by _reasoningReplayScope, so the rebuild must rebind + // that scope to the refreshed credential BEFORE buildRequest runs. Restoring under the key + // whose selection just lapsed, then sending under the newly selected key, would hand the + // first credential's recorded result to the second credential's upstream. + let now = 0; + let resumePacing: (() => void) | undefined; + const queued = Promise.withResolvers(); + setProviderRequestPacingRuntimeForTest({ + now: () => now, + setTimer(callback, delayMs) { + resumePacing = () => { now += delayMs; callback(); }; + queued.resolve(); + return callback; + }, + clearTimer() {}, + enqueueMicrotask: queueMicrotask, + }); + const seen: { authorization: string | null; input: Record[] }[] = []; + upstream = Bun.serve({ hostname: "127.0.0.1", port: 0, async fetch(req) { + const body = await req.json() as { input?: unknown }; + seen.push({ + authorization: req.headers.get("authorization"), + input: Array.isArray(body.input) ? body.input as Record[] : [], + }); + return Response.json({ + id: "resp_keyrace", object: "response", status: "completed", model: "test", + output: [{ type: "message", id: "msg_keyrace", role: "assistant", status: "completed", + content: [{ type: "output_text", text: "done", annotations: [] }] }], + usage: { input_tokens: 1, output_tokens: 1, total_tokens: 2 }, + }); + } }); + const baseUrl = `http://127.0.0.1:${upstream.port}/v1`; + const config = { port: 0, hostname: "127.0.0.1", defaultProvider: "pooled", providers: { pooled: { + adapter: "openai-responses", baseUrl, allowPrivateNetwork: true, + authMode: "key", apiKey: "synthetic-first", + apiKeyPool: [{ id: "first", key: "synthetic-first" }, { id: "second", key: "synthetic-second" }], + webSearchBridge: { enabled: true, backend: "ollama" }, + requestPacing: { enabled: true, minIntervalMs: 100 }, + } } } as OcxConfig; + saveConfig(config); + + // Seed the bridged-search memo under the identity the FIRST key binds: same loopback + // principal and thread the request below carries, but the lapsed credential. + const cellId = "ws_keyrace"; + rememberBridgeSearchReplay( + bridgeSearchReplayScope({ + clientPrincipalId: "loopback", + clientThreadId: "thread-keyrace", + current: { + providerName: "pooled", + providerDestinationIdentity: reasoningReplayDestinationIdentity(baseUrl), + adapterName: "openai-responses", + modelId: "test", + credentialIdentity: reasoningReplayKeyCredentialIdentity({ apiKey: "synthetic-first" }), + }, + }), + cellId, + { callId: "call_ws_1", name: "web_search", + argumentsText: "{\"query\":\"opencodex release\"}", output: "cached bridged result" }, + ); + + const server = startServer(0); + const abort = new AbortController(); + try { + await waitForProviderRequestSlot("pooled", config.providers.pooled); + const pending = fetch(new URL("/v1/responses", server.url), { + method: "POST", + headers: { "content-type": "application/json", "thread-id": "thread-keyrace" }, + signal: abort.signal, + body: JSON.stringify({ + model: "pooled/test", stream: false, + input: [ + { role: "user", content: [{ type: "input_text", text: "what is the latest release?" }] }, + { type: "web_search_call", id: cellId, status: "completed", + action: { type: "search", query: "opencodex release", queries: ["opencodex release"] } }, + ], + }), + }); + await queued.promise; + const selected = await managementFetch(new URL("/api/providers/keys/active", server.url), { + method: "PUT", headers: { "content-type": "application/json" }, + body: JSON.stringify({ name: "pooled", id: "second" }), + }); + expect(selected.status).toBe(200); + await selected.text(); + resumePacing!(); + const response = await pending; + expect(response.status).toBe(200); + await response.text(); + expect(seen).toHaveLength(1); + expect(seen[0]!.authorization).toBe("Bearer synthetic-second"); + // Rebound before rebuild: the memo lookup misses under the new credential, so the hosted + // cell reaches the second key's upstream verbatim instead of the first key's result. + expect(seen[0]!.input.some(item => item.type === "web_search_call" && item.id === cellId)).toBe(true); + expect(seen[0]!.input.some(item => item.call_id === "call_ws_1")).toBe(false); + } finally { + abort.abort(); + await server.stop(true); + resetProviderRequestPacingForTest(); + } +}); diff --git a/tests/web-search/web-search-bridge-replay.test.ts b/tests/web-search/web-search-bridge-replay.test.ts index 3d5a6aab28c..2b98b3f63e5 100644 --- a/tests/web-search/web-search-bridge-replay.test.ts +++ b/tests/web-search/web-search-bridge-replay.test.ts @@ -20,7 +20,7 @@ import { } from "../../src/responses/bridge-search-replay-cache"; import { createResponsesPassthroughAdapter as createResponsesPassthroughAdapterProduction } from "../../src/adapters/openai-responses"; import { withTestTranslatorBudget } from "../helpers/translator-budget"; -import type { OcxProviderConfig } from "../../src/types"; +import type { OcxProviderConfig, OcxReasoningReplayScopeRef } from "../../src/types"; const createResponsesPassthroughAdapter = (...args: Parameters) => withTestTranslatorBudget(createResponsesPassthroughAdapterProduction(...args)); @@ -28,6 +28,21 @@ const createResponsesPassthroughAdapter = (...args: Parameters = {}): OcxReasoningReplayScopeRef { + return { + clientPrincipalId: "principal-a", + clientThreadId: "thread-a", + current: { + providerName: "bridge-a", + providerDestinationIdentity: overrides.providerDestinationIdentity ?? GATEWAY_BASE_URL, + adapterName: "openai-responses", + modelId: "glm-4.7", + credentialIdentity: "key-a", + ...overrides, + }, + }; +} + function frame(type: string, payload: Record): string { return "event: " + type + "\ndata: " + JSON.stringify({ type, ...payload }); } @@ -115,7 +130,7 @@ async function runBridgedMixedLeg(baseUrl: string, result = "opencodex 2.50.0 sh throw new Error("a mixed leg must not send a continuation"); }, execute: async () => ({ text: result, sources: [{ url: "https://example.test/rel", title: "Releases" }] }), - destinationScope: bridgeSearchReplayScope(baseUrl), + destinationScope: bridgeSearchReplayScope(replayScope({ providerDestinationIdentity: baseUrl })), }); const body = await new Response(stream).text(); const added = clientEvents(body).find(event => @@ -154,7 +169,7 @@ describe("bridged web_search replay to the destination", () => { const restored = restoreBridgedWebSearchCalls( nextTurnBody(cellId), - bridgeSearchReplayScope(GATEWAY_BASE_URL), + bridgeSearchReplayScope(replayScope()), ) as { input: Record[] }; // The item type the destination never produced is gone, replaced in place by the exchange @@ -188,7 +203,7 @@ describe("bridged web_search replay to the destination", () => { throw new Error("a mixed leg must not send a continuation"); }, execute: async () => ({ text: "", sources: [], error: "backend refused" }), - destinationScope: bridgeSearchReplayScope(GATEWAY_BASE_URL), + destinationScope: bridgeSearchReplayScope(replayScope()), }); const body = await new Response(stream).text(); const added = clientEvents(body).find(event => @@ -198,7 +213,7 @@ describe("bridged web_search replay to the destination", () => { const restored = restoreBridgedWebSearchCalls( nextTurnBody(cellId), - bridgeSearchReplayScope(GATEWAY_BASE_URL), + bridgeSearchReplayScope(replayScope()), ) as { input: Record[] }; expect(restored.input[2]).toEqual({ type: "function_call_output", @@ -209,7 +224,7 @@ describe("bridged web_search replay to the destination", () => { test("a cell this proxy never executed is left exactly as the caller sent it", () => { const body = nextTurnBody("ws_never-recorded"); - const restored = restoreBridgedWebSearchCalls(body, bridgeSearchReplayScope(GATEWAY_BASE_URL)); + const restored = restoreBridgedWebSearchCalls(body, bridgeSearchReplayScope(replayScope())); // Same reference: a miss allocates nothing and invents nothing. expect(restored).toBe(body); }); @@ -217,7 +232,26 @@ describe("bridged web_search replay to the destination", () => { test("a search recorded for one destination is not replayed into another", async () => { const cellId = await runBridgedMixedLeg(GATEWAY_BASE_URL); const body = nextTurnBody(cellId); - expect(restoreBridgedWebSearchCalls(body, bridgeSearchReplayScope(OTHER_BASE_URL))).toBe(body); + expect(restoreBridgedWebSearchCalls(body, bridgeSearchReplayScope(replayScope({ providerDestinationIdentity: OTHER_BASE_URL })))).toBe(body); + }); + + test("a cell cannot cross any conversation or serving-identity boundary", async () => { + const cellId = await runBridgedMixedLeg(GATEWAY_BASE_URL); + const body = nextTurnBody(cellId); + const mismatchedScopes: OcxReasoningReplayScopeRef[] = [ + { ...replayScope(), clientPrincipalId: "principal-b" }, + { ...replayScope(), clientThreadId: "thread-b" }, + replayScope({ providerName: "bridge-b" }), + replayScope({ adapterName: "other-adapter" }), + replayScope({ modelId: "other-model" }), + replayScope({ credentialIdentity: "key-b" }), + ]; + for (const scope of mismatchedScopes) { + expect(restoreBridgedWebSearchCalls(body, bridgeSearchReplayScope(scope))).toBe(body); + } + expect(bridgeSearchReplayScope(undefined)).toBeUndefined(); + expect(bridgeSearchReplayScope({ clientThreadId: "thread-a" })).toBeUndefined(); + expect(bridgeSearchReplayScope({ clientPrincipalId: "principal-a", clientThreadId: "thread-a" })).toBeUndefined(); }); test("an expired entry behaves exactly like a miss", async () => { @@ -226,13 +260,13 @@ describe("bridged web_search replay to the destination", () => { const cellId = await runBridgedMixedLeg(GATEWAY_BASE_URL); const body = nextTurnBody(cellId); // Still inside the TTL. - expect(restoreBridgedWebSearchCalls(body, bridgeSearchReplayScope(GATEWAY_BASE_URL))).not.toBe(body); + expect(restoreBridgedWebSearchCalls(body, bridgeSearchReplayScope(replayScope()))).not.toBe(body); clockMs += 61 * 60 * 1000; - expect(restoreBridgedWebSearchCalls(body, bridgeSearchReplayScope(GATEWAY_BASE_URL))).toBe(body); + expect(restoreBridgedWebSearchCalls(body, bridgeSearchReplayScope(replayScope()))).toBe(body); }); test("a call id the body already carries is never duplicated", () => { - const scope = bridgeSearchReplayScope(GATEWAY_BASE_URL); + const scope = bridgeSearchReplayScope(replayScope()); rememberBridgeSearchReplay(scope, "ws_dup", { callId: "call_2", name: "web_search", @@ -244,7 +278,7 @@ describe("bridged web_search replay to the destination", () => { }); test("an unbridged provider is never given a scope to restore from", () => { - const scope = bridgeSearchReplayScope(GATEWAY_BASE_URL); + const scope = bridgeSearchReplayScope(replayScope()); rememberBridgeSearchReplay(scope, "ws_unbridged", { callId: "call_1", name: "web_search", @@ -274,6 +308,7 @@ describe("the Responses passthrough adapter", () => { stream: true, options: {}, _rawBody: nextTurnBody(cellId), + _reasoningReplayScope: replayScope(), }, { headers: new Headers() }); return (JSON.parse(request.body) as { input: Record[] }).input; } diff --git a/tests/web-search/web-search-sidecar-429.test.ts b/tests/web-search/web-search-sidecar-429.test.ts index 6956b855333..dcfe530122a 100644 --- a/tests/web-search/web-search-sidecar-429.test.ts +++ b/tests/web-search/web-search-sidecar-429.test.ts @@ -31,14 +31,20 @@ describe("web-search sidecar 429 replays", () => { return new Response("data: [DONE]\n\n", { headers: { "content-type": "text/event-stream" } }); } - function searchWith(fetchImpl: () => Promise) { + function searchWith( + fetchImpl: () => Promise, + timeoutMs = 30_000, + recordOutcome?: (outcome: number | "connect_error" | "connect_neutral" | "timeout") => void, + ) { globalThis.fetch = fetchImpl as unknown as typeof fetch; return runOpenAiWebSearch( "current docs", { type: "web_search" }, sidecarProvider(), new Headers({ authorization: "Bearer selected-token" }), - { model: "gpt-5.6-luna", reasoning: "low", timeoutMs: 30_000 }, + { model: "gpt-5.6-luna", reasoning: "low", timeoutMs }, + undefined, + recordOutcome, ); } @@ -72,4 +78,42 @@ describe("web-search sidecar 429 replays", () => { expect(calls).toBe(1); expect(outcome.error).toContain("429"); }); + + test("a Retry-After that cannot fit the sidecar deadline preserves the 429", async () => { + let calls = 0; + const recorded: Array = []; + const outcome = await searchWith(async () => { + calls += 1; + return new Response("slow down", { status: 429, headers: { "retry-after": "0.1" } }); + }, 50, value => recorded.push(value)); + expect(calls).toBe(1); + expect(outcome.error).toContain("429"); + expect(recorded).toEqual([429]); + }); + + test("a deadline expiring during pre-retry body cleanup preserves the 429", async () => { + // The never-settling body is a worse leak than the other mocks leave behind: restore + // fetch so a later file's shared search loop does not inherit a 1s release per retry. + const originalFetch = globalThis.fetch; + try { + let calls = 0; + const recorded: Array = []; + const outcome = await searchWith(async () => { + calls += 1; + // A cancel() that never settles makes the bounded 1s release run to its cap; the + // remaining deadline then cannot fit the backoff, so the wait ends mid-sleep. The + // observed 429 must survive that expiry instead of being recorded as a timeout. + const body = new ReadableStream({ + start: controller => controller.enqueue(new TextEncoder().encode("rate limited")), + cancel: () => new Promise(() => {}), + }); + return new Response(body, { status: 429, headers: { "retry-after": "1" } }); + }, 1_500, value => recorded.push(value)); + expect(calls).toBe(1); + expect(outcome.error).toContain("429"); + expect(recorded).toEqual([429]); + } finally { + globalThis.fetch = originalFetch; + } + }); }); From 7e826dc0891fe308126c5284f2f535bdf9fe6125 Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Tue, 22 Sep 2026 06:24:00 +0900 Subject: [PATCH 02/12] fix(web-search): scope replay cells to key-resolved loopback principals (cherry picked from commit 7ffc8b858f4021ffe04f205c6bf24768bd6c1c75) --- src/server/responses/request-prepare.ts | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/src/server/responses/request-prepare.ts b/src/server/responses/request-prepare.ts index c9f04d73d73..e250a332bcb 100644 --- a/src/server/responses/request-prepare.ts +++ b/src/server/responses/request-prepare.ts @@ -20,7 +20,7 @@ import { sessionIdHeaderFromRequest, reasoningReplayConversationIdFromResponsesRequest, } from "../request-log-conversation"; -import { contextPrincipalIdOf } from "../auth-cors"; +import { resolveContextPrincipal } from "../auth-cors"; import { isShadowSourceModel, shadowSourceModelPrefix, @@ -409,7 +409,10 @@ export async function prepareResponsesRequest( } } if (parsed._reasoningReplayScope) { - const clientPrincipalId = contextPrincipalIdOf(options.admission) + // Scope replay cells to the caller principal. On loopback, admission carries no identity, + // so resolve it from an opencodex API key the caller volunteered (same rule as context + // history ownership); keyless loopback callers still share the "loopback" bucket. + const clientPrincipalId = resolveContextPrincipal(req, config, options.admission) ?? (options.admission?.kind === "loopback" ? "loopback" : undefined); parsed._reasoningReplayScope = { ...parsed._reasoningReplayScope, clientPrincipalId }; } From 7f45883fb59550f5a9cdee7416364bcee99b0901 Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Mon, 21 Sep 2026 03:57:44 +0900 Subject: [PATCH 03/12] fix(responses): isolate Cursor combo shadow calls (cherry picked from commit 3e5a4dd9bb6d7a1bc274e63dd569169ea0d972f1) --- src/server/responses/core-options.ts | 2 ++ src/server/responses/request-prepare.ts | 6 ++++++ .../responses-shadow-intercept.test.ts | 20 ++++++++++++++++++- 3 files changed, 27 insertions(+), 1 deletion(-) diff --git a/src/server/responses/core-options.ts b/src/server/responses/core-options.ts index 5334cf4e02d..110f56f7240 100644 --- a/src/server/responses/core-options.ts +++ b/src/server/responses/core-options.ts @@ -108,6 +108,8 @@ export interface HandleResponsesOptions { callerDirectAuth?: CallerDirectAuth | null; /** Internal recursion guard; callers outside this module must not set it. */ comboAttempt?: boolean; + /** Internal handoff: this combo was selected by shadow-call interception. */ + shadowCallIntercepted?: boolean; compactionRoutingOverride?: CompactionRoutingOverride | null; /** Internal combo handoff for one parent-validated continuation snapshot. */ comboReplaySnapshot?: { diff --git a/src/server/responses/request-prepare.ts b/src/server/responses/request-prepare.ts index e250a332bcb..ada3c423a14 100644 --- a/src/server/responses/request-prepare.ts +++ b/src/server/responses/request-prepare.ts @@ -226,6 +226,7 @@ export async function prepareResponsesRequest( // hops — which only exist inside that loop — are unreachable (#4129). Rewrite the selector // here instead, before comboIdFromRawBody reads `model`, and identify the combo by CONFIG // LOOKUP so the check can never observe a one-candidate collapse. + let shadowCallIntercepted = false; if (!options.comboAttempt && !options.compactionRoutingOverride && body && typeof body === "object" && !Array.isArray(body)) { const shadowIntercept = config.shadowCallIntercept; const rawShadowModel = (body as { model?: unknown }).model; @@ -233,6 +234,7 @@ export async function prepareResponsesRequest( && isShadowSourceModel(rawShadowModel, shadowIntercept.sourceModels)) { const shadowComboId = resolveComboId(config, shadowIntercept.model); if (shadowComboId && Object.hasOwn(config.combos ?? {}, shadowComboId)) { + shadowCallIntercepted = true; (body as Record).model = shadowIntercept.model; // Same rule as the late intercept site: record the operator-configured prefix that // matched, never the caller's raw model string. Matching is by prefix, so the raw @@ -248,6 +250,9 @@ export async function prepareResponsesRequest( options.onRequestBodyRead?.(); return requestDispatchers.handleComboResponses(req, body, comboId, config, logCtx, { ...options, + // Concrete combo child selectors no longer match the shadow source model. Carry the + // interception decision explicitly so provider-specific helper isolation still applies. + shadowCallIntercepted, // The original request body was accepted above. Combo children are synthetic // replays and must not repeat the caller-owned timeout transition. onRequestBodyRead: undefined, @@ -370,6 +375,7 @@ export async function prepareResponsesRequest( } } if (cursorClientThreadId) parsed._cursorClientThreadId = cursorClientThreadId; + if (options.shadowCallIntercepted === true) parsed._cursorIsolateConversation = true; } catch (err) { if (isTranslatorBudgetExceededError(err)) { return formatErrorResponse(413, "request_too_large", "request translation buffer exceeded the safe limit", { diff --git a/tests/responses/responses-shadow-intercept.test.ts b/tests/responses/responses-shadow-intercept.test.ts index d5cbdd680c4..ab3800136cf 100644 --- a/tests/responses/responses-shadow-intercept.test.ts +++ b/tests/responses/responses-shadow-intercept.test.ts @@ -4,7 +4,7 @@ * default follows modern clients, while sourceModels keeps an escape hatch. */ import { afterEach, describe, expect, test } from "bun:test"; -import { mkdtempSync} from "node:fs"; +import { mkdtempSync, readFileSync } from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; import { handleResponses, isShadowSourceModel } from "../../src/server/responses"; @@ -15,6 +15,7 @@ import type { OcxConfig } from "../../src/types"; import { catalogConvergenceFactory } from "../helpers/catalog-convergence"; import { removeTreeWithRetry } from "../helpers/remove-tree"; import { acquireOwnedSpendHome } from "../helpers/owned-spend-home"; +import { repoPath } from "../helpers/repo-root"; const originalFetch = globalThis.fetch; let releaseSpendHome: (() => void) | undefined; @@ -307,6 +308,23 @@ function chatOk(text: string): Response { } describe("a combo shadow-call target enters the failover loop (#4129)", () => { + test("carries helper conversation isolation into concrete combo children", () => { + const prepare = readFileSync(repoPath("src/server/responses/request-prepare.ts"), "utf8"); + const comboDispatch = prepare.slice( + prepare.indexOf("const comboId = !options.comboAttempt"), + prepare.indexOf("let unreadableEncryptedAgentTask"), + ); + const parsedHandoff = prepare.slice( + prepare.indexOf("if (cursorClientThreadId) parsed._cursorClientThreadId"), + prepare.indexOf("} catch (err)", prepare.indexOf("if (cursorClientThreadId) parsed._cursorClientThreadId")), + ); + + expect(comboDispatch).toContain("shadowCallIntercepted,"); + expect(parsedHandoff).toContain( + "if (options.shadowCallIntercepted === true) parsed._cursorIsolateConversation = true;", + ); + }); + test("a helper call rewritten to a combo hops past a 429 to the second target", async () => { takeSpendHome(); const urls: string[] = []; From c4fa8c8d8f972355179e6435c46256f8c61b7d22 Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Mon, 21 Sep 2026 03:58:15 +0900 Subject: [PATCH 04/12] fix(responses): repair terminal-less bridged search legs (cherry picked from commit 1fb6005335a9aee7f61b9dda1171ccb6b7e8272e) --- src/server/responses/passthrough-delivery.ts | 29 ++++++++++---------- tests/responses/passthrough-abort.test.ts | 9 +++++- 2 files changed, 23 insertions(+), 15 deletions(-) diff --git a/src/server/responses/passthrough-delivery.ts b/src/server/responses/passthrough-delivery.ts index b4941fbd6dd..7e56236ecc5 100644 --- a/src/server/responses/passthrough-delivery.ts +++ b/src/server/responses/passthrough-delivery.ts @@ -380,12 +380,22 @@ export async function deliverPassthroughResponse( }); // Capture the binding that actually served the first leg, after its permitted reselection. const webSearchBridgeBinding = requestBindings.get(nativeExchange.request); - // The bridge wraps the RAW upstream body, so terminal repair below still owns the single - // client-facing terminal — the bridge drops the terminal of every intercepted leg. - const upstreamSseBody = webSearchBridgePlan + // Repair must observe the raw first leg before the bridge suppresses an intercepted search + // lifecycle. Otherwise a provider that leaves that complete call open never arms repair's + // grace timer, so the bridge cannot execute the search or begin its continuation. + let passthroughSseBody = terminalRepairPolicy + ? relayResponsesSseWithTerminalRepair( + upstreamResponse.body, + upstream, + terminalRepairPolicy, + translatorBudget, + options.responsesTerminalRepairScheduler, + ) + : upstreamResponse.body; + passthroughSseBody = webSearchBridgePlan ? createPassthroughWebSearchBridgeStream({ plan: webSearchBridgePlan, - firstLeg: upstreamResponse.body, + firstLeg: passthroughSseBody, requestBody: nativeExchange.request.body, // Continuation legs replay the same built request with the executed search appended. // The first leg already passed the recovery ladder, the outbound size ceiling, and the @@ -432,16 +442,7 @@ export async function deliverPassthroughResponse( onFinalize: () => releaseCodexAuthContextProbeLease(openAiSidecar?.authContext), signal: upstream.signal, }) - : upstreamResponse.body; - const passthroughSseBody = terminalRepairPolicy - ? relayResponsesSseWithTerminalRepair( - upstreamSseBody, - upstream, - terminalRepairPolicy, - translatorBudget, - options.responsesTerminalRepairScheduler, - ) - : upstreamSseBody; + : passthroughSseBody; const repairConfig = route.provider.responsesItemIdRepair; // Grok Build renders deltas live but reconstructs its durable assistant // turn from the completed response snapshot. Native Responses streams diff --git a/tests/responses/passthrough-abort.test.ts b/tests/responses/passthrough-abort.test.ts index bbdd622d3db..cc50132307b 100644 --- a/tests/responses/passthrough-abort.test.ts +++ b/tests/responses/passthrough-abort.test.ts @@ -62,8 +62,15 @@ describe("passthrough relayWithAbort (RC2, passthrough path)", () => { // The captured static policy now supplies the repair decision; the real platform gate and // pure native relay invariants below are unchanged. expect(sseBranch).toContain("const terminalRepairPolicy = route.staticPolicy.model.responsesTerminalRepair;"); - expect(sseBranch).toContain("const passthroughSseBody = terminalRepairPolicy"); + expect(sseBranch).toContain("let passthroughSseBody = terminalRepairPolicy"); expect(sseBranch).toContain(": upstreamResponse.body;"); + // Repair has to wrap the raw first leg before the bridge hides its completed web-search call; + // otherwise a terminal-less open leg cannot trigger the repair timer and continuation stalls. + const terminalRepair = sseBranch.indexOf("relayResponsesSseWithTerminalRepair("); + const webSearchBridge = sseBranch.indexOf("createPassthroughWebSearchBridgeStream({"); + expect(terminalRepair).toBeGreaterThanOrEqual(0); + expect(webSearchBridge).toBeGreaterThan(terminalRepair); + expect(sseBranch.slice(webSearchBridge)).toContain("firstLeg: passthroughSseBody,"); // Native tee stays inside the bounded observer. The production owner passes // the raw stream and disconnect signal before any client-side rewrite. expect(sseBranch).toMatch(/const \[nativeBody, inspectBody\] = teeWithBoundedInspection\(passthroughSseBody, \{ clientGoneSignal \}\)/); From 8d46989165059c163eda0cb1df23f33c5b8dcb8e Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Mon, 21 Sep 2026 05:43:15 +0900 Subject: [PATCH 05/12] fix(responses): repair continuation legs and prove behavior in tests (cherry picked from commit 3948970a6459b803e5b1be4604fdeed78566f74f) --- src/server/responses/passthrough-delivery.ts | 56 +++--- tests/responses/passthrough-abort.test.ts | 166 ++++++++++++++++++ .../responses-shadow-intercept.test.ts | 40 +++++ 3 files changed, 242 insertions(+), 20 deletions(-) diff --git a/src/server/responses/passthrough-delivery.ts b/src/server/responses/passthrough-delivery.ts index 7e56236ecc5..8d7972cd5c4 100644 --- a/src/server/responses/passthrough-delivery.ts +++ b/src/server/responses/passthrough-delivery.ts @@ -400,26 +400,42 @@ export async function deliverPassthroughResponse( // Continuation legs replay the same built request with the executed search appended. // The first leg already passed the recovery ladder, the outbound size ceiling, and the // host circuit; a KEY-auth destination has no OAuth refresh to replay on a later leg. - send: (continuationBody: string) => fetchWithHeaderTimeout( - nativeExchange.request.url, - { method: nativeExchange.request.method, headers: nativeExchange.request.headers, body: continuationBody }, - upstream.signal, - connectMs, - true, - providerFetch(route.provider, options.codexWsRuntimeIdentity, { - // Pacing can outlive a manual selection change. A continuation must retain the - // first leg's key and appended search result, never rebuild from the original turn. - beforeDispatch: () => { - if (webSearchBridgeBinding?.kind !== "api-key" - || !providerApiKeySelectionIsCurrent(config, route.providerName, webSearchBridgeBinding.provider)) { - throw new Error("API key selection changed during a web-search continuation"); - } - }, - providerName: route.providerName, - modelId: route.modelId, - }), - false, - ), + send: async (continuationBody: string) => { + const continuation = await fetchWithHeaderTimeout( + nativeExchange.request.url, + { method: nativeExchange.request.method, headers: nativeExchange.request.headers, body: continuationBody }, + upstream.signal, + connectMs, + true, + providerFetch(route.provider, options.codexWsRuntimeIdentity, { + // Pacing can outlive a manual selection change. A continuation must retain the + // first leg's key and appended search result, never rebuild from the original turn. + beforeDispatch: () => { + if (webSearchBridgeBinding?.kind !== "api-key" + || !providerApiKeySelectionIsCurrent(config, route.providerName, webSearchBridgeBinding.provider)) { + throw new Error("API key selection changed during a web-search continuation"); + } + }, + providerName: route.providerName, + modelId: route.modelId, + }), + false, + ); + // The same provider can leave a complete continuation open without a terminal, which + // stalls the bridge's decide loop exactly like the first leg — so every leg gets the + // same repair, not only the intercepted first one. + if (!terminalRepairPolicy || !continuation.ok || !continuation.body) return continuation; + return new Response( + relayResponsesSseWithTerminalRepair( + continuation.body, + upstream, + terminalRepairPolicy, + translatorBudget, + options.responsesTerminalRepairScheduler, + ), + continuation, + ); + }, execute: createPassthroughWebSearchBridgeExecutor(webSearchBridgePlan, { providerApiKey: route.provider.apiKey ?? "", auth: webSearchBridgeAuth, diff --git a/tests/responses/passthrough-abort.test.ts b/tests/responses/passthrough-abort.test.ts index cc50132307b..74cfcddca73 100644 --- a/tests/responses/passthrough-abort.test.ts +++ b/tests/responses/passthrough-abort.test.ts @@ -2,6 +2,9 @@ import { describe, expect, test } from "bun:test"; import { consumeForInspection, linkAbortSignal, relaySseWithFailedTail, relaySseWithHeartbeat, relayWithAbort } from "../../src/server"; import { pathToFileURL } from "node:url"; import { repoRoot } from "../helpers/repo-root"; +import { relayResponsesSseWithTerminalRepair, type ResponsesTerminalRepairScheduler } from "../../src/server/responses-terminal-repair"; +import { createPassthroughWebSearchBridgeStream, type PassthroughWebSearchBridgePlan } from "../../src/web-search/passthrough-bridge"; +import { createTestTranslatorBudget } from "../helpers/translator-budget"; const root = pathToFileURL(repoRoot() + "/"); @@ -627,3 +630,166 @@ describe("passthrough relayWithAbort (RC2, passthrough path)", () => { expect(upstream.signal.reason).toBe("replacement turn"); }); }); + +/** + * The reported stall: a provider emits a complete intercepted `web_search` call but never sends + * a terminal and holds the leg open. With repair wrapped around the raw first leg, the grace + * timer still arms and the bridge can execute the search and continue upstream. + */ +describe("terminal repair ahead of the passthrough web-search bridge", () => { + class ManualScheduler implements ResponsesTerminalRepairScheduler { + private current = 0; + private nextId = 1; + private readonly jobs = new Map void }>(); + + nowMs(): number { return this.current; } + + schedule(callback: () => void, delayMs: number): unknown { + const id = this.nextId++; + this.jobs.set(id, { at: this.current + delayMs, callback }); + return id; + } + + cancel(handle: unknown): void { + this.jobs.delete(handle as number); + } + + advance(ms: number): void { + this.current += ms; + for (;;) { + const due = [...this.jobs.entries()] + .filter(([, job]) => job.at <= this.current) + .sort((left, right) => left[1].at - right[1].at); + if (due.length === 0) return; + for (const [id, job] of due) { + if (!this.jobs.delete(id)) continue; + job.callback(); + } + } + } + + pending(): number { return this.jobs.size; } + } + + const searchCall = { + type: "function_call", + id: "fc_1", + status: "completed", + call_id: "call_1", + name: "web_search", + arguments: "{\"query\":\"opencodex release\"}", + }; + + const preamble = { + type: "message", + id: "msg_1", + status: "completed", + role: "assistant", + content: [{ type: "output_text", text: "Let me look that up." }], + }; + + const answer = { + type: "message", + id: "msg_2", + role: "assistant", + content: [{ type: "output_text", text: "The current release is 2.50.0." }], + }; + + function frame(type: string, payload: Record): string { + return "event: " + type + "\ndata: " + JSON.stringify({ type, ...payload }); + } + + function sseBody(...blocks: string[]): string { + return blocks.concat("data: [DONE]").join("\n\n") + "\n\n"; + } + + /** Every output item complete, no terminal event, and the leg is never closed. */ + function terminallessSearchLeg(): ReadableStream { + const text = [ + frame("response.created", { response: { id: "resp_1", status: "in_progress" } }), + frame("response.output_item.added", { output_index: 0, item: { ...preamble, content: [] } }), + frame("response.output_item.done", { output_index: 0, item: preamble }), + frame("response.output_item.added", { output_index: 1, item: { ...searchCall, arguments: "" } }), + frame("response.function_call_arguments.done", { + output_index: 1, + item_id: "fc_1", + arguments: searchCall.arguments, + }), + frame("response.output_item.done", { output_index: 1, item: searchCall }), + ].join("\n\n") + "\n\n"; + return new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode(text)); + }, + }); + } + + function answerLeg(): ReadableStream { + return streamFromChunks([new TextEncoder().encode(sseBody( + frame("response.created", { response: { id: "resp_2", status: "in_progress" } }), + frame("response.output_item.added", { output_index: 0, item: { ...answer, content: [] } }), + frame("response.output_item.done", { output_index: 0, item: answer }), + frame("response.completed", { + response: { id: "resp_2", status: "completed", output: [answer] }, + }), + ))]); + } + + test("a terminal-less first search leg reaches the bridge once repaired", async () => { + const scheduler = new ManualScheduler(); + const upstream = new AbortController(); + const plan: PassthroughWebSearchBridgePlan = { + backend: "ollama", + endpoint: "https://ollama.com/api/web_search", + maxSearches: 3, + timeoutMs: 60_000, + }; + const sent: string[] = []; + const executed: string[][] = []; + const stream = createPassthroughWebSearchBridgeStream({ + plan, + firstLeg: relayResponsesSseWithTerminalRepair( + terminallessSearchLeg(), + upstream, + { graceMs: 5_000 }, + createTestTranslatorBudget(), + scheduler, + ), + requestBody: JSON.stringify({ + model: "glm-4.7", + stream: true, + input: [{ role: "user", content: [{ type: "input_text", text: "what is the latest release?" }] }], + tools: [{ type: "web_search" }], + }), + send: async (body) => { + sent.push(body); + return new Response(answerLeg(), { + headers: { "content-type": "text/event-stream" }, + }); + }, + execute: async (queries) => { + executed.push(queries); + return { text: "opencodex 2.50.0 shipped", sources: [] }; + }, + }); + + const bodyPromise = new Response(stream).text(); + // The bridge is pull-driven: let it drain the pushed leg frames so repair arms the timer. + for (let i = 0; i < 1_000 && scheduler.pending() === 0; i += 1) { + await new Promise(resolve => setTimeout(resolve, 0)); + } + expect(scheduler.pending()).toBe(1); + scheduler.advance(5_000); + const body = await bodyPromise; + + expect(executed).toEqual([["opencodex release"]]); + expect(sent).toHaveLength(1); + const events = body + .split(/\r?\n/) + .filter(line => line.startsWith("data:")) + .map(line => line.slice(5).trim()) + .filter(payload => payload.length > 0 && payload !== "[DONE]") + .map(payload => JSON.parse(payload) as Record); + expect(events.some(event => event.type === "response.completed")).toBe(true); + }); +}); diff --git a/tests/responses/responses-shadow-intercept.test.ts b/tests/responses/responses-shadow-intercept.test.ts index ab3800136cf..fd5befeb7d6 100644 --- a/tests/responses/responses-shadow-intercept.test.ts +++ b/tests/responses/responses-shadow-intercept.test.ts @@ -16,6 +16,8 @@ import { catalogConvergenceFactory } from "../helpers/catalog-convergence"; import { removeTreeWithRetry } from "../helpers/remove-tree"; import { acquireOwnedSpendHome } from "../helpers/owned-spend-home"; import { repoPath } from "../helpers/repo-root"; +import { createTestTranslatorBudget } from "../helpers/translator-budget"; +import { prepareResponsesRequest } from "../../src/server/responses/request-prepare"; const originalFetch = globalThis.fetch; let releaseSpendHome: (() => void) | undefined; @@ -308,6 +310,44 @@ function chatOk(text: string): Response { } describe("a combo shadow-call target enters the failover loop (#4129)", () => { + test("a combo child of a shadow-intercepted call gets Cursor conversation isolation", async () => { + const config = comboInterceptConfig([{ provider: "xai", model: "grok-4.5" }]); + const logCtx: RequestLogContext = { model: "", provider: "" }; + const mkreq = () => new Request("http://localhost/v1/responses", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ + model: "grok-4.5", + input: [{ type: "message", role: "user", content: [{ type: "input_text", text: "hi" }] }], + stream: false, + }), + }); + const dispatchers = { + handleResponses: () => Promise.reject(new Error("unexpected recursion")), + handleComboResponses: () => Promise.reject(new Error("unexpected combo dispatch")), + }; + const admission = () => ({ pendingHostAdmissionLease: null, authCtx: { kind: "main", accountId: null } }) as never; + + const intercepted = await prepareResponsesRequest( + { req: mkreq(), config, logCtx, options: { comboAttempt: true, shadowCallIntercepted: true, translatorBudget: createTestTranslatorBudget() } }, + admission(), + dispatchers, + ); + expect(intercepted).not.toBeInstanceOf(Response); + if (intercepted instanceof Response) throw new Error("expected a prepared request, got HTTP " + intercepted.status); + expect(intercepted.parsed._cursorIsolateConversation).toBe(true); + + // A plain combo child (no interception marker) must not be isolated. + const plain = await prepareResponsesRequest( + { req: mkreq(), config, logCtx, options: { comboAttempt: true, translatorBudget: createTestTranslatorBudget() } }, + admission(), + dispatchers, + ); + expect(plain).not.toBeInstanceOf(Response); + if (plain instanceof Response) throw new Error("expected a prepared request, got HTTP " + plain.status); + expect(plain.parsed._cursorIsolateConversation).not.toBe(true); + }); + test("carries helper conversation isolation into concrete combo children", () => { const prepare = readFileSync(repoPath("src/server/responses/request-prepare.ts"), "utf8"); const comboDispatch = prepare.slice( From 3f3fdf17f4f480294a3cafaebc2b99d52016d214 Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Mon, 21 Sep 2026 11:11:56 +0900 Subject: [PATCH 06/12] fix(responses): close shadow combo intersection and continuation coverage Apply the existing source-target non-intersection rule before early combo interception, and add production-path behavioral coverage for terminal-less continuation repair. (cherry picked from commit 4bfcc0a8ea657b4f659355385a87c3b11a542365) --- src/server/responses/request-prepare.ts | 23 ++- tests/responses/passthrough-abort.test.ts | 161 ++++++++++++++++++ .../responses-shadow-intercept.test.ts | 12 +- 3 files changed, 181 insertions(+), 15 deletions(-) diff --git a/src/server/responses/request-prepare.ts b/src/server/responses/request-prepare.ts index ada3c423a14..53931fbf38c 100644 --- a/src/server/responses/request-prepare.ts +++ b/src/server/responses/request-prepare.ts @@ -234,14 +234,21 @@ export async function prepareResponsesRequest( && isShadowSourceModel(rawShadowModel, shadowIntercept.sourceModels)) { const shadowComboId = resolveComboId(config, shadowIntercept.model); if (shadowComboId && Object.hasOwn(config.combos ?? {}, shadowComboId)) { - shadowCallIntercepted = true; - (body as Record).model = shadowIntercept.model; - // Same rule as the late intercept site: record the operator-configured prefix that - // matched, never the caller's raw model string. Matching is by prefix, so the raw - // value is caller-controlled and reaches usage.jsonl and /api/logs. - logCtx.shadowCallRewrittenFrom = sanitizeLogMetadataString( - shadowSourceModelPrefix(rawShadowModel, shadowIntercept.sourceModels), - ); + const sourcePrefix = shadowSourceModelPrefix(rawShadowModel, shadowIntercept.sourceModels)!; + let sourceIdentity = { providerName: OPENAI_CODEX_PROVIDER_ID, modelId: sourcePrefix }; + try { + const resolvedSource = routeConcreteModel(config, rawShadowModel); + sourceIdentity = { providerName: resolvedSource.providerName, modelId: sourcePrefix }; + } catch { /* Native Codex helper calls remain OpenAI-owned without an enabled OpenAI route. */ } + const targetRoute = routeModel(config, shadowIntercept.model, evidenceFromBody(body)); + if (shouldInterceptShadowCall(rawShadowModel, shadowIntercept.sourceModels, sourceIdentity, targetRoute)) { + shadowCallIntercepted = true; + (body as Record).model = shadowIntercept.model; + // Same rule as the late intercept site: record the operator-configured prefix that + // matched, never the caller's raw model string. Matching is by prefix, so the raw + // value is caller-controlled and reaches usage.jsonl and /api/logs. + logCtx.shadowCallRewrittenFrom = sanitizeLogMetadataString(sourcePrefix); + } } } } diff --git a/tests/responses/passthrough-abort.test.ts b/tests/responses/passthrough-abort.test.ts index 74cfcddca73..5d5818bd6c3 100644 --- a/tests/responses/passthrough-abort.test.ts +++ b/tests/responses/passthrough-abort.test.ts @@ -4,6 +4,8 @@ import { pathToFileURL } from "node:url"; import { repoRoot } from "../helpers/repo-root"; import { relayResponsesSseWithTerminalRepair, type ResponsesTerminalRepairScheduler } from "../../src/server/responses-terminal-repair"; import { createPassthroughWebSearchBridgeStream, type PassthroughWebSearchBridgePlan } from "../../src/web-search/passthrough-bridge"; +import { deliverPassthroughResponse } from "../../src/server/responses/passthrough-delivery"; +import { routedProviderConfig } from "../../src/router"; import { createTestTranslatorBudget } from "../helpers/translator-budget"; const root = pathToFileURL(repoRoot() + "/"); @@ -735,6 +737,20 @@ describe("terminal repair ahead of the passthrough web-search bridge", () => { ))]); } + /** Complete answer output, no terminal event, and the continuation remains open. */ + function terminallessAnswerLeg(): ReadableStream { + const text = [ + frame("response.created", { response: { id: "resp_2", status: "in_progress" } }), + frame("response.output_item.added", { output_index: 0, item: { ...answer, content: [] } }), + frame("response.output_item.done", { output_index: 0, item: { ...answer, status: "completed" } }), + ].join("\n\n") + "\n\n"; + return new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode(text)); + }, + }); + } + test("a terminal-less first search leg reaches the bridge once repaired", async () => { const scheduler = new ManualScheduler(); const upstream = new AbortController(); @@ -792,4 +808,149 @@ describe("terminal repair ahead of the passthrough web-search bridge", () => { .map(payload => JSON.parse(payload) as Record); expect(events.some(event => event.type === "response.completed")).toBe(true); }); + + /** + * The production continuation sender in deliverPassthroughResponse must apply the same + * terminal repair to every leg, not only the first one. This drives the real function: + * the first leg is a terminal-less intercepted web_search call, the ollama search fetch is + * stubbed, and the provider's own fetch returns a terminal-less continuation — which only + * reaches the client when the sender's repair wrap synthesizes response.completed. + */ + test("deliverPassthroughResponse repairs a terminal-less continuation leg", async () => { + const scheduler = new ManualScheduler(); + const upstream = new AbortController(); + const originalFetch = globalThis.fetch; + + const provider = routedProviderConfig("bridge-test", { + adapter: "openai-responses", + baseUrl: "https://bridge-test.example/v1", + authMode: "key", + apiKey: "test-bridge-key", + webSearchBridge: { + enabled: true, + backend: "ollama", + endpoint: "https://bridge-test.example/api/web_search", + maxSearches: 3, + timeoutMs: 60_000, + }, + fetch: (async () => new Response(terminallessAnswerLeg(), { + headers: { "content-type": "text/event-stream" }, + })) as unknown as typeof globalThis.fetch, + } as never); + const config = { + providers: { "bridge-test": provider }, + maxUpstreamBodyBytes: 8 * 1024 * 1024, + }; + const upstreamRequest = { + url: "https://bridge-test.example/v1/responses", + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ + model: "glm-4.7", + stream: true, + input: [{ role: "user", content: [{ type: "input_text", text: "what is the latest release?" }] }], + tools: [{ type: "web_search" }], + }), + }; + const requestBindings = new WeakMap(); + requestBindings.set(upstreamRequest, { kind: "api-key", provider }); + + globalThis.fetch = (async () => Response.json({ + results: [{ url: "https://example.com/release", title: "Release notes", content: "2.50.0 shipped" }], + })) as typeof globalThis.fetch; + try { + const response = await deliverPassthroughResponse( + { + logCtx: { model: "", provider: "" }, + config, + options: { responsesTerminalRepairScheduler: scheduler }, + req: new Request("http://localhost/v1/responses", { method: "POST" }), + }, + { authCtx: { kind: "main", accountId: null } }, + { + parsed: { + modelId: "glm-4.7", + stream: true, + options: {}, + _webSearch: { type: "web_search" }, + }, + route: { + providerName: "bridge-test", + provider, + modelId: "glm-4.7", + staticPolicy: { model: { responsesTerminalRepair: { graceMs: 5_000 } } }, + }, + subagentQuotaFailureModel: undefined, + subagentFallbackAccountId: undefined, + clientRequestedStream: true, + translatorBudget: createTestTranslatorBudget(), + }, + { requestBindings }, + { openAiSidecar: undefined }, + { + plaintextV2AgentMessageToolNames: new Set(), + commitReasoningReplayServingRoute: () => {}, + routedMuseToolNameAliases: new Map(), + routedNamespaceToolAliases: new Map(), + plaintextV2AgentMessageAliasedToolNames: new Set(), + recordTerminalOutcomes: false, + responseCompletionCancelled: false, + }, + { + upstreamResponse: new Response(terminallessSearchLeg(), { + headers: { "content-type": "text/event-stream" }, + }), + codexSafetyBufferingOptions: undefined, + upstream, + request: upstreamRequest, + connectMs: 5_000, + imageGenCallAliases: new Map(), + selfNamedNamespaceScrubAuthorization: undefined, + authorizedBareNamespaceToolAliases: new Map(), + rememberPassthroughResponseChecked: () => {}, + routedCustomToolNames: new Set(), + routedCustomToolRepairNames: new Set(), + declaredWireToolNames: new Set(), + routedToolSearchNames: new Set(), + outboundRequestBody: undefined, + functionRepairSchemas: new Map(), + undeclaredToolGuardActive: false, + declaredNamelessClientCallTypes: new Set(), + providerExecutedCallTypes: new Set(), + declaredBareWireToolNames: new Set(), + rememberPassthroughResponse: false, + noteInspectedPayload: () => {}, + normalizeFunctionCompletionJson: (text: string) => text, + }, + ); + + expect(response.ok).toBe(true); + const bodyPromise = response.text(); + // First leg: repair arms once every output item is complete and the grace timer + // synthesizes the terminal that lets the bridge dispatch its continuation. + for (let i = 0; i < 1_000 && scheduler.pending() === 0; i += 1) { + await new Promise(resolve => setTimeout(resolve, 0)); + } + expect(scheduler.pending()).toBe(1); + scheduler.advance(5_000); + // Continuation leg: the production sender wraps the fetch result in the same repair, + // so its own terminal-less body re-arms the timer instead of stalling the stream. + for (let i = 0; i < 1_000 && scheduler.pending() === 0; i += 1) { + await new Promise(resolve => setTimeout(resolve, 0)); + } + expect(scheduler.pending()).toBe(1); + scheduler.advance(5_000); + const body = await bodyPromise; + + const events = body + .split(/\r?\n/) + .filter(line => line.startsWith("data:")) + .map(line => line.slice(5).trim()) + .filter(payload => payload.length > 0 && payload !== "[DONE]") + .map(payload => JSON.parse(payload) as Record); + expect(events.some(event => event.type === "response.completed")).toBe(true); + } finally { + globalThis.fetch = originalFetch; + } + }); }); diff --git a/tests/responses/responses-shadow-intercept.test.ts b/tests/responses/responses-shadow-intercept.test.ts index fd5befeb7d6..fc07e8bb8a5 100644 --- a/tests/responses/responses-shadow-intercept.test.ts +++ b/tests/responses/responses-shadow-intercept.test.ts @@ -397,7 +397,7 @@ describe("a combo shadow-call target enters the failover loop (#4129)", () => { .toEqual(["xai/grok-4.5", "alt/grok-4.5"]); }); - test("a combo whose first target intersects the source still routes as a combo", async () => { + test("a combo whose selected target intersects the source is not intercepted", async () => { takeSpendHome(); const urls: string[] = []; const logCtx: RequestLogContext = { model: "", provider: "" }; @@ -423,12 +423,10 @@ describe("a combo shadow-call target enters the failover loop (#4129)", () => { // A healthy first target still costs exactly one upstream call. expect(urls).toHaveLength(1); expect(urls[0]).toContain("api.x.ai"); - expect(logCtx.provider).toBe("combo"); - expect(logCtx.comboId).toBe("shadow"); - expect(logCtx.routeDecision?.routeKind).toBe("combo"); - // Red before the fix: shouldInterceptShadowCall saw the collapsed pick as a self-target, - // skipped the rewrite, and the request left as a plain native route with no marker. - expect(logCtx.shadowCallRewrittenFrom).toBe("custom-helper"); + expect(logCtx.provider).toBe("xai"); + expect(logCtx.comboId).toBeUndefined(); + expect(logCtx.routeDecision?.routeKind).not.toBe("combo"); + expect(logCtx.shadowCallRewrittenFrom).toBeUndefined(); }); test("a non-combo replacement still takes the ordinary late intercept", async () => { From 973a4ac70281add989e7b6643a08463f25baacce Mon Sep 17 00:00:00 2001 From: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Mon, 21 Sep 2026 05:45:55 +0000 Subject: [PATCH 07/12] test(web-search): preserve complete but open bridge-leg coverage Carry the end-to-end handleResponses regression and transport contracts for repairing both the first and continuation search legs. The corresponding production changes are already preserved by the earlier terminal-repair carries; keep this broader integration coverage without applying that implementation twice. Source commit: b1044e7b22520a9fac0051477ea18bed91454823 Co-authored-by: Epinephrine --- structure/providers-and-adapters.md | 5 +- structure/transports/streaming-health.md | 8 + .../web-search-passthrough-bridge.test.ts | 171 ++++++++++++++++++ 3 files changed, 183 insertions(+), 1 deletion(-) diff --git a/structure/providers-and-adapters.md b/structure/providers-and-adapters.md index 2256b0ce1ba..536b348d78e 100644 --- a/structure/providers-and-adapters.md +++ b/structure/providers-and-adapters.md @@ -164,7 +164,10 @@ the configured entry, reference, revision, resolved key, authentication mode, an disabled or removed provider fails the same check. Drift produces the bridge's failed terminal without another provider request, and an unchanged binding resends the built request with its executed search result appended, never re-entering the initial reselection/rebuild path. Initial -dispatch keeps its normal reselection policy. `tests/web-search/web-search-passthrough-bridge.test.ts` +dispatch keeps its normal reselection policy. When the route's registry policy carries a +terminal-repair grace (`modelResponsesTerminalRepair`), the response body of every successful +continuation is wrapped by the same repair that saw the raw first leg, so a complete leg the +destination leaves open still ends that leg on schedule instead of stalling the turn. `tests/web-search/web-search-passthrough-bridge.test.ts` covers drift during search, while pacing, and before first-leg headers return, plus successful first-dispatch reselection and result preservation. diff --git a/structure/transports/streaming-health.md b/structure/transports/streaming-health.md index 9246b1e9695..e6b9d393f38 100644 --- a/structure/transports/streaming-health.md +++ b/structure/transports/streaming-health.md @@ -248,6 +248,14 @@ as `response.incomplete`, never synthetic success. The repair shares the per-tur budget, preserves backpressure, and composes ahead of item-id/snapshot rewrites so HTTP/SSE and WebSocket clients observe the same canonical lifecycle. +When the hosted-search bridge is also armed, repair wraps the raw first leg BEFORE the bridge: +the bridge suppresses an intercepted `web_search` lifecycle, so a complete call whose leg never +closes would otherwise leave the grace timer unarmed and the turn stalled. The same wrap applies +to every continuation leg the bridge's `send` returns — each leg gets its own grace window on the +shared abort controller — so a terminal-less continuation cannot stall the bridged turn either. +`tests/web-search/web-search-passthrough-bridge.test.ts` drives both legs through `handleResponses` +with an injected scheduler and proves search execution, continuation dispatch, and final terminal. + `ws-bridge.ts` preserves upstream `failed` and `incomplete` status values in the final WebSocket frame rather than always emitting `response.completed`. If the response status is `failed`, a `response.failed` frame is sent; otherwise `response.completed` carries through the original status. diff --git a/tests/web-search/web-search-passthrough-bridge.test.ts b/tests/web-search/web-search-passthrough-bridge.test.ts index 05d06f4332a..8777c183247 100644 --- a/tests/web-search/web-search-passthrough-bridge.test.ts +++ b/tests/web-search/web-search-passthrough-bridge.test.ts @@ -26,6 +26,9 @@ import { providerWebSearchBridgeConfigError, validateConfigCandidate } from "../ import { mapOllamaSearchResponse } from "../../src/web-search/ollama-executor"; import { UNDECLARED_TOOL_CALL_ERROR_CODE } from "../../src/server/responses-undeclared-tool-guard"; import { handleResponses } from "../../src/server/responses"; +import { providerConfigSeed } from "../../src/providers/derive"; +import { getProviderRegistryEntry } from "../../src/providers/registry"; +import type { ResponsesTerminalRepairScheduler } from "../../src/server/responses-terminal-repair"; import { resetProviderRequestPacingForTest, setProviderRequestPacingRuntimeForTest, @@ -1469,6 +1472,174 @@ describe("the reported turn, end to end through handleResponses", () => { item.type === "function_call" && item.name === "web_search")).toBe(true); }); + test("a complete but terminal-less leg still repairs, on the first leg AND the continuation", async () => { + // Repair is registry-gated, so only a registry-keyed provider arms it: deepseek carries + // modelResponsesTerminalRepair for the V4 flash ids. The fixture legs below emit a fully + // complete item lifecycle and then stay open — the reported stall — with no terminal and + // no [DONE]. Before the fix the repaired first leg could fire the search, but the raw + // continuation leg never got a grace window, so the turn still hung. + class ManualScheduler implements ResponsesTerminalRepairScheduler { + private current = 0; + private nextId = 1; + private readonly jobs = new Map void }>(); + nowMs(): number { return this.current; } + schedule(callback: () => void, delayMs: number): unknown { + const id = this.nextId++; + this.jobs.set(id, { at: this.current + delayMs, callback }); + return id; + } + cancel(handle: unknown): void { this.jobs.delete(handle as number); } + pending(): number { return this.jobs.size; } + advance(ms: number): void { + this.current += ms; + for (const [id, job] of [...this.jobs.entries()]) { + if (job.at > this.current || !this.jobs.delete(id)) continue; + job.callback(); + } + } + } + + const openSse = (): { stream: ReadableStream; push: (text: string) => void; end: () => void } => { + const encoder = new TextEncoder(); + let controller: ReadableStreamDefaultController | null = null; + return { + stream: new ReadableStream({ start(next) { controller = next; } }), + push(text) { controller?.enqueue(encoder.encode(text)); }, + end() { try { controller?.close(); } catch { /* already closed */ } }, + }; + }; + + // Every item must reach a COMPLETE output_item.done or repair never arms — the status + // field is what isCompleteItem actually requires. + const donePreamble = { ...preamble, status: "completed" }; + const doneSearchCall = { ...searchCall, status: "completed" }; + const doneAnswer = { ...answer, status: "completed" }; + const blocks = (...frames: string[]): string => frames.join("\n\n") + "\n\n"; + const openSearchLeg = blocks( + frame("response.created", { response: { id: "resp_1", status: "in_progress" } }), + frame("response.output_item.added", { output_index: 0, item: { ...donePreamble, content: [] } }), + frame("response.output_item.done", { output_index: 0, item: donePreamble }), + frame("response.output_item.added", { output_index: 1, item: { ...doneSearchCall, arguments: "" } }), + frame("response.function_call_arguments.done", { + output_index: 1, item_id: "fc_1", arguments: searchCall.arguments, + }), + frame("response.output_item.done", { output_index: 1, item: doneSearchCall }), + ); + const openAnswerLeg = blocks( + frame("response.created", { response: { id: "resp_2", status: "in_progress" } }), + frame("response.output_item.added", { output_index: 0, item: { ...doneAnswer, content: [] } }), + frame("response.output_item.done", { output_index: 0, item: doneAnswer }), + ); + + const firstLeg = openSse(); + const continuationLeg = openSse(); + const scheduler = new ManualScheduler(); + const outbound: string[] = []; + let searches = 0; + const savedFetch = globalThis.fetch; + globalThis.fetch = (async (input: unknown, init?: RequestInit) => { + const url = typeof input === "string" + ? input + : input instanceof URL ? input.href : (input as Request).url; + if (url.includes("api.exa.ai/search")) { + searches += 1; + return new Response(JSON.stringify({ + results: [{ title: "Releases", url: "https://example.test/rel", content: "opencodex 2.50.0", text: "opencodex 2.50.0" }], + }), { headers: { "content-type": "application/json" } }); + } + outbound.push(String(init?.body ?? "")); + return new Response(outbound.length === 1 ? firstLeg.stream : continuationLeg.stream, { + headers: { "content-type": "text/event-stream" }, + }); + }) as unknown as typeof fetch; + const cfg = { + port: 0, + defaultProvider: "deepseek", + providers: { + deepseek: { + ...providerConfigSeed(getProviderRegistryEntry("deepseek")!), + apiKey: "fixture-key", + webSearchBridge: { enabled: true, backend: "exa" }, + }, + }, + webSearchSidecar: { exaApiKey: "exa-canary" }, + } as unknown as OcxConfig; + const releaseSpendHome = acquireOwnedSpendHome(); + const decoder = new TextDecoder(); + const readUntil = async (reader: ReadableStreamDefaultReader, pattern: string): Promise => { + let out = ""; + while (!out.includes(pattern)) { + const { done, value } = await reader.read(); + if (done) throw new Error(`stream closed before ${pattern}`); + out += decoder.decode(value, { stream: true }); + } + return out; + }; + const flush = async (condition: () => boolean): Promise => { + for (let attempts = 0; attempts < 50 && !condition(); attempts += 1) await Bun.sleep(0); + }; + try { + const response = await handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", + headers: { "content-type": "application/json", authorization: "Bearer caller-inbound" }, + body: JSON.stringify({ + model: "deepseek/deepseek-v4-flash", + stream: true, + input: [{ role: "user", content: [{ type: "input_text", text: "what is the latest release?" }] }], + tools: [{ type: "web_search" }], + }), + }), cfg, { model: "", provider: "" }, { + responsesTerminalRepairScheduler: scheduler, + }); + const reader = response.body!.getReader(); + try { + // First leg: the complete search lifecycle streams through while the leg stays open. + firstLeg.push(openSearchLeg); + const opened = await readUntil(reader, "web_search_call"); + expect(opened).toContain("\"type\":\"web_search_call\""); + await flush(() => scheduler.pending() === 1); + expect(scheduler.pending()).toBe(1); + // The grace window is what ends the leg — before it fires, no search may run. + expect(searches).toBe(0); + scheduler.advance(5_000); + await flush(() => searches === 1 && outbound.length === 2); + expect(searches).toBe(1); + expect(outbound).toHaveLength(2); + const continued = JSON.parse(outbound[1]!) as { input: Record[] }; + expect(continued.input.some(item => item.type === "function_call_output" + && String(item.output).includes("opencodex 2.50.0"))).toBe(true); + + // Continuation leg: a complete answer that also never sends its terminal. Without + // repair on send() this is where the turn hangs. + continuationLeg.push(openAnswerLeg); + await flush(() => scheduler.pending() === 1); + expect(scheduler.pending()).toBe(1); + scheduler.advance(5_000); + const rest = await Promise.race([ + (async () => { + let out = ""; + for (;;) { + const { done, value } = await reader.read(); + if (done) return out + decoder.decode(); + out += decoder.decode(value, { stream: true }); + } + })(), + new Promise((_, reject) => setTimeout(() => reject(new Error("continuation never repaired")), 5_000)), + ]); + expect(rest).toContain("response.completed"); + expect(rest).toContain("The current release is 2.50.0."); + expect(rest).toContain("[DONE]"); + } finally { + try { await reader.cancel(); } catch { /* already closed */ } + firstLeg.end(); + continuationLeg.end(); + } + } finally { + releaseSpendHome(); + globalThis.fetch = savedFetch; + } + }); + const selectionChanges: Array<[string, (ocxConfig: OcxConfig) => void]> = [ ["selection revision with an unchanged key", cfg => { cfg.providers.fixture!.apiKeySelectionRevision = "selection-after"; From bb49c9f5821cc5fb44d08a78524633cb14378486 Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Tue, 22 Sep 2026 22:58:42 +0900 Subject: [PATCH 08/12] test(web-search): keep repaired replay within caller and serving scope --- .../web-search-passthrough-bridge.test.ts | 29 ++++++++++++++++++- 1 file changed, 28 insertions(+), 1 deletion(-) diff --git a/tests/web-search/web-search-passthrough-bridge.test.ts b/tests/web-search/web-search-passthrough-bridge.test.ts index 8777c183247..e1b057596cb 100644 --- a/tests/web-search/web-search-passthrough-bridge.test.ts +++ b/tests/web-search/web-search-passthrough-bridge.test.ts @@ -29,6 +29,8 @@ import { handleResponses } from "../../src/server/responses"; import { providerConfigSeed } from "../../src/providers/derive"; import { getProviderRegistryEntry } from "../../src/providers/registry"; import type { ResponsesTerminalRepairScheduler } from "../../src/server/responses-terminal-repair"; +import { bridgeSearchReplayScope, clearBridgeSearchReplayCacheForTests, peekBridgeSearchReplay } from "../../src/responses/bridge-search-replay-cache"; +import { reasoningReplayDestinationIdentity, reasoningReplayKeyCredentialIdentity } from "../../src/responses/reasoning-replay-cache"; import { resetProviderRequestPacingForTest, setProviderRequestPacingRuntimeForTest, @@ -1473,6 +1475,7 @@ describe("the reported turn, end to end through handleResponses", () => { }); test("a complete but terminal-less leg still repairs, on the first leg AND the continuation", async () => { + clearBridgeSearchReplayCacheForTests(); // Repair is registry-gated, so only a registry-keyed provider arms it: deepseek carries // modelResponsesTerminalRepair for the V4 flash ids. The fixture legs below emit a fully // complete item lifecycle and then stay open — the reported stall — with no terminal and @@ -1581,7 +1584,7 @@ describe("the reported turn, end to end through handleResponses", () => { try { const response = await handleResponses(new Request("http://localhost/v1/responses", { method: "POST", - headers: { "content-type": "application/json", authorization: "Bearer caller-inbound" }, + headers: { "content-type": "application/json", authorization: "Bearer caller-inbound", "thread-id": "thread-repaired-search" }, body: JSON.stringify({ model: "deepseek/deepseek-v4-flash", stream: true, @@ -1590,6 +1593,7 @@ describe("the reported turn, end to end through handleResponses", () => { }), }), cfg, { model: "", provider: "" }, { responsesTerminalRepairScheduler: scheduler, + admission: { kind: "loopback" }, }); const reader = response.body!.getReader(); try { @@ -1629,6 +1633,28 @@ describe("the reported turn, end to end through handleResponses", () => { expect(rest).toContain("response.completed"); expect(rest).toContain("The current release is 2.50.0."); expect(rest).toContain("[DONE]"); + const hosted = clientEvents(opened + rest).find(event => + event.type === "response.output_item.added" + && (event.item as Record | undefined)?.type === "web_search_call"); + const cellId = (hosted?.item as Record | undefined)?.id; + expect(typeof cellId).toBe("string"); + const scope = { + clientPrincipalId: "loopback", clientThreadId: "thread-repaired-search", + current: { + providerName: "deepseek", adapterName: "openai-responses", modelId: "deepseek-v4-flash", + providerDestinationIdentity: reasoningReplayDestinationIdentity(cfg.providers.deepseek!.baseUrl), + credentialIdentity: reasoningReplayKeyCredentialIdentity({ apiKey: "fixture-key" }), + }, + }; + // Results produced by repaired legs retain the same caller/serving boundary as + // ordinary search legs; knowing the emitted cell id does not widen that boundary. + expect(peekBridgeSearchReplay(bridgeSearchReplayScope(scope), cellId as string)?.output) + .toContain("opencodex 2.50.0"); + for (const changedScope of [ + { ...scope, clientPrincipalId: "another-caller" }, + { ...scope, clientThreadId: "another-thread" }, + { ...scope, current: { ...scope.current, credentialIdentity: "another-key" } }, + ]) expect(peekBridgeSearchReplay(bridgeSearchReplayScope(changedScope), cellId as string)).toBeUndefined(); } finally { try { await reader.cancel(); } catch { /* already closed */ } firstLeg.end(); @@ -1637,6 +1663,7 @@ describe("the reported turn, end to end through handleResponses", () => { } finally { releaseSpendHome(); globalThis.fetch = savedFetch; + clearBridgeSearchReplayCacheForTests(); } }); From ae52669293c4fd3c4b8b94c7f0580f7a345ecd9d Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Tue, 22 Sep 2026 23:19:19 +0900 Subject: [PATCH 09/12] test(server): reap fixture ACL workers before removing failover homes --- tests/server/server-key-failover-e2e.test.ts | 12 ++++++++++-- 1 file changed, 10 insertions(+), 2 deletions(-) diff --git a/tests/server/server-key-failover-e2e.test.ts b/tests/server/server-key-failover-e2e.test.ts index 5279288c90c..6b5366b48db 100644 --- a/tests/server/server-key-failover-e2e.test.ts +++ b/tests/server/server-key-failover-e2e.test.ts @@ -5,6 +5,9 @@ import { join } from "node:path"; import { apiKeyAccountLogLabel } from "../../src/codex/account-label"; import { readUsageEntries, resetUsageReadCacheForTests } from "../../src/usage/log"; import { loadConfig, saveConfig } from "../../src/config"; +import { flushConfigDirHardeningForTests } from "../../src/config/paths"; +import { flushNativeMainStartupReleases } from "../../src/codex/native-profile-startup"; +import { flushWindowsSecretAclReapsBeforeRemoval } from "../../src/lib/windows-secret-acl"; import { clearKeyCooldowns, getKeyCooldownUntil, rotateKeyOn429 } from "../../src/providers/key-failover"; import { deriveXaiConvId } from "../../src/providers/xai-transport"; import { @@ -43,9 +46,14 @@ beforeEach(() => { clearBridgeSearchReplayCacheForTests(); }); -afterEach(() => { - upstream?.stop(true); +afterEach(async () => { + await upstream?.stop(true); upstream = null; + await flushNativeMainStartupReleases(); + await flushConfigDirHardeningForTests(); + // Caller-facing ACL deadlines do not prove that their child released this home. + if (testDir) await flushWindowsSecretAclReapsBeforeRemoval(testDir); + if (isolatedCodexHome) await flushWindowsSecretAclReapsBeforeRemoval(isolatedCodexHome.path); if (previousHome === undefined) delete process.env.OPENCODEX_HOME; else process.env.OPENCODEX_HOME = previousHome; isolatedCodexHome?.restore(); From 421ba780ae89ded56b27ffb414a3ac055eeb7e35 Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Tue, 22 Sep 2026 23:20:43 +0900 Subject: [PATCH 10/12] fix(tests): reap sandbox ACL children before cleanup Keep test-home isolation, guard arming, and lock admission ahead of cleanup dependencies. Await config/native producers and the exact sandbox reap barrier in afterAll; the synchronous exit fallback defers undrained roots to ownership-checked recovery. Add a deterministic caller-belt/reap ordering regression without changing production ACL behavior or test timeouts. (cherry picked from commit fa96c7bca6780bb4d871ff5b5b0bb1187f3b9d16) --- scripts/test-layout/layout.json | 1 + structure/ops/docs-and-release.md | 7 +- .../ci-workflows/test-sandbox-cleanup.test.ts | 92 +++++++++++++++++++ tests/fixtures/test-layout-expected.json | 1 + tests/helpers/test-sandbox-cleanup.ts | 42 +++++++++ tests/preload.ts | 32 ++++--- 6 files changed, 161 insertions(+), 14 deletions(-) create mode 100644 tests/ci-workflows/test-sandbox-cleanup.test.ts create mode 100644 tests/helpers/test-sandbox-cleanup.ts diff --git a/scripts/test-layout/layout.json b/scripts/test-layout/layout.json index 7117b4679ea..17ad4c947ae 100644 --- a/scripts/test-layout/layout.json +++ b/scripts/test-layout/layout.json @@ -1469,6 +1469,7 @@ "terminal-guard-server.test.ts": "server", "terminal-guard.test.ts": "server", "test-home-guard.test.ts": "ci-workflows", + "test-sandbox-cleanup.test.ts": "ci-workflows", "test-runner.test.ts": "ci-workflows", "thought-signature-credential-scope.test.ts": "responses", "token-estimate.test.ts": "lib", diff --git a/structure/ops/docs-and-release.md b/structure/ops/docs-and-release.md index 5495a6728f0..83c33144962 100644 --- a/structure/ops/docs-and-release.md +++ b/structure/ops/docs-and-release.md @@ -392,7 +392,12 @@ storage-policy and api-usage families to its dedicated jobs. The Windows step di user-scoped test-run queue with `OCX_TEST_NO_QUEUE=1`: the batches already run sequentially in one dedicated job, and queueing a new batch behind a surviving process from the preceding batch spends the process timeout without executing tests. The per-process home isolation and live-home/service -manager guards remain active because the preload installs them before the lock boundary. A test +manager guards remain active because the preload installs them before the lock boundary. +`tests/preload.ts` awaits config hardening and native-main startup releases, then the sandbox's +registered ACL child reaps before removing that root. Its synchronous exit fallback leaves an +undrained root for ownership-checked stale recovery instead of blocking child cleanup with removal +retries. `tests/ci-workflows/test-sandbox-cleanup.test.ts` pins that ordering with a delayed reap. +A test failure, a process timeout and a Bun runtime crash each fail their job on the first occurrence; the batch runner still sweeps a crashed or timed-out batch one file per process, but only to attribute a failure the shard has already taken. The aggregate `ci` gate derives, from the event and the `changes` outputs, which diff --git a/tests/ci-workflows/test-sandbox-cleanup.test.ts b/tests/ci-workflows/test-sandbox-cleanup.test.ts new file mode 100644 index 00000000000..a446cc98054 --- /dev/null +++ b/tests/ci-workflows/test-sandbox-cleanup.test.ts @@ -0,0 +1,92 @@ +import { expect, test } from "bun:test"; +import { mkdtempSync, writeFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { + flushWindowsSecretAclReapsBeforeRemoval, + hardenSecretPathAsync, + resetHardenedStateForTests, + setAsyncIcaclsBeltSchedulerForTests, + setAsyncIcaclsRunnerForTests, + setNowForTests, + setPlatformForTests, + windowsSecretAclReapPendingAtOrBelow, + type IcaclsResult, +} from "../../src/lib/windows-secret-acl"; +import { setSyntheticWindowsPrincipalForTests } from "../../src/lib/windows-user-principal"; +import { removeTreeWithRetry } from "../helpers/remove-tree"; +import { createTestSandboxCleanup } from "../helpers/test-sandbox-cleanup"; + +test("sandbox removal waits for producers and the actual reap after the caller belt fires", async () => { + const root = mkdtempSync(join(tmpdir(), "ocx-sandbox-reap-")); + const target = join(root, "fixture.json"); + writeFileSync(target, "{}"); + const started = Promise.withResolvers(); + const reaped = Promise.withResolvers(); + const producers = Promise.withResolvers(); + let fireBelt: (() => void) | undefined; + let removals = 0; + let barrierEntered = false; + let clock = 0; + setNowForTests(() => clock); + setPlatformForTests("win32"); + setSyntheticWindowsPrincipalForTests("*S-1-5-21-1-2-3-1001"); + setAsyncIcaclsRunnerForTests(async () => { started.resolve(); return reaped.promise; }); + setAsyncIcaclsBeltSchedulerForTests(callback => { fireBelt = callback; return () => {}; }); + const timeout = { success: false, exitCode: null, timedOut: true, stdout: "" }; + const hardening = hardenSecretPathAsync(target, { required: false, deadlineMs: 10_000 }); + let pending: Promise | undefined; + try { + await started.promise; + clock = 10_001; // The real belt fires after the hardening deadline; no diagnostic probe remains. + fireBelt!(); + expect((await hardening).ok).toBe(false); + expect(windowsSecretAclReapPendingAtOrBelow(root)).toBe(true); + const cleanup = createTestSandboxCleanup({ + drainProducers: () => producers.promise, + waitForReaps: () => { barrierEntered = true; return flushWindowsSecretAclReapsBeforeRemoval(root); }, + hasPendingReaps: () => windowsSecretAclReapPendingAtOrBelow(root), + remove: () => { removals++; }, + }); + cleanup.onExit(); + expect(removals).toBe(0); + pending = cleanup.afterAll(); + expect(barrierEntered).toBe(false); + producers.resolve(); + await Promise.resolve(); + expect(barrierEntered).toBe(true); + cleanup.onExit(); + expect(removals).toBe(0); + // A live event loop, not sleepSync retries, must deliver this exit observation. + setTimeout(() => reaped.resolve(timeout), 0); + await pending; + expect(removals).toBe(1); + expect(windowsSecretAclReapPendingAtOrBelow(root)).toBe(false); + await cleanup.afterAll(); + cleanup.onExit(); + expect(removals).toBe(1); + } finally { + producers.resolve(); + reaped.resolve(timeout); + await hardening; + await pending; + await flushWindowsSecretAclReapsBeforeRemoval(root); + setAsyncIcaclsRunnerForTests(null); + setAsyncIcaclsBeltSchedulerForTests(null); + setPlatformForTests(null); + setNowForTests(null); + setSyntheticWindowsPrincipalForTests(null); + resetHardenedStateForTests(); + removeTreeWithRetry(root); + } +}); + +test("exit defers a root whose producers have not been drained even without a registered reap", () => { + let removals = 0; + const cleanup = createTestSandboxCleanup({ + drainProducers: async () => {}, waitForReaps: async () => {}, + hasPendingReaps: () => false, remove: () => { removals++; }, + }); + cleanup.onExit(); + expect(removals).toBe(0); +}); diff --git a/tests/fixtures/test-layout-expected.json b/tests/fixtures/test-layout-expected.json index 1fd592857c0..745743b6f84 100644 --- a/tests/fixtures/test-layout-expected.json +++ b/tests/fixtures/test-layout-expected.json @@ -1298,6 +1298,7 @@ "terminal-guard-server.test.ts": "server", "terminal-guard.test.ts": "server", "test-home-guard.test.ts": "ci-workflows", + "test-sandbox-cleanup.test.ts": "ci-workflows", "test-runner.test.ts": "ci-workflows", "thought-signature-credential-scope.test.ts": "responses", "token-estimate.test.ts": "lib", diff --git a/tests/helpers/test-sandbox-cleanup.ts b/tests/helpers/test-sandbox-cleanup.ts new file mode 100644 index 00000000000..e9496ff5658 --- /dev/null +++ b/tests/helpers/test-sandbox-cleanup.ts @@ -0,0 +1,42 @@ +/** Keep removal behind producer settlement and actual handle release, not a caller timeout. */ +export function createTestSandboxCleanup(options: { + drainProducers(): Promise; + waitForReaps(): Promise; + hasPendingReaps(): boolean; + remove(): void; +}) { + let complete = false; + let drained = false; + let running = false; + let cleanup: Promise | undefined; + const remove = () => { + if (complete) return; + try { + options.remove(); + complete = true; + } catch { + // A later run can reclaim the marked root; failure never grants stronger removal. + } + }; + return { + afterAll(): Promise { + return cleanup ??= (async () => { + running = true; + try { + await options.drainProducers(); + await options.waitForReaps(); + drained = true; + if (!options.hasPendingReaps()) remove(); + } finally { + running = false; + } + })(); + }, + onExit(): void { + // Exit cannot await. Before the async barrier finishes, even a not-yet-registered + // reaper may still be produced. Leave that root for ownership-checked stale recovery. + if (!drained || running || options.hasPendingReaps()) return; + remove(); + }, + }; +} diff --git a/tests/preload.ts b/tests/preload.ts index 50e65eff16d..bcbb1f265ee 100644 --- a/tests/preload.ts +++ b/tests/preload.ts @@ -117,16 +117,22 @@ if (process.platform === "win32" && lockPath && runLock.owner) { // Clean up only the root this preload created. The `bun run test` wrapper owns its own. // Bun test workers do not reliably run process `exit` handlers, so the test lifecycle hook -// is primary; the process hook remains a best-effort fallback for setup failures. -let cleanupComplete = false; -const cleanupIsolatedRoot = () => { - if (cleanupComplete) return; - try { - isolated.cleanup(); - cleanupComplete = true; - } catch { - // The wrapper contains this root, and a later bare run reclaims it after the grace period. - } -}; -afterAll(cleanupIsolatedRoot); -process.on("exit", cleanupIsolatedRoot); +// is primary; the process hook retries only an already-drained root. Setup failures +// leave an ownership-marked root for stale recovery rather than blocking on child handles. +// Load cleanup dependencies only AFTER home isolation, guard arming, and run-lock admission. +const { createTestSandboxCleanup } = await import("./helpers/test-sandbox-cleanup"); +const { flushWindowsSecretAclReapsBeforeRemoval, windowsSecretAclReapPendingAtOrBelow } = + await import("../src/lib/windows-secret-acl"); +const cleanup = createTestSandboxCleanup({ + drainProducers: async () => { + const { flushConfigDirHardeningForTests } = await import("../src/config/paths"); + const { flushNativeMainStartupReleases } = await import("../src/codex/native-profile-startup"); + await flushConfigDirHardeningForTests(); + await flushNativeMainStartupReleases(); + }, + waitForReaps: () => flushWindowsSecretAclReapsBeforeRemoval(isolated.root), + hasPendingReaps: () => windowsSecretAclReapPendingAtOrBelow(isolated.root), + remove: () => isolated.cleanup(), +}); +afterAll(cleanup.afterAll); +process.on("exit", cleanup.onExit); From 6b122cd2f024fb56f25667dad235575b5c50f001 Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Wed, 23 Sep 2026 00:36:24 +0900 Subject: [PATCH 11/12] test(server): own key-failover fixture shutdown across timeouts A timed-out pacing test left its proxy and request in a local finally while afterEach stopped the mock upstream and removed the state home. Own that work in the fixture, share one listener-stop promise, and settle cancellation before releasing the homes. Preserve ordinary assertion errors while consuming the expected late teardown abort. Use the existing test-only icacls runners for the synthetic HTTP fixture so cold Windows permission setup does not consume its unchanged five-second behavior budget. Preserve real listener and SQLite ownership checks, including reacquisition from a different home after cancellation. Resolve cleanup owners during protected preload setup, and drain native releases before config work. Validation: the original explicit eight-file wrapper passes 199 tests with 976 assertions; typecheck, structure, privacy and the file-size ratchet pass. This does not claim real OS ACL coverage or repair the separate changed-mode empty-selection result. --- structure/ops/docs-and-release.md | 10 +- .../ci-workflows/test-sandbox-cleanup.test.ts | 40 +++++- tests/helpers/test-sandbox-cleanup.ts | 43 ++++++ tests/preload.ts | 8 +- tests/server/server-key-failover-e2e.test.ts | 133 ++++++++++++++---- 5 files changed, 202 insertions(+), 32 deletions(-) diff --git a/structure/ops/docs-and-release.md b/structure/ops/docs-and-release.md index 83c33144962..e8be5f30617 100644 --- a/structure/ops/docs-and-release.md +++ b/structure/ops/docs-and-release.md @@ -393,10 +393,18 @@ user-scoped test-run queue with `OCX_TEST_NO_QUEUE=1`: the batches already run s dedicated job, and queueing a new batch behind a surviving process from the preceding batch spends the process timeout without executing tests. The per-process home isolation and live-home/service manager guards remain active because the preload installs them before the lock boundary. -`tests/preload.ts` awaits config hardening and native-main startup releases, then the sandbox's +`tests/preload.ts` resolves cleanup dependencies after home/lock admission and before test cases, +then teardown awaits native-main startup releases and config hardening, followed by the sandbox's registered ACL child reaps before removing that root. Its synchronous exit fallback leaves an undrained root for ownership-checked stale recovery instead of blocking child cleanup with removal retries. `tests/ci-workflows/test-sandbox-cleanup.test.ts` pins that ordering with a delayed reap. +Case fixtures own their proxy listeners and cancellable asynchronous work independently of the +test runner's deadline. The key-failover fixture settles both before stopping its upstream mock, +draining producers and ACL reaps, or restoring/removing either home; its body and teardown share +one stop promise so a timed-out request cannot retain the previous home's spend-ledger lease. +The HTTP/key fixture uses the existing test-only icacls runners for deterministic synthetic-home +preparation; it does not claim real OS ACL coverage. Those runners return to their defaults after +producer/reap settlement, while dedicated ACL tests and the real SQLite lease remain authoritative. A test failure, a process timeout and a Bun runtime crash each fail their job on the first occurrence; the batch runner still sweeps a crashed or timed-out batch one file per process, but only to attribute a diff --git a/tests/ci-workflows/test-sandbox-cleanup.test.ts b/tests/ci-workflows/test-sandbox-cleanup.test.ts index a446cc98054..e1acf8857b6 100644 --- a/tests/ci-workflows/test-sandbox-cleanup.test.ts +++ b/tests/ci-workflows/test-sandbox-cleanup.test.ts @@ -15,7 +15,7 @@ import { } from "../../src/lib/windows-secret-acl"; import { setSyntheticWindowsPrincipalForTests } from "../../src/lib/windows-user-principal"; import { removeTreeWithRetry } from "../helpers/remove-tree"; -import { createTestSandboxCleanup } from "../helpers/test-sandbox-cleanup"; +import { createTestCaseLifecycle, createTestSandboxCleanup } from "../helpers/test-sandbox-cleanup"; test("sandbox removal waits for producers and the actual reap after the caller belt fires", async () => { const root = mkdtempSync(join(tmpdir(), "ocx-sandbox-reap-")); @@ -90,3 +90,41 @@ test("exit defers a root whose producers have not been drained even without a re cleanup.onExit(); expect(removals).toBe(0); }); + +test("case teardown cancels late work and awaits its shared listener stop before the home is released", async () => { + const lifecycle = createTestCaseLifecycle(); + const entered = Promise.withResolvers(); + const released = Promise.withResolvers(); + let stopCalls = 0; + let settled = false; + const stop = lifecycle.ownStop(async () => { stopCalls++; await released.promise; }); + const work = lifecycle.run(async () => { + try { + await new Promise((_resolve, reject) => { + lifecycle.abort.signal.addEventListener("abort", () => reject(lifecycle.abort.signal.reason), { once: true }); + entered.resolve(); + }); + } finally { + await stop(); + settled = true; + } + }); + await entered.promise; + const cleanup = lifecycle.close(); + expect(lifecycle.close()).toBe(cleanup); + await Promise.resolve(); + expect(stopCalls).toBe(1); + expect(settled).toBe(false); + released.resolve(); + await cleanup; + expect(settled).toBe(true); + expect(stopCalls).toBe(1); + await expect(work).resolves.toBeUndefined(); +}); + +test("case ownership preserves ordinary assertion failures", async () => { + const lifecycle = createTestCaseLifecycle(); + const failure = new Error("fixture assertion failure"); + await expect(lifecycle.run(async () => { throw failure; })).rejects.toBe(failure); + await lifecycle.close(); +}); diff --git a/tests/helpers/test-sandbox-cleanup.ts b/tests/helpers/test-sandbox-cleanup.ts index e9496ff5658..c8c7929ace9 100644 --- a/tests/helpers/test-sandbox-cleanup.ts +++ b/tests/helpers/test-sandbox-cleanup.ts @@ -40,3 +40,46 @@ export function createTestSandboxCleanup(options: { }, }; } + +/** Own asynchronous case work independently of Bun's timeout on the returned promise. */ +export function createTestCaseLifecycle() { + const abort = new AbortController(); + const pending = new Set>(); + const stops: Array<() => Promise> = []; + let closing: Promise | undefined; + return { + abort, + ownStop(stop: () => void | Promise): () => Promise { + let stopped: Promise | undefined; + const once = () => stopped ??= Promise.resolve().then(stop); + stops.push(once); + return once; + }, + run(work: () => Promise): Promise { + abort.signal.throwIfAborted(); + const operation = Promise.resolve().then(work).catch(error => { + // Teardown cancellation is already owned by close(); a timed-out test must not + // throw its expected abort later as an unrelated error in the following case. + if (closing && abort.signal.aborted && error instanceof Error && error.name === "AbortError") return; + throw error; + }); + pending.add(operation); + // Keep a rejection observer after the runner's timeout. Return the original promise + // so an ordinary assertion/rejection still fails the test that owns it. + void operation.then(() => pending.delete(operation), () => pending.delete(operation)); + return operation; + }, + close(): Promise { + return closing ??= (async () => { + abort.abort(); + // Stop listeners while requests settle; a request may be waiting for that stop. + // Its finally and afterEach share the same stop promise instead of racing teardown. + const stopped = Promise.allSettled(stops.map(stop => stop())); + await Promise.allSettled([...pending]); + const results = await stopped; + const failures = results.flatMap(result => result.status === "rejected" ? [result.reason] : []); + if (failures.length) throw new AggregateError(failures, "test case listener cleanup failed"); + })(); + }, + }; +} diff --git a/tests/preload.ts b/tests/preload.ts index bcbb1f265ee..18d29c7d90f 100644 --- a/tests/preload.ts +++ b/tests/preload.ts @@ -123,12 +123,14 @@ if (process.platform === "win32" && lockPath && runLock.owner) { const { createTestSandboxCleanup } = await import("./helpers/test-sandbox-cleanup"); const { flushWindowsSecretAclReapsBeforeRemoval, windowsSecretAclReapPendingAtOrBelow } = await import("../src/lib/windows-secret-acl"); +// Resolve cleanup owners during protected setup, not for the first time inside a timed +// afterAll hook. Cleanup must drain existing producers rather than initialize their graph. +const { flushConfigDirHardeningForTests } = await import("../src/config/paths"); +const { flushNativeMainStartupReleases } = await import("../src/codex/native-profile-startup"); const cleanup = createTestSandboxCleanup({ drainProducers: async () => { - const { flushConfigDirHardeningForTests } = await import("../src/config/paths"); - const { flushNativeMainStartupReleases } = await import("../src/codex/native-profile-startup"); - await flushConfigDirHardeningForTests(); await flushNativeMainStartupReleases(); + await flushConfigDirHardeningForTests(); }, waitForReaps: () => flushWindowsSecretAclReapsBeforeRemoval(isolated.root), hasPendingReaps: () => windowsSecretAclReapPendingAtOrBelow(isolated.root), diff --git a/tests/server/server-key-failover-e2e.test.ts b/tests/server/server-key-failover-e2e.test.ts index 6b5366b48db..bfed950937d 100644 --- a/tests/server/server-key-failover-e2e.test.ts +++ b/tests/server/server-key-failover-e2e.test.ts @@ -7,7 +7,8 @@ import { readUsageEntries, resetUsageReadCacheForTests } from "../../src/usage/l import { loadConfig, saveConfig } from "../../src/config"; import { flushConfigDirHardeningForTests } from "../../src/config/paths"; import { flushNativeMainStartupReleases } from "../../src/codex/native-profile-startup"; -import { flushWindowsSecretAclReapsBeforeRemoval } from "../../src/lib/windows-secret-acl"; +import { flushWindowsSecretAclReapsBeforeRemoval, setAsyncIcaclsRunnerForTests, setIcaclsRunnerForTests } from "../../src/lib/windows-secret-acl"; +import { acquireSpendLedgerOwner } from "../../src/lib/spend-ledger-owner"; import { clearKeyCooldowns, getKeyCooldownUntil, rotateKeyOn429 } from "../../src/providers/key-failover"; import { deriveXaiConvId } from "../../src/providers/xai-transport"; import { @@ -20,11 +21,12 @@ import { clearBridgeSearchReplayCacheForTests, rememberBridgeSearchReplay, } from "../../src/responses/bridge-search-replay-cache"; -import { startServer } from "../../src/server"; +import { startServer as startProductionServer } from "../../src/server"; import { handleResponses } from "../../src/server/responses"; import type { OcxConfig } from "../../src/types"; import { installIsolatedCodexHome, type IsolatedCodexHome } from "../helpers/isolated-codex-home"; import { removeTreeWithRetry } from "../helpers/remove-tree"; +import { createTestCaseLifecycle } from "../helpers/test-sandbox-cleanup"; import { managementFetch } from "../helpers/management-auth"; import { resetProviderRequestPacingForTest, setProviderRequestPacingRuntimeForTest, waitForProviderRequestSlot } from "../../src/providers/request-pacing"; import { providerApiKeySelectionIsCurrent, resolveCurrentProviderApiKeyTransport } from "../../src/providers/api-key-selection"; @@ -35,8 +37,31 @@ let testDir = ""; let previousHome: string | undefined; let isolatedCodexHome: IsolatedCodexHome | null = null; let upstream: ReturnType | null = null; +let caseLifecycle = createTestCaseLifecycle(); + +function runCase(work: (lifecycle: ReturnType) => Promise) { + const lifecycle = caseLifecycle; + return lifecycle.run(() => work(lifecycle)); +} + +function startServer(...args: Parameters) { + caseLifecycle.abort.signal.throwIfAborted(); + const server = startProductionServer(...args); + const stop = server.stop.bind(server); + Object.defineProperty(server, "stop", { + configurable: true, + value: caseLifecycle.ownStop(() => stop(true)), + }); + return server; +} beforeEach(() => { + // These fixtures exercise real listeners, key selection and SQLite ownership, not OS ACLs. + // Keep permission ceremony out of the 5s HTTP behavior budget, using the existing test seam. + const aclOk = { success: true, exitCode: 0, timedOut: false, stdout: "" }; + setIcaclsRunnerForTests(() => aclOk); + setAsyncIcaclsRunnerForTests(async () => aclOk); + caseLifecycle = createTestCaseLifecycle(); previousHome = process.env.OPENCODEX_HOME; isolatedCodexHome = installIsolatedCodexHome("ocx-keyfail-e2e-codex-"); testDir = mkdtempSync(join(tmpdir(), "ocx-keyfail-e2e-")); @@ -47,21 +72,30 @@ beforeEach(() => { }); afterEach(async () => { - await upstream?.stop(true); - upstream = null; - await flushNativeMainStartupReleases(); - await flushConfigDirHardeningForTests(); - // Caller-facing ACL deadlines do not prove that their child released this home. - if (testDir) await flushWindowsSecretAclReapsBeforeRemoval(testDir); - if (isolatedCodexHome) await flushWindowsSecretAclReapsBeforeRemoval(isolatedCodexHome.path); - if (previousHome === undefined) delete process.env.OPENCODEX_HOME; - else process.env.OPENCODEX_HOME = previousHome; - isolatedCodexHome?.restore(); - isolatedCodexHome = null; - if (testDir) removeTreeWithRetry(testDir); - clearKeyCooldowns(); - clearReasoningReplayCacheForTests(); - clearBridgeSearchReplayCacheForTests(); + try { + // Bun's test timeout does not cancel the async body. End its requests and release + // the actual proxy's spend lease before stopping the upstream or deleting either home. + await caseLifecycle.close(); + resetProviderRequestPacingForTest(); + await upstream?.stop(true); + upstream = null; + await flushNativeMainStartupReleases(); + await flushConfigDirHardeningForTests(); + // Caller-facing ACL deadlines do not prove that their child released this home. + if (testDir) await flushWindowsSecretAclReapsBeforeRemoval(testDir); + if (isolatedCodexHome) await flushWindowsSecretAclReapsBeforeRemoval(isolatedCodexHome.path); + if (previousHome === undefined) delete process.env.OPENCODEX_HOME; + else process.env.OPENCODEX_HOME = previousHome; + isolatedCodexHome?.restore(); + isolatedCodexHome = null; + if (testDir) removeTreeWithRetry(testDir); + clearKeyCooldowns(); + clearReasoningReplayCacheForTests(); + clearBridgeSearchReplayCacheForTests(); + } finally { + setIcaclsRunnerForTests(null); + setAsyncIcaclsRunnerForTests(null); + } }); describe("server 429 key failover (end-to-end)", () => { @@ -103,7 +137,48 @@ describe("server 429 key failover (end-to-end)", () => { expect(providerApiKeySelectionIsCurrent(config, "current", current)).toBe(true); }); - test.each(["responses", "chat/completions"])("%s logs only the key selected after pacing", async surface => { + test("cancelling a queued proxy request releases the actual spend lease before the next home", async () => { + const lifecycle = caseLifecycle; + const queued = Promise.withResolvers(); + setProviderRequestPacingRuntimeForTest({ + now: () => 0, setTimer(callback) { queued.resolve(); return callback; }, + clearTimer() {}, enqueueMicrotask: queueMicrotask, + }); + let sends = 0; + upstream = Bun.serve({ hostname: "127.0.0.1", port: 0, fetch() { sends++; return Response.json({}); } }); + const config = { port: 0, hostname: "127.0.0.1", defaultProvider: "cancelled", providers: { cancelled: { + adapter: "openai-chat", baseUrl: `http://127.0.0.1:${upstream.port}/v1`, allowPrivateNetwork: true, + authMode: "key", apiKey: "synthetic-cancel-key", requestPacing: { enabled: true, minIntervalMs: 100 }, + } } } as OcxConfig; + saveConfig(config); + const server = startServer(0); + const nextHome = mkdtempSync(join(tmpdir(), "ocx-keyfail-next-")); + try { + await waitForProviderRequestSlot("cancelled", config.providers.cancelled); + const pending = lifecycle.run(async () => { + try { + const response = await fetch(new URL("/v1/responses", server.url), { + method: "POST", headers: { "content-type": "application/json" }, signal: lifecycle.abort.signal, + body: JSON.stringify({ model: "cancelled/test", stream: false, input: "hello" }), + }); + await response.text(); + } finally { await server.stop(true); } + }); + await Promise.race([queued.promise, pending.then(() => { throw new Error("request never reached pacing"); })]); + expect(() => acquireSpendLedgerOwner(nextHome)).toThrow("different state directory"); + await lifecycle.close(); + await pending; + const nextLease = acquireSpendLedgerOwner(nextHome); + nextLease.release(); + expect(sends).toBe(0); + } finally { + await lifecycle.close(); + await flushWindowsSecretAclReapsBeforeRemoval(nextHome); + removeTreeWithRetry(nextHome); + } + }); + + test.each(["responses", "chat/completions"])("%s logs only the key selected after pacing", surface => runCase(async lifecycle => { resetUsageReadCacheForTests(); let now = 0; let resumePacing: (() => void) | undefined; @@ -133,7 +208,7 @@ describe("server 429 key failover (end-to-end)", () => { } } } as OcxConfig; saveConfig(config); const server = startServer(0); - const abort = new AbortController(); + const abort = lifecycle.abort; try { await waitForProviderRequestSlot("paced", config.providers.paced); const pending = fetch(new URL(`/v1/${surface}`, server.url), { @@ -141,10 +216,12 @@ describe("server 429 key failover (end-to-end)", () => { body: JSON.stringify({ model: "paced/test", stream: false, ...(surface === "responses" ? { input: "hello" } : { messages: [{ role: "user", content: "hello" }] }) }), }); - await queued.promise; + await Promise.race([queued.promise, pending.then(response => { + throw new Error(`request settled before its pacing barrier (${response.status})`); + })]); expect(seen).toHaveLength(0); const selected = await managementFetch(new URL("/api/providers/keys/active", server.url), { - method: "PUT", headers: { "content-type": "application/json" }, + method: "PUT", headers: { "content-type": "application/json" }, signal: abort.signal, body: JSON.stringify({ name: "paced", id: "second" }), }); expect(selected.status).toBe(200); @@ -165,7 +242,7 @@ describe("server 429 key failover (end-to-end)", () => { await server.stop(true); resetProviderRequestPacingForTest(); } - }); + })); test.each(["responses", "chat/completions"])("%s carries the configured env-key identity through 429 recovery", async inbound => { const seen: string[] = []; @@ -1202,7 +1279,7 @@ test.each([false, true])("key refetch retains transient recovery metadata (strea } finally { await server.stop(true); } }); -test("a dispatch-time key switch rebuilds the bridged-search restore under the new credential", async () => { +test("a dispatch-time key switch rebuilds the bridged-search restore under the new credential", () => runCase(async lifecycle => { // Regression for the oauthDispatch rebuild order: the Responses adapter restores a replayed // web_search_call from the memo keyed by _reasoningReplayScope, so the rebuild must rebind // that scope to the refreshed credential BEFORE buildRequest runs. Restoring under the key @@ -1266,7 +1343,7 @@ test("a dispatch-time key switch rebuilds the bridged-search restore under the n ); const server = startServer(0); - const abort = new AbortController(); + const abort = lifecycle.abort; try { await waitForProviderRequestSlot("pooled", config.providers.pooled); const pending = fetch(new URL("/v1/responses", server.url), { @@ -1282,9 +1359,11 @@ test("a dispatch-time key switch rebuilds the bridged-search restore under the n ], }), }); - await queued.promise; + await Promise.race([queued.promise, pending.then(response => { + throw new Error(`request settled before its pacing barrier (${response.status})`); + })]); const selected = await managementFetch(new URL("/api/providers/keys/active", server.url), { - method: "PUT", headers: { "content-type": "application/json" }, + method: "PUT", headers: { "content-type": "application/json" }, signal: abort.signal, body: JSON.stringify({ name: "pooled", id: "second" }), }); expect(selected.status).toBe(200); @@ -1304,4 +1383,4 @@ test("a dispatch-time key switch rebuilds the bridged-search restore under the n await server.stop(true); resetProviderRequestPacingForTest(); } -}); +})); From 6ea3a95c212c94103177fd17ffa37b28f05ce843 Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Wed, 23 Sep 2026 00:55:27 +0900 Subject: [PATCH 12/12] docs(tests): scope shared cleanup contract independently of case fixtures Keep key-failover HTTP fixture execution evidence in the PR and validation manifest. This source document describes common cancellation, listener, producer and reap ownership without importing CI-only workflow contracts. --- structure/ops/docs-and-release.md | 12 +++++------- 1 file changed, 5 insertions(+), 7 deletions(-) diff --git a/structure/ops/docs-and-release.md b/structure/ops/docs-and-release.md index e8be5f30617..beb88e92aa8 100644 --- a/structure/ops/docs-and-release.md +++ b/structure/ops/docs-and-release.md @@ -398,13 +398,11 @@ then teardown awaits native-main startup releases and config hardening, followed registered ACL child reaps before removing that root. Its synchronous exit fallback leaves an undrained root for ownership-checked stale recovery instead of blocking child cleanup with removal retries. `tests/ci-workflows/test-sandbox-cleanup.test.ts` pins that ordering with a delayed reap. -Case fixtures own their proxy listeners and cancellable asynchronous work independently of the -test runner's deadline. The key-failover fixture settles both before stopping its upstream mock, -draining producers and ACL reaps, or restoring/removing either home; its body and teardown share -one stop promise so a timed-out request cannot retain the previous home's spend-ledger lease. -The HTTP/key fixture uses the existing test-only icacls runners for deterministic synthetic-home -preparation; it does not claim real OS ACL coverage. Those runners return to their defaults after -producer/reap settlement, while dedicated ACL tests and the real SQLite lease remain authoritative. +`tests/helpers/test-sandbox-cleanup.ts` also exposes case-scoped lifecycle ownership: cancellation +starts listener stops while owned asynchronous work settles, and repeated close/stop calls share +one promise. Expected teardown aborts are observed without hiding ordinary assertion failures. +Callers settle that lifecycle before draining producers/reaps and restoring or removing a home; +the helper does not replace fixture-specific cleanup or claim OS ACL coverage for synthetic tests. A test failure, a process timeout and a Bun runtime crash each fail their job on the first occurrence; the batch runner still sweeps a crashed or timed-out batch one file per process, but only to attribute a