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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions scripts/test-layout/layout.json
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
8 changes: 4 additions & 4 deletions src/adapters/openai-responses/passthrough.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
30 changes: 20 additions & 10 deletions src/responses/bridge-search-replay-cache.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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;
Expand Down Expand Up @@ -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 {
Expand Down
2 changes: 2 additions & 0 deletions src/server/responses/core-options.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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?: {
Expand Down
92 changes: 54 additions & 38 deletions src/server/responses/passthrough-delivery.ts
Original file line number Diff line number Diff line change
Expand Up @@ -380,47 +380,72 @@ 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
// 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,
hostedTool: parsed._webSearch,
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) => {
Expand All @@ -433,16 +458,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
Expand Down
36 changes: 29 additions & 7 deletions src/server/responses/request-prepare.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ import {
sessionIdHeaderFromRequest,
reasoningReplayConversationIdFromResponsesRequest,
} from "../request-log-conversation";
import { resolveContextPrincipal } from "../auth-cors";
import {
isShadowSourceModel,
shadowSourceModelPrefix,
Expand Down Expand Up @@ -225,20 +226,29 @@ 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;
if (shadowIntercept?.enabled && shadowIntercept.model && typeof rawShadowModel === "string"
&& isShadowSourceModel(rawShadowModel, shadowIntercept.sourceModels)) {
const shadowComboId = resolveComboId(config, shadowIntercept.model);
if (shadowComboId && Object.hasOwn(config.combos ?? {}, shadowComboId)) {
(body as Record<string, unknown>).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<string, unknown>).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);
}
}
}
}
Expand All @@ -247,6 +257,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,
Expand Down Expand Up @@ -369,6 +382,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", {
Expand Down Expand Up @@ -407,6 +421,14 @@ export async function prepareResponsesRequest(
parsed._reasoningReplayScope = { clientThreadId: reasoningReplayConversationId };
}
}
if (parsed._reasoningReplayScope) {
// 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 };
}
// 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) {
Expand Down
7 changes: 5 additions & 2 deletions src/server/responses/request-transport.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 } : {}),
Expand All @@ -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");
Expand Down
2 changes: 2 additions & 0 deletions src/types/request.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading