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
2 changes: 2 additions & 0 deletions src/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1262,6 +1262,8 @@ const configSchema = z.object({
providerContextCapValues: z.record(z.string(), z.number().int().positive()).optional(),
contextCapValue: z.number().int().positive().optional(),
multiAgentGuidanceEnabled: z.boolean().optional(),
// Invalid hand edits disable only this experimental opt-in.
v2RoutedDelegationBridge: z.boolean().optional().catch(undefined),
Comment thread
coderabbitai[bot] marked this conversation as resolved.
// Invalid optional recovery config must not discard unrelated provider/account state.
agentTaskRecovery: agentTaskRecoverySchema.optional().catch(undefined),
// Same rationale: a bad notify section must not cost the operator their providers.
Expand Down
6 changes: 6 additions & 0 deletions src/responses/state.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2314,6 +2314,12 @@ export function markBodyNonPersistable(body: unknown): void {
if (body && typeof body === "object") nonPersistableBodies.add(body as object);
}

export function copyBodyNonPersistableMarker(source: unknown, target: unknown): void {
if (!source || typeof source !== "object" || Array.isArray(source)) return;
if (!target || typeof target !== "object" || Array.isArray(target)) return;
if (nonPersistableBodies.has(source as object)) nonPersistableBodies.add(target as object);
}

export function rememberResponseState(
requestBody: unknown,
response: { id?: unknown; output?: unknown; status?: unknown; incomplete_details?: unknown },
Expand Down
24 changes: 22 additions & 2 deletions src/server/management/agent-settings-routes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -253,6 +253,7 @@ export async function handleAgentSettingsRoutes(ctx: ManagementContext): Promise
// max_depth is V1-only upstream; this is the global-flag statement, derived
// server-side so no client can present it as an effective V2 limit.
agentsMaxDepthAppliesWhenV2Disabled: !enabled,
v2RoutedDelegationBridge: config.v2RoutedDelegationBridge === true,
});
}
if (url.pathname === "/api/v2" && req.method === "PUT") {
Expand All @@ -265,6 +266,7 @@ export async function handleAgentSettingsRoutes(ctx: ManagementContext): Promise
agentsMaxDepth?: unknown;
subagentDeveloperInstructions?: unknown;
multiAgentModeHintText?: unknown;
v2RoutedDelegationBridge?: unknown;
};
try { body = await readManagementJsonBody(req); } catch (error) { rethrowManagementBodyTooLarge(error); return jsonResponse({ error: "invalid JSON body" }, 400); }
const wantsFlag = body.enabled !== undefined;
Expand All @@ -275,8 +277,9 @@ export async function handleAgentSettingsRoutes(ctx: ManagementContext): Promise
const wantsMaxDepth = body.agentsMaxDepth !== undefined;
const wantsSubagentInstructions = body.subagentDeveloperInstructions !== undefined;
const wantsModeHintText = body.multiAgentModeHintText !== undefined;
if (!wantsFlag && !wantsThreads && !wantsMode && !wantsKeepNative && !wantsAgentsEnabled && !wantsMaxDepth && !wantsSubagentInstructions && !wantsModeHintText) {
return jsonResponse({ error: "body must set enabled, multiAgentMode, keepNativeChatGptOnV1, maxConcurrentThreadsPerSession, agentsEnabled, agentsMaxDepth, subagentDeveloperInstructions, and/or multiAgentModeHintText" }, 400);
const wantsV2RoutedDelegationBridge = body.v2RoutedDelegationBridge !== undefined;
if (!wantsFlag && !wantsThreads && !wantsMode && !wantsKeepNative && !wantsAgentsEnabled && !wantsMaxDepth && !wantsSubagentInstructions && !wantsModeHintText && !wantsV2RoutedDelegationBridge) {
return jsonResponse({ error: "body must set enabled, multiAgentMode, keepNativeChatGptOnV1, maxConcurrentThreadsPerSession, agentsEnabled, agentsMaxDepth, subagentDeveloperInstructions, multiAgentModeHintText, and/or v2RoutedDelegationBridge" }, 400);
}
if (wantsFlag && typeof body.enabled !== "boolean") return jsonResponse({ error: "body.enabled must be a boolean" }, 400);
if (wantsMode && body.multiAgentMode !== "v1" && body.multiAgentMode !== "default" && body.multiAgentMode !== "v2") {
Expand All @@ -285,6 +288,9 @@ export async function handleAgentSettingsRoutes(ctx: ManagementContext): Promise
if (wantsKeepNative && typeof body.keepNativeChatGptOnV1 !== "boolean") {
return jsonResponse({ error: "body.keepNativeChatGptOnV1 must be a boolean" }, 400);
}
if (wantsV2RoutedDelegationBridge && typeof body.v2RoutedDelegationBridge !== "boolean") {
return jsonResponse({ error: "body.v2RoutedDelegationBridge must be a boolean" }, 400);
}
if (wantsThreads && (typeof body.maxConcurrentThreadsPerSession !== "number" || !Number.isInteger(body.maxConcurrentThreadsPerSession) || body.maxConcurrentThreadsPerSession < 1)) {
return jsonResponse({ error: "body.maxConcurrentThreadsPerSession must be an integer >= 1" }, 400);
}
Expand Down Expand Up @@ -386,6 +392,19 @@ export async function handleAgentSettingsRoutes(ctx: ManagementContext): Promise
// asserts this route file contains no direct write primitive, and matches on the
// symbol name even inside a comment.
const scalarWrites: Array<{ field: string; run: () => { ok: true; changed: boolean } | { ok: false; error: string } }> = [];
if (wantsV2RoutedDelegationBridge) scalarWrites.push({
field: "v2RoutedDelegationBridge",
run: () => {
try {
if (body.v2RoutedDelegationBridge === false) deleteConfigTopLevelKey(config, "v2RoutedDelegationBridge");
else config.v2RoutedDelegationBridge = true;
saveConfigPreservingClaudeCode(config);
return { ok: true, changed: true };
} catch (error) {
return { ok: false, error: error instanceof Error ? error.message : String(error) };
}
},
});
if (wantsAgentsEnabled) scalarWrites.push({ field: "agentsEnabled", run: () => setAgentsEnabled(body.agentsEnabled as boolean | null) });
if (wantsMaxDepth) scalarWrites.push({ field: "agentsMaxDepth", run: () => setAgentsMaxDepth(body.agentsMaxDepth as number | null) });
if (wantsSubagentInstructions) scalarWrites.push({ field: "subagentDeveloperInstructions", run: () => setSubagentDeveloperInstructions(body.subagentDeveloperInstructions as string | null) });
Expand Down Expand Up @@ -425,6 +444,7 @@ export async function handleAgentSettingsRoutes(ctx: ManagementContext): Promise
multiAgentModeHintText: getMultiAgentModeHintText(),
multiAgentModeHintRecommendation: MULTI_AGENT_MODE_HINT_RECOMMENDATION,
agentsMaxDepthAppliesWhenV2Disabled: !enabled,
v2RoutedDelegationBridge: config.v2RoutedDelegationBridge === true,
warnings,
catalogRefresh,
});
Expand Down
82 changes: 79 additions & 3 deletions src/server/responses/core.ts
Original file line number Diff line number Diff line change
Expand Up @@ -41,11 +41,13 @@ import { FORWARD_HEADERS, sanitizeReasoningInputContent } from "../../adapters/o
import { XaiToolSchemaCompatibilityError } from "../../adapters/xai-tool-schema";
import {
copyPreviousResponseReplayProvenance,
copyBodyNonPersistableMarker,
expandPreviousResponseInput,
markBodyNonPersistable,
previousResponseProviderState,
previousResponseReplayFailure,
previousResponseScopeMismatch,
previousResponseReplayPrefixLength,
rememberResponseState,
} from "../../responses/state";
import {
Expand Down Expand Up @@ -384,6 +386,14 @@ import {
import type { EffectiveSubagentRoster, SpawnAgentSurface } from "../../codex/catalog";

import { buildToolBridgeMaps, collabSurface, injectDeveloperMessage, multiAgentGuidanceText } from "./collaboration";
import { isMultiAgentV2Enabled } from "../../codex/features";
import {
createV2RoutedDelegationSseRewrite,
injectV2RoutedDelegationBridge,
rewriteV2RoutedDelegationCallsInJson,
type V2RoutedDelegationBridgeContext,
} from "./v2-routed-delegation-bridge";
import { decideV2RoutedDelegationBridge } from "./v2-routed-delegation-policy";
import { mapCodexAuthContextErrorToResponse, nativeMainRefreshFailureResponse } from "./codex-auth-error";
import { hasUnreadableEncryptedAgentTask, looksLikeBackendCiphertext, sanitizeEncryptedContentInPlace } from "./encrypted-payload";
import { fetchWithHeaderTimeout, providerFetch, safeHostLabel, safeOriginLabel, storedPoolReplayDispatchNotifier, type ProviderFetchOptions } from "./fetch-helpers";
Expand Down Expand Up @@ -509,6 +519,18 @@ function readProviderContinuationOwner(
return { kind: "valid", owner: { ...owner } };
}

function requestHasSubagentMarker(headers: Headers): boolean {
if (headers.has("x-openai-subagent")) return true;
const metadata = headers.get("x-codex-turn-metadata");
if (!metadata) return false;
try {
const parsedMetadata = JSON.parse(metadata) as { subagent_kind?: unknown };
return typeof parsedMetadata.subagent_kind === "string" && parsedMetadata.subagent_kind.length > 0;
} catch {
return false;
}
}

function providerContinuationPayload(
state: OcxProviderContinuationState | undefined,
): OcxProviderContinuationState | undefined {
Expand Down Expand Up @@ -3502,6 +3524,7 @@ async function handleResponsesInner(
return formatErrorResponse(400, "invalid_request_error", err instanceof Error ? err.message : String(err));
}
options.onRequestBodyRead?.();
let v2RoutedDelegationBridge: V2RoutedDelegationBridgeContext | undefined;
const responseStateOptions = (force = false): { force?: boolean; clientThreadId?: string } => ({
...(force ? { force: true } : {}),
...(parsed._clientThreadId ? { clientThreadId: parsed._clientThreadId } : {}),
Expand Down Expand Up @@ -3546,6 +3569,7 @@ async function handleResponsesInner(

let route: RouteResult;
let credentialDomainWasRewritten = false;
let shadowIntercepted = false;
try {
// A `compaction_trigger` turn may name a bare native model the operator has
// no canonical OpenAI route for (#2901). Only the initial compaction route
Expand All @@ -3567,6 +3591,7 @@ async function handleResponsesInner(
} catch { /* Native Codex helper calls remain OpenAI-owned without an enabled OpenAI route. */ }
const targetRoute = resolveRoute(_sci.model);
if (shouldInterceptShadowCall(parsed.modelId, _sci.sourceModels, sourceIdentity, targetRoute)) {
shadowIntercepted = true;
credentialDomainWasRewritten = true;
const _sciOriginal = parsed.modelId;
parsed.modelId = _sci.model;
Expand Down Expand Up @@ -3899,6 +3924,44 @@ async function handleResponsesInner(
);
}

const bridgeEnabled = config.v2RoutedDelegationBridge === true;
const bridgeDecision = bridgeEnabled ? decideV2RoutedDelegationBridge({
enabled: true,
inboundWire,
multiAgentMode: config.multiAgentMode,
upstreamV2Enabled: isMultiAgentV2Enabled(),
canonicalNativeRoute: isCanonicalOpenAiForwardProvider(route.provider),
hasSubagentMarker: requestHasSubagentMarker(req.headers),
threadSpawn,
comboAttempt: options.comboAttempt === true,
compaction: parsed._compactionRequest === true,
shadowRoute: shadowIntercepted,
collaborationSurface: collabSurface(parsed),
body: parsed._rawBody,
replayPrefixLength: previousResponseReplayPrefixLength(parsed._rawBody),
}) : undefined;
if (bridgeDecision) {
Object.assign(logCtx, {
v2BridgeDecision: bridgeDecision.decision,
...(bridgeDecision.active ? { v2BridgeScope: bridgeDecision.scope } : {}),
});
}
if (bridgeDecision?.active) {
try {
v2RoutedDelegationBridge = injectV2RoutedDelegationBridge(parsed);
if (v2RoutedDelegationBridge) {
copyPreviousResponseReplayProvenance(
parsed._rawBody,
v2RoutedDelegationBridge.requestStateBody,
);
copyBodyNonPersistableMarker(parsed._rawBody, v2RoutedDelegationBridge.requestStateBody);
toolBridgeMaps = buildToolBridgeMaps(parsed, translatorBudget);
}
} catch (error) {
return formatErrorResponse(400, "invalid_request_error", error instanceof Error ? error.message : String(error));
}
}

// Captured before normalization: whether the CLIENT asked for SSE. The
// transport-neutral upstream-streaming policy below may force a bounded JSON
// upstream for reliability (#875); the answer must then be reframed to SSE
Expand Down Expand Up @@ -4700,7 +4763,7 @@ async function handleResponsesInner(
&& (!parsed.previousResponseId || parsed._previousResponseInputExpanded === true);
const rememberPassthroughResponse = passthroughRecordEligible
? (response: { id?: unknown; output?: unknown; status?: unknown }) =>
rememberResponseState(parsed._rawBody, response, undefined, responseStateOptions(true))
rememberResponseState(v2RoutedDelegationBridge?.requestStateBody ?? parsed._rawBody, response, undefined, responseStateOptions(true))
Comment thread
coderabbitai[bot] marked this conversation as resolved.
: undefined;
if (parsed.previousResponseId && !parsed._previousResponseInputExpanded) {
console.warn(
Expand Down Expand Up @@ -4968,8 +5031,15 @@ async function handleResponsesInner(
response: { id?: unknown; output?: unknown; status?: unknown; model?: unknown },
) => {
if (inspectionSawUndeclaredTool) return;
const namespaceRestored = restoreRoutedNamespaceCalls(response, routedNamespaceToolAliases).value;
const bridgeRestored = v2RoutedDelegationBridge
? JSON.parse(rewriteV2RoutedDelegationCallsInJson(
JSON.stringify(namespaceRestored),
v2RoutedDelegationBridge,
)) as typeof response
: namespaceRestored;
const restored = restoreRoutedCustomCalls(
restoreAuthorizedBareNamespaceToolCalls(restoreRoutedNamespaceCalls(response, routedNamespaceToolAliases).value),
restoreAuthorizedBareNamespaceToolCalls(bridgeRestored),
routedCustomToolNames,
routedCustomToolRepairNames,
declaredWireToolNames,
Expand Down Expand Up @@ -5984,7 +6054,9 @@ async function handleResponsesInner(
// injection at the block level, after payload rewrites. Defaults come
// from the finalized OUTBOUND body — the normalized internal tool shapes
// are not the Responses wire shapes the snapshot must mirror.
const bridgeSseRewrite = createV2RoutedDelegationSseRewrite(v2RoutedDelegationBridge);
const blockRewrites = [
bridgeSseRewrite ? payloadRewriteAsBlockRewrite(bridgeSseRewrite) : undefined,
payloadRewrites.length > 0
? payloadRewriteAsBlockRewrite(composeSsePayloadRewrites(...payloadRewrites))
: undefined,
Expand Down Expand Up @@ -6235,8 +6307,12 @@ async function handleResponsesInner(
restoredNamespace,
authorizedBareNamespaceToolAliases,
);
const restored = restoreRoutedCustomCallsInJson(
const bridgeNormalized = rewriteV2RoutedDelegationCallsInJson(
restoredAuthorizedBareNamespace,
v2RoutedDelegationBridge,
);
const restored = restoreRoutedCustomCallsInJson(
bridgeNormalized,
routedCustomToolNames,
routedCustomToolRepairNames,
declaredWireToolNames,
Expand Down
Loading
Loading