diff --git a/src/query.ts b/src/query.ts index 716d6622..997560ba 100644 --- a/src/query.ts +++ b/src/query.ts @@ -1,4 +1,5 @@ // biome-ignore-all assist/source/organizeImports: ANT-ONLY import markers must not be reordered +import { OpenAICodexTurnState } from './services/openaiAuth/turnState.js' import type { ToolResultBlockParam, ToolUseBlock, @@ -264,6 +265,9 @@ async function* queryLoop( skipCacheWrite, } = params const deps = params.deps ?? productionDeps() + // One object for the whole agentic turn, including tool continuations and + // retries. Never put server routing state on the reusable ToolUseContext. + using openAITurnState = new OpenAICodexTurnState(params.toolUseContext.abortController.signal) // Mutable cross-iteration state. The loop body destructures this at the top // of each iteration so reads stay bare-name (`messages`, `toolUseContext`). @@ -699,6 +703,7 @@ async function* queryLoop( c => c.type === 'pending', ), queryTracking, + openAITurnState, effortValue: appState.effortValue, effortValueOverridesEnv: toolUseContext.options.effortValueOverridesEnv, diff --git a/src/query/openaiTurnState.test.ts b/src/query/openaiTurnState.test.ts new file mode 100644 index 00000000..36a233b6 --- /dev/null +++ b/src/query/openaiTurnState.test.ts @@ -0,0 +1,195 @@ +import { afterAll, beforeAll, describe, expect, test } from 'bun:test' +import { randomUUID } from 'node:crypto' +import { mkdtemp, rm } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { z } from 'zod' +import { createSandboxedTestEnvironment } from '../../scripts/pr/test-environment.js' +import type { Tool, ToolUseContext } from '../Tool.js' +import type { QueryParams } from '../query.js' +import type { OpenAICodexTurnState } from '../services/openaiAuth/turnState.js' + +let query: typeof import('../query.js')['query'] +let getDefaultAppState: typeof import('../state/AppStateStore.js')['getDefaultAppState'] +let createAssistantMessage: typeof import('../utils/messages.js')['createAssistantMessage'] +let createUserMessage: typeof import('../utils/messages.js')['createUserMessage'] +let asSystemPrompt: typeof import('../utils/systemPromptType.js')['asSystemPrompt'] +let root: string +let originalCwd: string +let originalEnv: NodeJS.ProcessEnv +let bootstrap: typeof import('../bootstrap/state.js') +let originalPaths: { cwd: string, originalCwd: string, projectRoot: string } + +beforeAll(async () => { + originalCwd = process.cwd() + originalEnv = { ...process.env } + root = await mkdtemp(join(tmpdir(), 'query-openai-turn-state-')) + const env = createSandboxedTestEnvironment(root, { + CLAUDE_CODE_SIMPLE: '1', + CLAUDE_CODE_DISABLE_AUTO_MEMORY: '1', + CLAUDE_CODE_ENABLE_PROMPT_SUGGESTION: '0', + }) + for (const key of Object.keys(process.env)) delete process.env[key] + Object.assign(process.env, env) + bootstrap = await import('../bootstrap/state.js') + originalPaths = { + cwd: bootstrap.getCwdState(), + originalCwd: bootstrap.getOriginalCwd(), + projectRoot: bootstrap.getProjectRoot(), + } + bootstrap.setCwdState(root) + bootstrap.setOriginalCwd(root) + bootstrap.setProjectRoot(root) + process.chdir(root) + query = (await import('../query.js')).query + getDefaultAppState = (await import('../state/AppStateStore.js')).getDefaultAppState + const messages = await import('../utils/messages.js') + createAssistantMessage = messages.createAssistantMessage + createUserMessage = messages.createUserMessage + asSystemPrompt = (await import('../utils/systemPromptType.js')).asSystemPrompt +}) + +afterAll(async () => { + process.chdir(originalCwd) + bootstrap.setCwdState(originalPaths.cwd) + bootstrap.setOriginalCwd(originalPaths.originalCwd) + bootstrap.setProjectRoot(originalPaths.projectRoot) + for (const key of Object.keys(process.env)) delete process.env[key] + Object.assign(process.env, originalEnv) + await rm(root, { recursive: true, force: true }) +}) + +function context(tools: Tool[] = []): ToolUseContext { + let state = getDefaultAppState() + return { + options: { + commands: [], debug: false, mainLoopModel: 'gpt-6-astra', tools, + verbose: false, thinkingConfig: { type: 'disabled' }, + mcpClients: [], mcpResources: {}, isNonInteractiveSession: true, + agentDefinitions: { activeAgents: [], allAgents: [] }, + }, + abortController: new AbortController(), + readFileState: new Map(), + getAppState: () => state, + setAppState: update => { state = update(state) }, + setInProgressToolUseIDs: () => {}, + setResponseLength: () => {}, + updateFileHistoryState: () => {}, + updateAttributionState: () => {}, + messages: [], + } as ToolUseContext +} + +function params( + toolUseContext: ToolUseContext, + callModel: NonNullable['callModel'], +): QueryParams { + return { + messages: [createUserMessage({ content: 'fixture' })], + systemPrompt: asSystemPrompt([]), userContext: {}, systemContext: {}, + canUseTool: async (_tool, input) => ({ behavior: 'allow', updatedInput: input }), + toolUseContext, querySource: 'sdk', maxTurns: 3, + deps: { + callModel, + microcompact: async messages => ({ messages }), + autocompact: async () => ({}), + uuid: randomUUID, + }, + } +} + +async function drain(generator: ReturnType) { + while (true) { + const next = await generator.next() + if (next.done) return next.value + } +} + +describe('query OpenAI routing state lifetime', () => { + test('keeps routing state through a real tool continuation and resets it for the next user turn', async () => { + let toolCalls = 0 + const tool = { + name: 'FixtureTool', inputSchema: z.object({}), maxResultSizeChars: 1000, + isConcurrencySafe: () => true, isReadOnly: () => true, isEnabled: () => true, + userFacingName: () => 'fixture', description: async () => 'fixture', + call: async () => { + toolCalls++ + return { data: 'fixture-result' } + }, + mapToolResultToToolResultBlockParam: (data: string, id: string) => ({ + type: 'tool_result', tool_use_id: id, content: data, + }), + } as unknown as Tool + const toolUseContext = context([tool]) + const states: OpenAICodexTurnState[] = [] + let modelCalls = 0 + const callModel: NonNullable['callModel'] = async function* (request) { + modelCalls++ + expect(request.options.openAITurnState).toBeDefined() + const state = request.options.openAITurnState! + states.push(state) + if (modelCalls === 1) { + expect(state.get()).toBeUndefined() + state.capture('first-user-turn') + yield createAssistantMessage({ content: [{ + type: 'tool_use', id: 'fixture-tool-call', name: 'FixtureTool', input: {}, + }] }) + } else { + if (modelCalls === 2) { + expect(state).toBe(states[0]) + expect(state.get()).toBe('first-user-turn') + expect(toolCalls).toBe(1) + expect(request.messages.some(message => message.type === 'user' + && Array.isArray(message.message.content) + && message.message.content.some(block => block.type === 'tool_result' + && block.tool_use_id === 'fixture-tool-call' + && block.content === 'fixture-result'))).toBe(true) + } else { + expect(state).not.toBe(states[0]) + expect(state.get()).toBeUndefined() + state.capture('second-user-turn') + } + yield createAssistantMessage({ content: 'complete' }) + } + } + + expect(await drain(query(params(toolUseContext, callModel)))).toEqual({ reason: 'completed' }) + expect(modelCalls).toBe(2) + expect(toolUseContext.abortController.signal.aborted).toBe(false) + expect(states[0]!.get()).toBeUndefined() + states[0]!.capture('late-response') + expect(states[0]!.get()).toBeUndefined() + + expect(await drain(query(params(toolUseContext, callModel)))).toEqual({ reason: 'completed' }) + expect(modelCalls).toBe(3) + expect(states[2]!.get()).toBeUndefined() + expect(toolCalls).toBe(1) + }) + + test.each(['abort', 'close'] as const)('clears routing state when a running query is interrupted by %s', async interruption => { + const toolUseContext = context() + let state: OpenAICodexTurnState | undefined + const generator = query(params(toolUseContext, async function* (request) { + state = request.options.openAITurnState + expect(state).toBeDefined() + state!.capture('in-flight-turn') + yield createAssistantMessage({ content: 'partial response' }) + })) + try { + while (!state) { + expect((await generator.next()).done).toBe(false) + } + expect(state.get()).toBe('in-flight-turn') + if (interruption === 'abort') { + toolUseContext.abortController.abort() + expect(state.get()).toBeUndefined() + } + await generator.return({ reason: 'completed' }) + expect(state.get()).toBeUndefined() + state.capture('late-response') + expect(state.get()).toBeUndefined() + } finally { + await generator.return({ reason: 'completed' }) + } + }) +}) diff --git a/src/server/__tests__/providers.test.ts b/src/server/__tests__/providers.test.ts index f84f7cdd..00b87464 100644 --- a/src/server/__tests__/providers.test.ts +++ b/src/server/__tests__/providers.test.ts @@ -16,6 +16,7 @@ import { traceCaptureService, } from '../services/traceCaptureService.js' import type { CreateProviderInput } from '../types/provider.js' +import { buildComputerUseTools } from '../../vendor/computer-use-mcp/tools.js' // ─── Test helpers ───────────────────────────────────────────────────────────── @@ -1605,6 +1606,67 @@ describe('ProviderService', () => { }) describe('handleProxyRequest', () => { + test('preserves optional Computer Use parameters in the final Responses proxy request', async () => { + const originalFetch = globalThis.fetch + const computerTools = buildComputerUseTools().filter(tool => + ['get_app_state', 'click'].includes(tool.name), + ) + const originalSchemas = structuredClone(computerTools.map(tool => tool.inputSchema)) + const calls: Array<{ url: string; body: Record }> = [] + globalThis.fetch = mock(async (input: string | URL | Request, init?: RequestInit) => { + calls.push({ url: String(input), body: JSON.parse(String(init?.body)) }) + return Response.json({ + id: 'resp_computer_schema', + object: 'response', + created_at: 0, + model: 'gpt-6-astra', + status: 'completed', + output: [], + }) + }) as typeof fetch + + try { + const svc = new ProviderService() + const provider = await svc.addProvider(sampleInput({ apiFormat: 'openai_responses' })) + await svc.activateProvider(provider.id) + const req = new Request('http://localhost:3456/proxy/v1/messages', { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ + model: 'gpt-6-astra', + max_tokens: 64, + messages: [{ role: 'user', content: 'Inspect Blender' }], + tools: computerTools.map(tool => ({ + name: tool.name, + description: tool.description, + input_schema: tool.inputSchema, + })), + }), + }) + + const response = await handleProxyRequest(req, new URL(req.url)) + expect(response.status).toBe(200) + await response.text() + expect(calls).toHaveLength(1) + expect(calls[0].url).toBe('https://api.example.com/v1/responses') + const outboundTools = calls[0].body.tools as Array<{ + name: string + strict?: boolean + parameters: Record + }> + expect(outboundTools).toHaveLength(computerTools.length) + for (const [index, tool] of outboundTools.entries()) { + expect(tool.name).toBe(computerTools[index].name) + expect(tool.strict).toBe(false) + expect(tool.parameters).toEqual(originalSchemas[index]) + expect(tool.parameters.required).toEqual(['app']) + } + expect(computerTools.map(tool => tool.inputSchema)).toEqual(originalSchemas) + } finally { + globalThis.fetch = originalFetch + } + }) + test('records a session trace for proxied OpenAI Chat calls', async () => { const originalFetch = globalThis.fetch const upstreamHeaders: Headers[] = [] diff --git a/src/server/__tests__/proxy-transform.test.ts b/src/server/__tests__/proxy-transform.test.ts index d74591d8..46c613d3 100644 --- a/src/server/__tests__/proxy-transform.test.ts +++ b/src/server/__tests__/proxy-transform.test.ts @@ -10,6 +10,7 @@ import { openaiResponsesToAnthropic } from '../proxy/transform/openaiResponsesTo import { stripLeadingBillingHeader } from '../proxy/transform/billingHeader.js' import { openaiUsageToAnthropic } from '../proxy/transform/usage.js' import { resolvePromptCacheKey } from '../proxy/promptCacheKey.js' +import { buildComputerUseTools } from '../../vendor/computer-use-mcp/tools.js' import type { AnthropicRequest, OpenAIChatResponse, OpenAIResponsesResponse } from '../proxy/transform/types.js' const BILLING_HEADER = 'x-anthropic-billing-header: cc_version=2.1.220.693; cc_entrypoint=cli; cch=00000;' @@ -1226,6 +1227,36 @@ describe('openaiChatToAnthropic', () => { // ─── anthropicToOpenaiResponses ───────────────────────────────── describe('anthropicToOpenaiResponses', () => { + test('keeps Computer Use optional fields optional on the Responses wire', () => { + const tools = buildComputerUseTools().map(tool => ({ + name: tool.name, + description: tool.description, + input_schema: tool.inputSchema, + })) + const original = structuredClone(tools) + const wire = JSON.parse(JSON.stringify(anthropicToOpenaiResponses({ + model: 'gpt-6-astra', + max_tokens: 100, + messages: [{ role: 'user', content: 'Read the Blender window' }], + tools, + }))) + + // Omitting strict lets Responses normalize every property to required. + // That made both non-nullable diff aliases mandatory while the executor + // correctly rejected their simultaneous presence (Blender regression). + for (const tool of wire.tools) { + expect(tool.strict).toBe(false) + expect(tool.parameters).toEqual(original.find(t => t.name === tool.name)!.input_schema) + } + const state = wire.tools.find((tool: { name: string }) => tool.name === 'get_app_state') + expect(state.parameters.required).toEqual(['app']) + expect(state.parameters.properties.disableDiff.type).toBe('boolean') + expect(state.parameters.properties.disable_diff.type).toBe('boolean') + const click = wire.tools.find((tool: { name: string }) => tool.name === 'click') + expect(click.parameters.required).toEqual(['app']) + expect(tools).toEqual(original) + }) + test('basic message', () => { const req: AnthropicRequest = { model: 'gpt-4o', @@ -1288,6 +1319,7 @@ describe('anthropicToOpenaiResponses', () => { name: 'get_weather', description: 'Get weather', parameters: { type: 'object', properties: { city: { type: 'string' } } }, + strict: false, }) }) diff --git a/src/server/__tests__/trace-capture.test.ts b/src/server/__tests__/trace-capture.test.ts index 1acc9383..e1cbd3cc 100644 --- a/src/server/__tests__/trace-capture.test.ts +++ b/src/server/__tests__/trace-capture.test.ts @@ -21,6 +21,8 @@ import { } from '../services/traceCaptureService.js' import { sessionService } from '../services/sessionService.js' import { createDumpPromptsFetch } from '../../services/api/dumpPrompts.js' +import { buildOpenAICodexFetch } from '../../services/openaiAuth/fetch.js' +import { clearOpenAIOAuthTokenCache } from '../../services/openaiAuth/storage.js' import { getTraceIndexDatabasePath } from '../services/localIndex/traceDatabase.js' let tmpDir: string @@ -809,6 +811,48 @@ describe('trace capture service', () => { } }) + test('audits trusted OAuth plaintext while the actual transport receives zstd bytes', async () => { + const originalFetch = globalThis.fetch + const overrides = { CC_HAHA_TRACE_API_CALLS: '1', OPENAI_CODEX_OAUTH_FILE: path.join(tmpDir, 'oauth-fixture.json'), CC_HAHA_OPENAI_REQUEST_COMPRESSION: 'true' } + const prior = Object.fromEntries(Object.keys(overrides).map(key => [key, process.env[key]])) + Object.assign(process.env, overrides) + clearOpenAIOAuthTokenCache() + await fs.writeFile(overrides.OPENAI_CODEX_OAUTH_FILE, JSON.stringify({ accessToken: 'fake-access-audit', refreshToken: 'fake-refresh-audit', expiresAt: Date.now() + 3600000 })) + let wireBytes = 0 + let plainBody = '' + try { + globalThis.fetch = (async (_input, init) => { + expect(new Headers(init?.headers).get('content-encoding')).toBe('zstd') + expect(init?.body).toBeInstanceOf(Uint8Array) + wireBytes = (init!.body as Uint8Array).byteLength + plainBody = Buffer.from(await Bun.zstdDecompress(init!.body as Uint8Array)).toString('utf8') + return Response.json({ id: 'resp_zstd_audit', object: 'response', model: 'gpt-6-astra', status: 'completed', output: [] }) + }) as typeof fetch + const traced = createDumpPromptsFetch('zstd-audit', { traceSessionId: 'session-zstd-audit' }) + const codex = buildOpenAICodexFetch(traced, 'test') + await (await codex('https://api.anthropic.com/v1/messages', { + method: 'POST', headers: { 'X-Claude-Code-Session-Id': 'fixture-root' }, + body: JSON.stringify({ model: 'gpt-6-astra', max_tokens: 16, messages: [{ role: 'user', content: 'Audit 中文 '.repeat(1000) }] }), + })).text() + const trace = await waitForTrace('session-zstd-audit', snapshot => Boolean(snapshot.calls[0]?.response)) + expect(trace.calls).toHaveLength(1) + const call = trace.calls[0] + expect(call.request.headers['content-encoding']).toBe('zstd') + expect(call.request.body.preview).toContain('Audit 中文') + expect(call.request.semantic?.request).toMatchObject({ model: 'gpt-6-astra', prompt_cache_key: 'fixture-root' }) + expect(call.metadata).toMatchObject({ requestEncoding: 'zstd', requestPlainBytes: Buffer.byteLength(plainBody), requestWireBytes: wireBytes }) + expect(wireBytes).toBeLessThan(Buffer.byteLength(plainBody)) + expect(JSON.stringify(trace)).not.toContain('fake-access-audit') + } finally { + globalThis.fetch = originalFetch + for (const key of Object.keys(overrides)) { + if (prior[key] === undefined) delete process.env[key] + else process.env[key] = prior[key] + } + clearOpenAIOAuthTokenCache() + } + }) + test('captures direct provider headers when fetch input is a Request', async () => { const originalFetch = globalThis.fetch const originalTraceEnv = process.env.CC_HAHA_TRACE_API_CALLS diff --git a/src/server/proxy/transform/anthropicToOpenaiResponses.ts b/src/server/proxy/transform/anthropicToOpenaiResponses.ts index a401ec0e..451b5d20 100644 --- a/src/server/proxy/transform/anthropicToOpenaiResponses.ts +++ b/src/server/proxy/transform/anthropicToOpenaiResponses.ts @@ -81,6 +81,10 @@ export function anthropicToOpenaiResponses( name: t.name, description: t.description, parameters: t.input_schema, + // Responses otherwise normalizes optional properties to required. + // Preserve Anthropic/MCP omission semantics, including alternative + // selectors and compatibility aliases that cannot coexist. + strict: false, })) if (tools.length > 0) { result.tools = tools diff --git a/src/server/proxy/transform/types.ts b/src/server/proxy/transform/types.ts index aab88f42..b79adfaa 100644 --- a/src/server/proxy/transform/types.ts +++ b/src/server/proxy/transform/types.ts @@ -149,6 +149,7 @@ export type OpenAIResponsesRequest = { name: string description?: string parameters?: Record + strict?: boolean }> tool_choice?: unknown reasoning?: { effort?: OpenAIReasoningEffort } diff --git a/src/services/api/claude.ts b/src/services/api/claude.ts index 61943720..b2e47303 100644 --- a/src/services/api/claude.ts +++ b/src/services/api/claude.ts @@ -1,4 +1,5 @@ -import type { +import { OpenAICodexTurnState } from '../openaiAuth/turnState.js'; +import type { BetaContentBlock, BetaContentBlockParam, BetaImageBlockParam, @@ -754,6 +755,7 @@ export function assistantMessageToMessageParam( } export type Options = { + openAITurnState?: OpenAICodexTurnState; getToolPermissionContext: () => Promise; model: string; toolChoice?: BetaToolChoiceTool | BetaToolChoiceAuto | undefined; @@ -802,6 +804,8 @@ export async function queryModelWithoutStreaming({ signal: AbortSignal; options: Options; }): Promise { + using ownedOpenAITurnState = options.openAITurnState ? undefined : new OpenAICodexTurnState(signal); + options = { ...options, openAITurnState: options.openAITurnState ?? ownedOpenAITurnState }; // Store the assistant message but continue consuming the generator to ensure // logAPISuccessAndDuration gets called (which happens after all yields) let assistantMessage: AssistantMessage | undefined; @@ -853,6 +857,8 @@ export async function* queryModelWithStreaming({ StreamEvent | AssistantMessage | SystemAPIErrorMessage | SystemStreamingFallbackMessage, void > { + using ownedOpenAITurnState = options.openAITurnState ? undefined : new OpenAICodexTurnState(signal); + options = { ...options, openAITurnState: options.openAITurnState ?? ownedOpenAITurnState }; return yield* withStreamingVCR(messages, async function* () { yield* withStreamRetry( () => @@ -911,6 +917,8 @@ export async function* executeNonStreamingRequest( model: string; fetchOverride?: Options["fetchOverride"]; source: string; + openAITurnState?: OpenAICodexTurnState; + agentId?: AgentId; }, retryOptions: { model: string; @@ -938,6 +946,8 @@ export async function* executeNonStreamingRequest( model: clientOptions.model, fetchOverride: clientOptions.fetchOverride, source: clientOptions.source, + openAITurnState: clientOptions.openAITurnState, + agentId: clientOptions.agentId, }), async (anthropic, attempt, context) => { const start = Date.now(); @@ -1956,6 +1966,8 @@ async function* queryModel( model: options.model, fetchOverride: options.fetchOverride, source: options.querySource, + openAITurnState: options.openAITurnState, + agentId: options.agentId, }), async (anthropic, attempt, context) => { attemptNumber = attempt; @@ -3001,7 +3013,7 @@ async function* queryModel( : "other") as AnalyticsMetadata_I_VERIFIED_THIS_IS_NOT_CODE_OR_FILEPATHS, }); const result = yield* executeNonStreamingRequest( - { model: options.model, source: options.querySource }, + { model: options.model, source: options.querySource, openAITurnState: options.openAITurnState, agentId: options.agentId }, { model: options.model, fallbackModel: options.fallbackModel, @@ -3111,7 +3123,7 @@ async function* queryModel( try { // Fall back to non-streaming mode const result = yield* executeNonStreamingRequest( - { model: options.model, source: options.querySource }, + { model: options.model, source: options.querySource, openAITurnState: options.openAITurnState, agentId: options.agentId }, { model: options.model, fallbackModel: options.fallbackModel, diff --git a/src/services/api/client.ts b/src/services/api/client.ts index b2c7c820..343ceb00 100644 --- a/src/services/api/client.ts +++ b/src/services/api/client.ts @@ -1,3 +1,4 @@ +import type { OpenAICodexTurnState } from '../openaiAuth/turnState.js' import Anthropic, { type ClientOptions } from '@anthropic-ai/sdk' import { normalizeAnthropicBaseUrl } from './anthropicBaseUrl.js' import { randomUUID } from 'crypto' @@ -204,12 +205,16 @@ export async function getAnthropicClient({ model, fetchOverride, source, + openAITurnState, + agentId, }: { apiKey?: string maxRetries: number model?: string fetchOverride?: ClientOptions['fetch'] source?: string + openAITurnState?: OpenAICodexTurnState + agentId?: string }): Promise { const containerId = process.env.CLAUDE_CODE_CONTAINER_ID const remoteSessionId = process.env.CLAUDE_CODE_REMOTE_SESSION_ID @@ -272,7 +277,7 @@ export async function getAnthropicClient({ const resolvedFetch = usingGrok ? buildGrokFetch(fetchOverride, source) : usingOpenAICodex - ? buildOpenAICodexFetch(fetchOverride, source) + ? buildOpenAICodexFetch(fetchOverride, source, openAITurnState, agentId) : buildFetch(fetchOverride, source) const stagingOAuthBaseUrl = process.env.USER_TYPE === 'ant' && isEnvTruthy(process.env.USE_STAGING_OAUTH) diff --git a/src/services/api/dumpPrompts.ts b/src/services/api/dumpPrompts.ts index 8f44fa7e..27d51a6d 100644 --- a/src/services/api/dumpPrompts.ts +++ b/src/services/api/dumpPrompts.ts @@ -1,4 +1,5 @@ import type { ClientOptions } from '@anthropic-ai/sdk' +import { getRequestBodyAudit } from './requestBodyAudit.js' import { createHash } from 'crypto' import { promises as fs } from 'fs' import { dirname, join } from 'path' @@ -240,15 +241,23 @@ export function createDumpPromptsFetch( let timestamp: string | undefined let traceStartedAtMs = 0 let traceRequestBody: unknown + const bodyAudit = getRequestBodyAudit(init?.body) + const encodingMetadata = bodyAudit ? { + requestEncoding: bodyAudit.requestEncoding, + requestPlainBytes: bodyAudit.requestPlainBytes, + requestWireBytes: bodyAudit.requestWireBytes, + } : {} if (init?.method === 'POST' && init.body) { timestamp = new Date().toISOString() traceStartedAtMs = Date.now() - traceRequestBody = init.body + traceRequestBody = bodyAudit?.plainBody ?? init.body // Parsing + stringifying the request (system prompt + tool schemas = MBs) // takes hundreds of ms. Defer so it doesn't block the actual API call — // this is debug tooling for /issue, not on the critical path. - setImmediate(dumpRequest, init.body as string, timestamp, state, filePath) + if (typeof traceRequestBody === 'string') { + setImmediate(dumpRequest, traceRequestBody, timestamp, state, filePath) + } } const requestUrl = getRequestUrl(input) @@ -298,6 +307,7 @@ export function createDumpPromptsFetch( bodySnapshot: createRequestPendingSnapshot(traceRequestBody), }, metadata: { + ...encodingMetadata, phase: 'api_call_started', }, }) @@ -312,6 +322,7 @@ export function createDumpPromptsFetch( severity: 'info', title: 'API call started', metadata: { + ...encodingMetadata, url: traceRequestUrl, }, }) @@ -344,6 +355,7 @@ export function createDumpPromptsFetch( }, error: err, metadata: { + ...encodingMetadata, phase: 'api_call_failed', ...(aborted ? { aborted: true } : {}), }, @@ -360,6 +372,7 @@ export function createDumpPromptsFetch( title: 'API call failed', message: err instanceof Error ? err.message : String(err), metadata: { + ...encodingMetadata, url: traceRequestUrl, ...(aborted ? { aborted: true } : {}), }, @@ -419,6 +432,7 @@ export function createDumpPromptsFetch( bodySnapshot: capture.snapshot, }, metadata: { + ...encodingMetadata, phase: 'api_call_completed', }, }) @@ -428,6 +442,7 @@ export function createDumpPromptsFetch( severity: response.ok ? 'info' : 'warning', title: 'API call completed', metadata: { + ...encodingMetadata, status: response.status, url: traceRequestUrl, }, @@ -451,6 +466,7 @@ export function createDumpPromptsFetch( }, error: abortError, metadata: { + ...encodingMetadata, phase: 'api_call_aborted', aborted: true, }, @@ -463,6 +479,7 @@ export function createDumpPromptsFetch( title: 'API call aborted', message: abortError.message, metadata: { + ...encodingMetadata, status: response.status, url: traceRequestUrl, durationMs, @@ -480,6 +497,7 @@ export function createDumpPromptsFetch( title: 'Response capture failed', message: captureFailure instanceof Error ? captureFailure.message : String(captureFailure), metadata: { + ...encodingMetadata, status: response.status, url: traceRequestUrl, }, @@ -498,6 +516,7 @@ export function createDumpPromptsFetch( }, error: captureFailure, metadata: { + ...encodingMetadata, phase: 'response_capture_failed', responseCaptureFailed: true, }, diff --git a/src/services/api/requestBodyAudit.ts b/src/services/api/requestBodyAudit.ts new file mode 100644 index 00000000..814ebc78 --- /dev/null +++ b/src/services/api/requestBodyAudit.ts @@ -0,0 +1,22 @@ +/** Only locally encoded bodies are associated with plaintext; never decompress caller input. */ +export type RequestBodyAudit = Readonly<{ + plainBody: string + requestEncoding: 'zstd' + requestPlainBytes: number + requestWireBytes: number +}> + +const encodedBodies = new WeakMap() + +export function registerEncodedRequestBody(body: Uint8Array, plainBody: string): void { + encodedBodies.set(body, Object.freeze({ + plainBody, + requestEncoding: 'zstd', + requestPlainBytes: Buffer.byteLength(plainBody), + requestWireBytes: body.byteLength, + })) +} + +export function getRequestBodyAudit(body: unknown): RequestBodyAudit | undefined { + return body instanceof Uint8Array ? encodedBodies.get(body) : undefined +} diff --git a/src/services/openaiAuth/fetch.test.ts b/src/services/openaiAuth/fetch.test.ts index e97d69fc..220353b0 100644 --- a/src/services/openaiAuth/fetch.test.ts +++ b/src/services/openaiAuth/fetch.test.ts @@ -1,21 +1,35 @@ -import { afterEach, beforeEach, describe, expect, test } from 'bun:test' +import { afterEach, beforeEach, describe, expect, spyOn, test } from 'bun:test' import * as fs from 'fs/promises' import * as os from 'os' import * as path from 'path' -import { configureEffortParams } from '../api/claude.js' +import { configureEffortParams, executeNonStreamingRequest } from '../api/claude.js' +import { asAgentId } from '../../types/ids.js' import { OPENAI_CODEX_API_ENDPOINT } from './client.js' import { buildOpenAICodexFetch } from './fetch.js' +import { OpenAICodexTurnState } from './turnState.js' import { OPENAI_CODEX_REASONING_EFFORT_ENV_KEY } from './models.js' import { clearOpenAIOAuthTokenCache } from './storage.js' +import { encodeOpenAIReasoningEnvelope } from '../../server/proxy/transform/openaiReasoning.js' +import { buildComputerUseTools } from '../../vendor/computer-use-mcp/tools.js' + +function readWireBody(init?: RequestInit): Record { + const body = init?.body + return JSON.parse(body instanceof Uint8Array + ? Buffer.from(Bun.zstdDecompressSync(body)).toString('utf8') + : String(body)) +} describe('buildOpenAICodexFetch', () => { let tmpDir: string let originalTokenFile: string | undefined let originalReasoningEffort: string | undefined + let originalCompression: string | undefined beforeEach(async () => { tmpDir = await fs.mkdtemp(path.join(os.tmpdir(), 'openai-codex-fetch-')) originalTokenFile = process.env.OPENAI_CODEX_OAUTH_FILE + originalCompression = process.env.CC_HAHA_OPENAI_REQUEST_COMPRESSION + delete process.env.CC_HAHA_OPENAI_REQUEST_COMPRESSION originalReasoningEffort = process.env[OPENAI_CODEX_REASONING_EFFORT_ENV_KEY] delete process.env[OPENAI_CODEX_REASONING_EFFORT_ENV_KEY] process.env.OPENAI_CODEX_OAUTH_FILE = path.join(tmpDir, 'openai-oauth.json') @@ -34,6 +48,8 @@ describe('buildOpenAICodexFetch', () => { }) afterEach(async () => { + if (originalCompression === undefined) delete process.env.CC_HAHA_OPENAI_REQUEST_COMPRESSION + else process.env.CC_HAHA_OPENAI_REQUEST_COMPRESSION = originalCompression if (originalTokenFile === undefined) { delete process.env.OPENAI_CODEX_OAUTH_FILE } else { @@ -77,7 +93,7 @@ describe('buildOpenAICodexFetch', () => { upstreamCalls.push({ url: String(input), headers: Object.fromEntries(headers.entries()), - body: JSON.parse(String(init?.body)) as Record, + body: readWireBody(init) as Record, proxy: (init as RequestInit & { proxy?: string } | undefined)?.proxy, }) return Response.json({ @@ -124,13 +140,421 @@ describe('buildOpenAICodexFetch', () => { }) }) + test('preserves optional Computer Use parameters in the final Codex OAuth request', async () => { + const computerTools = buildComputerUseTools().filter(tool => + ['get_app_state', 'click'].includes(tool.name), + ) + const originalSchemas = structuredClone(computerTools.map(tool => tool.inputSchema)) + const calls: Array<{ url: string; body: Record }> = [] + const codexFetch = buildOpenAICodexFetch(async (input, init) => { + calls.push({ url: String(input), body: readWireBody(init) }) + return Response.json({ + id: 'resp_computer_schema', + object: 'response', + created_at: 0, + model: 'gpt-6-astra', + status: 'completed', + output: [], + }) + }, 'test') + + const response = await codexFetch('https://api.anthropic.com/v1/messages', { + method: 'POST', + body: JSON.stringify({ + model: 'gpt-6-astra', + max_tokens: 64, + messages: [{ role: 'user', content: 'Inspect Blender' }], + tools: computerTools.map(tool => ({ + name: tool.name, + description: tool.description, + input_schema: tool.inputSchema, + })), + }), + }) + expect(response.status).toBe(200) + await response.text() + expect(calls).toHaveLength(1) + expect(calls[0].url).toBe(OPENAI_CODEX_API_ENDPOINT) + const outboundTools = calls[0].body.tools as Array<{ + name: string + strict?: boolean + parameters: Record + }> + expect(outboundTools).toHaveLength(computerTools.length) + for (const [index, tool] of outboundTools.entries()) { + expect(tool.name).toBe(computerTools[index].name) + expect(tool.strict).toBe(false) + expect(tool.parameters).toEqual(originalSchemas[index]) + expect(tool.parameters.required).toEqual(['app']) + } + expect(computerTools.map(tool => tool.inputSchema)).toEqual(originalSchemas) + }) + + test('routes repeated session requests to a stable cache key without changing multimodal reasoning or tool schemas', async () => { + const bodies: Array> = [] + const codexFetch = buildOpenAICodexFetch(async (_input, init) => { + bodies.push(readWireBody(init)) + return Response.json({ id: 'resp_cache', object: 'response', status: 'completed', output: [] }) + }, 'test') + const reasoning = { type: 'reasoning' as const, id: 'rs_cache', summary: [], encrypted_content: 'opaque-test-reasoning' } + const schema = { type: 'object', properties: { app: { type: 'string' }, disableDiff: { type: 'boolean' } }, required: ['app'] } + const body = JSON.stringify({ + model: 'gpt-6-astra', max_tokens: 64, + output_config: { effort: 'high' }, + metadata: { user_id: JSON.stringify({ session_id: 'sdk-json-metadata' }) }, + messages: [ + { role: 'assistant', content: [ + { type: 'redacted_thinking', data: encodeOpenAIReasoningEnvelope(reasoning) }, + { type: 'tool_use', id: 'call_image', name: 'get_app_state', input: { app: 'Blender' } }, + ] }, + { role: 'user', content: [{ type: 'tool_result', tool_use_id: 'call_image', content: [ + { type: 'text', text: 'Blender screenshot' }, + { type: 'image', source: { type: 'base64', media_type: 'image/png', data: 'AA==' } }, + ] }] }, + ], + tools: [{ name: 'get_app_state', description: 'Observe app', input_schema: schema }], + }) + for (const sessionId of ['session-one', 'session-one', 'session-two']) { + const response = await codexFetch('https://api.anthropic.com/v1/messages', { + method: 'POST', headers: { 'X-Claude-Code-Session-Id': sessionId }, body, + }) + await response.text() + } + expect(bodies.map(item => item.prompt_cache_key)).toEqual(['session-one', 'session-one', 'session-two']) + for (const outgoing of bodies) { + expect(outgoing.reasoning).toEqual({ effort: 'high' }) + expect(outgoing.include).toEqual(['reasoning.encrypted_content']) + expect(outgoing.input).toEqual([ + reasoning, + { type: 'function_call', call_id: 'call_image', name: 'get_app_state', arguments: JSON.stringify({ app: 'Blender' }) }, + { type: 'function_call_output', call_id: 'call_image', output: [ + { type: 'input_text', text: 'Blender screenshot' }, + { type: 'input_image', image_url: 'data:image/png;base64,AA==' }, + ] }, + ]) + expect(outgoing.tools).toMatchObject([{ strict: false, parameters: schema }]) + } + }) + + test('uses existing cache identity precedence and Request headers without inventing a key', async () => { + const keys: unknown[] = [] + const codexFetch = buildOpenAICodexFetch(async (_input, init) => { + keys.push(readWireBody(init).prompt_cache_key) + return Response.json({ id: 'resp_cache_identity', object: 'response', status: 'completed', output: [] }) + }, 'test') + for (const entry of [ + { metadata: { user_id: 'user_test_session_metadata-session' }, header: 'header-session', initHeader: 'init-session' }, + { metadata: { session_id: 'metadata-field-session' }, header: 'header-session' }, + { header: 'request-session' }, + { header: 'request-session', initHeader: 'init-session' }, + {}, + ]) { + const request = new Request('https://api.anthropic.com/v1/messages', { + method: 'POST', + headers: entry.header ? { 'x-claude-code-session-id': entry.header } : {}, + body: JSON.stringify({ model: 'gpt-6-astra', max_tokens: 64, metadata: entry.metadata, messages: [{ role: 'user', content: 'Hello' }] }), + }) + const response = await codexFetch(request, entry.initHeader ? { headers: { 'X-Claude-Code-Session-Id': entry.initHeader } } : undefined) + await response.text() + } + expect(keys).toEqual(['metadata-session', 'metadata-field-session', 'request-session', 'init-session', undefined]) + }) + + test('replays the first server routing state across same-turn client recreation without rotating it', async () => { + using state = new OpenAICodexTurnState(new AbortController().signal) + const sent: Array = [] + const upstream: typeof fetch = async (_input, init) => { + sent.push(new Headers(init?.headers).get('x-codex-turn-state')) + return Response.json({ id: 'resp_state', object: 'response', status: 'completed', output: [] }, { + headers: { 'x-codex-turn-state': `server-state-${sent.length}` }, + }) + } + for (let attempt = 0; attempt < 3; attempt++) { + // Auth retry/client recreation and tool continuation share the explicit turn object. + const fetch = buildOpenAICodexFetch(upstream, 'test', state) + await (await fetch('https://api.anthropic.com/v1/messages', { + method: 'POST', body: JSON.stringify({ model: 'gpt-6-astra', messages: [{ role: 'user', content: 'Hello' }] }), + })).text() + } + expect(sent).toEqual([null, 'server-state-1', 'server-state-1']) + expect(JSON.stringify(state)).toBe('{}') + }) + + test('isolates concurrent turns with the same session key and discards canceled or finished state', async () => { + const controller = new AbortController() + using first = new OpenAICodexTurnState(controller.signal) + using second = new OpenAICodexTurnState(new AbortController().signal) + const firstSent: Array = [] + const secondSent: Array = [] + const upstream = (name: string, sent: Array): typeof fetch => async (_input, init) => { + sent.push(new Headers(init?.headers).get('x-codex-turn-state')) + return Response.json({ id: 'resp_isolated', object: 'response', status: 'completed', output: [] }, { + headers: { 'x-codex-turn-state': name }, + }) + } + const request = async (state: OpenAICodexTurnState, name: string, sent: Array) => { + const fetch = buildOpenAICodexFetch(upstream(name, sent), 'test', state) + await (await fetch('https://api.anthropic.com/v1/messages', { + method: 'POST', headers: { 'X-Claude-Code-Session-Id': 'shared-session' }, + body: JSON.stringify({ model: 'gpt-6-astra', messages: [{ role: 'user', content: 'Hello' }] }), + })).text() + } + await Promise.all([request(first, 'first', firstSent), request(second, 'second', secondSent)]) + await Promise.all([request(first, 'first-next', firstSent), request(second, 'second-next', secondSent)]) + expect(firstSent).toEqual([null, 'first']) + expect(secondSent).toEqual([null, 'second']) + controller.abort() + expect(first.get()).toBeUndefined() + first.capture('late-canceled-response') + expect(first.get()).toBeUndefined() + second[Symbol.dispose]() + second.capture('late-finished-response') + expect(second.get()).toBeUndefined() + using nextTurn = new OpenAICodexTurnState(new AbortController().signal) + await request(nextTurn, 'new-turn', firstSent) + expect(firstSent).toEqual([null, 'first', null]) + }) + + test('does not capture routing state from an error or a response arriving after cancellation', async () => { + const controller = new AbortController() + using state = new OpenAICodexTurnState(controller.signal) + const sent: Array = [] + let attempt = 0 + const fetch = buildOpenAICodexFetch(async (_input, init) => { + sent.push(new Headers(init?.headers).get('x-codex-turn-state')) + if (attempt++ === 0) return Response.json({ error: 'temporary' }, { status: 503, headers: { 'x-codex-turn-state': 'error-state' } }) + controller.abort() + return Response.json({ id: 'late', object: 'response', status: 'completed', output: [] }, { headers: { 'x-codex-turn-state': 'late-state' } }) + }, 'test', state) + for (let i = 0; i < 2; i++) await (await fetch('https://api.anthropic.com/v1/messages', { + method: 'POST', body: JSON.stringify({ model: 'gpt-6-astra', messages: [{ role: 'user', content: 'Hello' }] }), + })).text() + expect(sent).toEqual([null, null]) + expect(state.get()).toBeUndefined() + }) + + test('keeps turn routing through streaming SDK cleanup while a real root-turn abort clears it', async () => { + const root = new AbortController() + const request = new AbortController() + using state = new OpenAICodexTurnState(root.signal) + const sent: Array = [] + const fetch = buildOpenAICodexFetch(async (_input, init) => { + sent.push(new Headers(init?.headers).get('x-codex-turn-state')) + return new Response([ + 'event: response.completed', + 'data: {"response":{"id":"resp_sticky_stream","object":"response","model":"gpt-6-astra","status":"completed","output":[],"usage":{"input_tokens":1,"output_tokens":0,"total_tokens":1}}}', + '', '', + ].join('\n'), { headers: { 'Content-Type': 'text/event-stream', 'x-codex-turn-state': 'first-stream-state' } }) + }, 'test', state) + const body = JSON.stringify({ model: 'gpt-6-astra', stream: true, messages: [{ role: 'user', content: 'Hello' }] }) + await (await fetch('https://api.anthropic.com/v1/messages', { method: 'POST', body, signal: request.signal })).text() + request.abort() // SDK stream disposal is not the end of the agentic turn. + expect(state.get()).toBe('first-stream-state') + await (await fetch('https://api.anthropic.com/v1/messages', { method: 'POST', body })).text() + expect(sent).toEqual([null, 'first-stream-state']) + root.abort() + expect(state.get()).toBeUndefined() + }) + + test('threads routing state through the actual SDK client factory across tool continuations', async () => { + const { getAnthropicClient } = await import('../api/client.js') + const overrides = { CC_HAHA_OPENAI_OAUTH_PROVIDER: '1', CC_HAHA_GROK_OAUTH_PROVIDER: '', CLAUDE_CONFIG_DIR: tmpDir, CLAUDE_CODE_SIMPLE: '1' } + const previous = Object.fromEntries(Object.keys(overrides).map(key => [key, process.env[key]])) + Object.assign(process.env, overrides) + using state = new OpenAICodexTurnState(new AbortController().signal) + const sent: Array = [] + try { + for (let i = 0; i < 2; i++) { + const client = await getAnthropicClient({ + maxRetries: 0, model: 'gpt-6-astra', openAITurnState: state, + fetchOverride: async (_input, init) => { + sent.push(new Headers(init?.headers).get('x-codex-turn-state')) + return Response.json({ id: 'resp_factory', object: 'response', status: 'completed', output: [] }, { + headers: { 'x-codex-turn-state': 'factory-turn-state' }, + }) + }, + }) + await client.messages.create({ model: 'gpt-6-astra', max_tokens: 16, messages: [{ role: 'user', content: 'Hello' }] }) + } + expect(sent).toEqual([null, 'factory-turn-state']) + } finally { + for (const key of Object.keys(overrides)) { + if (previous[key] === undefined) delete process.env[key] + else process.env[key] = previous[key] + } + } + }) + + test('sends stable root and branch identity independent of cache overrides and turn state', async () => { + const sent: Array<{ session: string | null; thread: string | null; cache: string | undefined }> = [] + const upstream: typeof fetch = async (_input, init) => { + const headers = new Headers(init?.headers) + sent.push({ session: headers.get('session-id'), thread: headers.get('thread-id'), cache: readWireBody(init).prompt_cache_key }) + return Response.json({ id: 'resp_identity', object: 'response', status: 'completed', output: [] }) + } + const root = '01234567-89ab-4cde-8fab-0123456789ab' + const other = '11234567-89ab-4cde-8fab-0123456789ab' + for (const [session, agent] of [[root, undefined], [root, 'a0000000000000001'], [root, 'a0000000000000001'], [root, 'a0000000000000002'], [other, 'a0000000000000001']] as const) { + using state = new OpenAICodexTurnState(new AbortController().signal) + const codexFetch = buildOpenAICodexFetch(upstream, 'test', state, agent) + await (await codexFetch('https://api.anthropic.com/v1/messages', { + method: 'POST', headers: { 'X-Claude-Code-Session-Id': session }, + body: JSON.stringify({ model: 'gpt-6-astra', max_tokens: 16, metadata: { session_id: 'explicit-cache-override' }, messages: [{ role: 'user', content: 'Hello' }] }), + })).text() + } + expect(sent[0]).toEqual({ session: root, thread: root, cache: 'explicit-cache-override' }) + expect(sent[1]).toEqual(sent[2]) + expect(sent[1].session).toBe(root) + expect(sent[1].thread).toMatch(/^[0-9a-f]{8}-[0-9a-f]{4}-5[0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/) + expect(sent[1].thread).not.toBe(root) + expect(sent[3].thread).not.toBe(sent[1].thread) + expect(sent[4].session).toBe(other) + expect(sent[4].thread).not.toBe(sent[1].thread) + }) + + test('reads identity from Request headers with RequestInit precedence and never invents missing identity', async () => { + const sent: Array<[string | null, string | null]> = [] + const codexFetch = buildOpenAICodexFetch(async (_input, init) => { + const headers = new Headers(init?.headers) + sent.push([headers.get('session-id'), headers.get('thread-id')]) + return Response.json({ id: 'resp_identity_headers', object: 'response', status: 'completed', output: [] }) + }, 'test') + const body = JSON.stringify({ model: 'gpt-6-astra', metadata: { session_id: 'cache-only' }, messages: [{ role: 'user', content: 'Hello' }] }) + const request = () => new Request('https://api.anthropic.com/v1/messages', { method: 'POST', headers: { 'X-Claude-Code-Session-Id': 'request-root' }, body }) + await (await codexFetch(request())).text() + await (await codexFetch(request(), { headers: { 'X-Claude-Code-Session-Id': 'init-root' } })).text() + await (await codexFetch('https://api.anthropic.com/v1/messages', { method: 'POST', body })).text() + expect(sent).toEqual([['request-root', 'request-root'], ['init-root', 'init-root'], [null, null]]) + }) + + test('forwards explicit agent identity through recreated SDK clients and SDK retries', async () => { + const { getAnthropicClient } = await import('../api/client.js') + const overrides = { CC_HAHA_OPENAI_OAUTH_PROVIDER: '1', CC_HAHA_GROK_OAUTH_PROVIDER: '', CLAUDE_CONFIG_DIR: tmpDir, CLAUDE_CODE_SIMPLE: '1' } + const previous = Object.fromEntries(Object.keys(overrides).map(key => [key, process.env[key]])) + Object.assign(process.env, overrides) + const sent: Array<[string | null, string | null]> = [] + try { + for (const agentId of ['a0000000000000001', 'a0000000000000001', 'a0000000000000002', undefined]) { + const client = await getAnthropicClient({ + maxRetries: 1, model: 'gpt-6-astra', agentId, + fetchOverride: async (_input, init) => { + const headers = new Headers(init?.headers) + sent.push([headers.get('session-id'), headers.get('thread-id')]) + if (sent.length === 1) return Response.json({ error: { message: 'retry fixture' } }, { status: 500, headers: { 'retry-after-ms': '1' } }) + return Response.json({ id: 'resp_identity_factory', object: 'response', status: 'completed', output: [] }) + }, + }) + await client.messages.create({ model: 'gpt-6-astra', max_tokens: 16, messages: [{ role: 'user', content: 'Hello' }] }) + } + expect(sent).toHaveLength(5) + expect(sent[0][0]).toBeTruthy() + expect(sent[0][1]).not.toBe(sent[0][0]) + expect(sent[0]).toEqual(sent[1]) + expect(sent[1]).toEqual(sent[2]) + expect(sent[3][0]).toBe(sent[0][0]) + expect(sent[3][1]).not.toBe(sent[0][1]) + expect(sent[4][1]).toBe(sent[4][0]) + } finally { + for (const key of Object.keys(overrides)) { + if (previous[key] === undefined) delete process.env[key] + else process.env[key] = previous[key] + } + } + }) + + test('retains branch identity through the non-streaming fallback client factory', async () => { + const overrides = { CC_HAHA_OPENAI_OAUTH_PROVIDER: '1', CC_HAHA_GROK_OAUTH_PROVIDER: '', CLAUDE_CONFIG_DIR: tmpDir, CLAUDE_CODE_SIMPLE: '1' } + const previous = Object.fromEntries(Object.keys(overrides).map(key => [key, process.env[key]])) + Object.assign(process.env, overrides) + const sent: Array<[string | null, string | null]> = [] + try { + const run = executeNonStreamingRequest({ + model: 'gpt-6-astra', source: 'test', agentId: asAgentId('a0000000000000001'), + fetchOverride: async (_input, init) => { + const headers = new Headers(init?.headers) + sent.push([headers.get('session-id'), headers.get('thread-id')]) + return Response.json({ id: 'resp_fallback_identity', object: 'response', status: 'completed', output: [] }) + }, + }, { model: 'gpt-6-astra', thinkingConfig: { type: 'disabled' }, signal: new AbortController().signal }, + () => ({ model: 'gpt-6-astra', max_tokens: 16, messages: [{ role: 'user', content: 'Hello' }] }), + () => {}, () => {}) + let next = await run.next() + while (!next.done) next = await run.next() + expect(sent).toHaveLength(1) + expect(sent[0][0]).toBeTruthy() + expect(sent[0][1]).toBeTruthy() + expect(sent[0][1]).not.toBe(sent[0][0]) + } finally { + for (const key of Object.keys(overrides)) { + if (previous[key] === undefined) delete process.env[key] + else process.env[key] = previous[key] + } + } + }) + + test('compresses the final OAuth wire request with zstd without replaying transport failures', async () => { + let calls = 0 + const codexFetch = buildOpenAICodexFetch(async (_input, init) => { + calls++ + expect(new Headers(init?.headers).get('content-encoding')).toBe('zstd') + expect(init?.body).toBeInstanceOf(Uint8Array) + const decoded = await Bun.zstdDecompress(init!.body as Uint8Array) + const body = JSON.parse(Buffer.from(decoded).toString('utf8')) + expect(body.model).toBe('gpt-6-astra') + throw new Error('transport fixture failure') + }, 'test') + await expect(codexFetch('https://api.anthropic.com/v1/messages', { + method: 'POST', body: JSON.stringify({ model: 'gpt-6-astra', messages: [{ role: 'user', content: 'Hello' }] }), + })).rejects.toThrow('transport fixture failure') + expect(calls).toBe(1) + }) + + test('never submits after cancellation during compression and only falls back before transport', async () => { + let release!: (body: Uint8Array) => void + const compressor = spyOn(Bun, 'zstdCompress').mockImplementation(() => new Promise(resolve => { release = resolve })) + let submitted = 0 + const codexFetch = buildOpenAICodexFetch(async (_input, init) => { + submitted++ + expect(typeof init?.body).toBe('string') + expect(new Headers(init?.headers).has('content-encoding')).toBe(false) + return Response.json({ id: 'resp_local_fallback', object: 'response', status: 'completed', output: [] }) + }, 'test') + const body = JSON.stringify({ model: 'gpt-6-astra', messages: [{ role: 'user', content: 'Hello' }] }) + try { + const abort = new AbortController() + const pending = codexFetch('https://api.anthropic.com/v1/messages', { method: 'POST', body, signal: abort.signal }) + while (!release) await new Promise(resolve => setTimeout(resolve, 1)) + abort.abort() + release(new Uint8Array([1])) + await expect(pending).rejects.toThrow() + expect(submitted).toBe(0) + compressor.mockRejectedValue(new Error('local codec failure')) + await (await codexFetch('https://api.anthropic.com/v1/messages', { method: 'POST', body })).text() + expect(submitted).toBe(1) + } finally { compressor.mockRestore() } + }) + + test('leaves non-messages requests untouched without calling the encoder', async () => { + const encoder = spyOn(Bun, 'zstdCompress') + const init = { method: 'POST', body: 'unchanged fixture' } + try { + const codexFetch = buildOpenAICodexFetch(async (input, received) => { + expect(String(input)).toBe('https://example.test/other') + expect(received).toBe(init) + return new Response('ok') + }, 'test') + await (await codexFetch('https://example.test/other', init)).text() + expect(encoder).not.toHaveBeenCalled() + } finally { encoder.mockRestore() } + }) + test('uses streamed Codex responses even for non-streaming Anthropic callers', async () => { const upstreamCalls: Array<{ url: string body: Record }> = [] const fetchOverride: typeof fetch = async (input, init) => { - const body = JSON.parse(String(init?.body)) as Record + const body = readWireBody(init) as Record upstreamCalls.push({ url: String(input), body, @@ -169,7 +593,7 @@ describe('buildOpenAICodexFetch', () => { test('marks streaming responses as OpenAI OAuth and preserves encrypted reasoning', async () => { const upstreamBodies: Array> = [] const fetchOverride: typeof fetch = async (_input, init) => { - upstreamBodies.push(JSON.parse(String(init?.body)) as Record) + upstreamBodies.push(readWireBody(init) as Record) return new Response([ 'event: response.created', 'data: {"response":{"id":"resp_reasoning","model":"gpt-5.6-terra","status":"in_progress"}}', @@ -277,7 +701,7 @@ describe('buildOpenAICodexFetch', () => { test('applies defaults and validates request and session efforts for the final request', async () => { const upstreamBodies: Array> = [] const fetchOverride: typeof fetch = async (_input, init) => { - const body = JSON.parse(String(init?.body)) as Record + const body = readWireBody(init) as Record upstreamBodies.push(body) return Response.json({ id: `resp_${upstreamBodies.length}`, @@ -336,7 +760,7 @@ describe('buildOpenAICodexFetch', () => { test('keeps Agent request effort above Desktop session effort without synthesizing high', async () => { const upstreamBodies: Array> = [] const fetchOverride: typeof fetch = async (_input, init) => { - const body = JSON.parse(String(init?.body)) as Record + const body = readWireBody(init) as Record upstreamBodies.push(body) return Response.json({ id: `resp_${upstreamBodies.length}`, @@ -400,7 +824,7 @@ describe('buildOpenAICodexFetch', () => { process.env[OPENAI_CODEX_REASONING_EFFORT_ENV_KEY] = 'high' const upstreamBodies: Array> = [] const fetchOverride: typeof fetch = async (_input, init) => { - const body = JSON.parse(String(init?.body)) as Record + const body = readWireBody(init) as Record // Let the two requests overlap and finish in the opposite order from // their launch. Per-request effort must not depend on process.env writes. await new Promise(resolve => diff --git a/src/services/openaiAuth/fetch.ts b/src/services/openaiAuth/fetch.ts index 430e0c7a..b2f9bc89 100644 --- a/src/services/openaiAuth/fetch.ts +++ b/src/services/openaiAuth/fetch.ts @@ -1,5 +1,8 @@ import type { ClientOptions } from '@anthropic-ai/sdk' import { randomUUID } from 'crypto' +import type { OpenAICodexTurnState } from './turnState.js' +import { resolveOpenAIRequestIdentity } from './requestIdentity.js' +import { encodeOpenAIRequestBody } from './requestCompression.js' import { getOpenAIPolicyError } from './policyError.js' import { OPENAI_CODEX_API_ENDPOINT, @@ -14,6 +17,7 @@ import { resolveOpenAIReasoningEffortWithPriority, } from './models.js' import { getOpenAIOAuthTokens } from './storage.js' +import { resolvePromptCacheKey } from '../../server/proxy/promptCacheKey.js' import { anthropicToOpenaiResponses } from '../../server/proxy/transform/anthropicToOpenaiResponses.js' import { openaiResponsesToAnthropic } from '../../server/proxy/transform/openaiResponsesToAnthropic.js' import { openaiResponsesStreamToAnthropic } from '../../server/proxy/streaming/openaiResponsesStreamToAnthropic.js' @@ -32,6 +36,8 @@ export function shouldUseOpenAICodexAuth(): boolean { export function buildOpenAICodexFetch( fetchOverride: ClientOptions['fetch'], source: string | undefined, + turnState?: OpenAICodexTurnState, + agentId?: string, ): ClientOptions['fetch'] { const inner = fetchOverride ?? globalThis.fetch @@ -44,12 +50,15 @@ export function buildOpenAICodexFetch( const originalBody = await readAnthropicBody(input, init) const mappedModel = resolveOpenAICodexModel(originalBody.model) + const sessionId = readSessionId(input, init) + const identity = resolveOpenAIRequestIdentity(sessionId, agentId) + const cacheKey = resolvePromptCacheKey(originalBody, sessionId) const transformedBody = anthropicToOpenaiResponses( { ...originalBody, model: mappedModel, }, - { preserveOpenAIReasoning: true }, + { preserveOpenAIReasoning: true, cacheKey }, ) // Keep a valid native request-scoped value ahead of the transformed value, // the session env, and the model default. The generic transformer preserves @@ -91,6 +100,12 @@ export function buildOpenAICodexFetch( headers.set('Authorization', `Bearer ${tokens.accessToken}`) headers.set('originator', OPENAI_CODEX_ORIGINATOR) headers.set('User-Agent', OPENAI_CODEX_TOKEN_USER_AGENT) + if (identity) { + headers.set('session-id', identity.sessionId) + headers.set('thread-id', identity.threadId) + } + const routingState = turnState?.get() + if (routingState) headers.set('x-codex-turn-state', routingState) if (tokens.accountId) { headers.set('ChatGPT-Account-Id', tokens.accountId) } @@ -99,6 +114,7 @@ export function buildOpenAICodexFetch( `[API REQUEST] ${url.pathname} remapped_to=OpenAI/Codex model=${mappedModel} source=${source ?? 'unknown'} request_id=${randomUUID()}`, ) + const wireBody = await encodeOpenAIRequestBody(JSON.stringify(upstreamBody), headers, init?.signal) const upstreamAbort = transformedBody.stream ? createTerminalAwareAbortBridge(init?.signal) : null @@ -108,7 +124,7 @@ export function buildOpenAICodexFetch( ...init, method: 'POST', headers, - body: JSON.stringify(upstreamBody), + body: wireBody, signal: upstreamAbort?.signal ?? init?.signal, }) } catch (error) { @@ -147,6 +163,10 @@ export function buildOpenAICodexFetch( ) } + // The server's first successful response pins this agentic turn. Do not + // rotate the value on continuations or persist it into another user turn. + if (!init?.signal?.aborted) turnState?.capture(upstream.headers.get('x-codex-turn-state')) + if (transformedBody.stream) { if (!upstream.body) { upstreamAbort?.dispose() @@ -260,6 +280,12 @@ function isEventStreamResponse(response: Response): boolean { .includes('text/event-stream') } +function readSessionId(input: RequestInfo | URL, init?: RequestInit): string | null { + const header = 'x-claude-code-session-id' + const fromInit = init?.headers ? new Headers(init.headers).get(header) : null + return fromInit || (input instanceof Request ? input.headers.get(header) : null) +} + async function readAnthropicBody( input: RequestInfo | URL, init?: RequestInit, diff --git a/src/services/openaiAuth/requestCompression.test.ts b/src/services/openaiAuth/requestCompression.test.ts new file mode 100644 index 00000000..1c34247a --- /dev/null +++ b/src/services/openaiAuth/requestCompression.test.ts @@ -0,0 +1,69 @@ +import { afterEach, describe, expect, spyOn, test } from 'bun:test' +import { encodeOpenAIRequestBody } from './requestCompression.js' +import { getRequestBodyAudit } from '../api/requestBodyAudit.js' + +const prior = process.env.CC_HAHA_OPENAI_REQUEST_COMPRESSION +afterEach(() => { + if (prior === undefined) delete process.env.CC_HAHA_OPENAI_REQUEST_COMPRESSION + else process.env.CC_HAHA_OPENAI_REQUEST_COMPRESSION = prior +}) + +describe('Codex request compression', () => { + test('encodes level 3 exactly and associates only the original encoded object with plaintext', async () => { + delete process.env.CC_HAHA_OPENAI_REQUEST_COMPRESSION + const compress = spyOn(Bun, 'zstdCompress') + try { + const plain = JSON.stringify({ fixture: '中文 synthetic '.repeat(1000) }) + const headers = new Headers({ 'Content-Type': 'application/json', 'Content-Length': '1' }) + const body = await encodeOpenAIRequestBody(plain, headers) + expect(compress).toHaveBeenCalledWith(Buffer.from(plain), { level: 3 }) + expect(body).toBeInstanceOf(Uint8Array) + expect(Buffer.from(await Bun.zstdDecompress(body as Uint8Array)).toString('utf8')).toBe(plain) + expect(headers.get('content-encoding')).toBe('zstd') + expect(headers.has('content-length')).toBe(false) + expect(getRequestBodyAudit(body)).toEqual({ plainBody: plain, requestEncoding: 'zstd', requestPlainBytes: Buffer.byteLength(plain), requestWireBytes: (body as Uint8Array).byteLength }) + expect(getRequestBodyAudit(new Uint8Array(body as Uint8Array))).toBeUndefined() + expect(getRequestBodyAudit(plain)).toBeUndefined() + } finally { compress.mockRestore() } + }) + + test('honors compatibility opt-out and never double encodes an existing content encoding', async () => { + process.env.CC_HAHA_OPENAI_REQUEST_COMPRESSION = 'false' + expect(await encodeOpenAIRequestBody('fixture', new Headers())).toBe('fixture') + delete process.env.CC_HAHA_OPENAI_REQUEST_COMPRESSION + const headers = new Headers({ 'Content-Encoding': 'gzip' }) + expect(await encodeOpenAIRequestBody('fixture', headers)).toBe('fixture') + expect(headers.get('content-encoding')).toBe('gzip') + }) + + test('falls back only on local encoding errors and never swallows cancellation', async () => { + delete process.env.CC_HAHA_OPENAI_REQUEST_COMPRESSION + const compress = spyOn(Bun, 'zstdCompress').mockRejectedValue(new Error('local fixture')) + try { + const headers = new Headers() + expect(await encodeOpenAIRequestBody('fixture', headers)).toBe('fixture') + expect(headers.has('content-encoding')).toBe(false) + const abort = new AbortController() + abort.abort() + await expect(encodeOpenAIRequestBody('fixture', headers, abort.signal)).rejects.toThrow() + expect(compress).toHaveBeenCalledTimes(1) + } finally { compress.mockRestore() } + }) + + test('checks cancellation after asynchronous encoding before registering or submitting', async () => { + delete process.env.CC_HAHA_OPENAI_REQUEST_COMPRESSION + let release!: (value: Uint8Array) => void + const compress = spyOn(Bun, 'zstdCompress').mockImplementation(() => new Promise(resolve => { release = resolve })) + try { + const abort = new AbortController() + const headers = new Headers() + const pending = encodeOpenAIRequestBody('fixture', headers, abort.signal) + abort.abort() + const result = new Uint8Array([1, 2, 3]) + release(result) + await expect(pending).rejects.toThrow() + expect(headers.has('content-encoding')).toBe(false) + expect(getRequestBodyAudit(result)).toBeUndefined() + } finally { compress.mockRestore() } + }) +}) diff --git a/src/services/openaiAuth/requestCompression.ts b/src/services/openaiAuth/requestCompression.ts new file mode 100644 index 00000000..55dd471a --- /dev/null +++ b/src/services/openaiAuth/requestCompression.ts @@ -0,0 +1,30 @@ +import { registerEncodedRequestBody } from '../api/requestBodyAudit.js' + +/** Called only after conversion to the fixed Codex OAuth endpoint. No network retries. */ +export async function encodeOpenAIRequestBody( + plainBody: string, + headers: Headers, + signal?: AbortSignal | null, +): Promise { + signal?.throwIfAborted() + // An explicit compatibility opt-out leaves the established plain JSON path. + if (/^(0|false|off)$/i.test(process.env.CC_HAHA_OPENAI_REQUEST_COMPRESSION ?? '') || + headers.has('Content-Encoding') || + typeof Bun === 'undefined' || typeof Bun.zstdCompress !== 'function') { + return plainBody + } + let body: Uint8Array + try { + body = await Bun.zstdCompress(Buffer.from(plainBody), { level: 3 }) + } catch { + // A local encoding failure occurs before submission; never replay a request + // after the transport has been invoked, including HTTP encoding rejection. + signal?.throwIfAborted() + return plainBody + } + signal?.throwIfAborted() + registerEncodedRequestBody(body, plainBody) + headers.set('Content-Encoding', 'zstd') + headers.delete('Content-Length') + return body +} diff --git a/src/services/openaiAuth/requestIdentity.test.ts b/src/services/openaiAuth/requestIdentity.test.ts new file mode 100644 index 00000000..4358c6db --- /dev/null +++ b/src/services/openaiAuth/requestIdentity.test.ts @@ -0,0 +1,26 @@ +import { describe, expect, test } from 'bun:test' +import { resolveOpenAIRequestIdentity } from './requestIdentity.js' + +describe('OpenAI request identity', () => { + test('uses root UUID as UUIDv5 namespace for stable branch identity', () => { + expect(resolveOpenAIRequestIdentity('01234567-89ab-4cde-8fab-0123456789ab', 'a0000000000000001')).toEqual({ + sessionId: '01234567-89ab-4cde-8fab-0123456789ab', + threadId: 'cd60d441-70fc-52c7-b989-79084eb60b3a', + }) + }) + + test('supports legacy session IDs without ambiguous concatenation or random state', () => { + const a = resolveOpenAIRequestIdentity(' legacy-root ', 'branch') + expect(a).toEqual(resolveOpenAIRequestIdentity('legacy-root', 'branch')) + expect(a?.sessionId).toBe('legacy-root') + expect(a?.threadId).not.toBe(resolveOpenAIRequestIdentity('legacy-rootbranch', '')?.threadId) + expect(a?.threadId).not.toBe(resolveOpenAIRequestIdentity('legacy-root', 'another')?.threadId) + expect(resolveOpenAIRequestIdentity('legacy-root')).toEqual({ sessionId: 'legacy-root', threadId: 'legacy-root' }) + }) + + test('missing root identity stays absent even when branch identity is supplied', () => { + for (const root of [undefined, null, '', ' ']) { + expect(resolveOpenAIRequestIdentity(root, 'branch')).toBeUndefined() + } + }) +}) diff --git a/src/services/openaiAuth/requestIdentity.ts b/src/services/openaiAuth/requestIdentity.ts new file mode 100644 index 00000000..eabab6e8 --- /dev/null +++ b/src/services/openaiAuth/requestIdentity.ts @@ -0,0 +1,31 @@ +import { createHash } from 'crypto' + +// UUID URL namespace, used only for legacy non-UUID session identifiers. +const URL_NAMESPACE = '6ba7b811-9dad-11d1-80b4-00c04fd430c8' +const UUID_PATTERN = /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i + +/** Stable transport identity. Cache overrides and per-turn routing are separate. */ +export function resolveOpenAIRequestIdentity( + rootSessionId: string | null | undefined, + agentId?: string, +): { sessionId: string; threadId: string } | undefined { + const sessionId = rootSessionId?.trim() + if (!sessionId) return undefined + if (!agentId) return { sessionId, threadId: sessionId } + + // A resumed branch keeps its agentId. Namespace by the root session so the + // same local agent ID in another conversation cannot alias this thread. + const uuidRoot = UUID_PATTERN.test(sessionId) + const namespace = uuidRoot ? sessionId : URL_NAMESPACE + const name = uuidRoot ? agentId : JSON.stringify(['cc-haha', sessionId, agentId]) + const bytes = createHash('sha1') + .update(Buffer.from(namespace.replaceAll('-', ''), 'hex')) + .update(name, 'utf8') + .digest() + .subarray(0, 16) + bytes[6] = (bytes[6]! & 0x0f) | 0x50 + bytes[8] = (bytes[8]! & 0x3f) | 0x80 + const hex = bytes.toString('hex') + const threadId = `${hex.slice(0, 8)}-${hex.slice(8, 12)}-${hex.slice(12, 16)}-${hex.slice(16, 20)}-${hex.slice(20)}` + return { sessionId, threadId } +} diff --git a/src/services/openaiAuth/subagentDefinition.integration.test.ts b/src/services/openaiAuth/subagentDefinition.integration.test.ts index 888376e4..ad6d9b8c 100644 --- a/src/services/openaiAuth/subagentDefinition.integration.test.ts +++ b/src/services/openaiAuth/subagentDefinition.integration.test.ts @@ -204,7 +204,7 @@ describe('Markdown subagent to OpenAI request integration', () => { const fetchOverride: typeof fetch = async (input, init) => { upstreamCalls.push({ url: String(input), - body: JSON.parse(String(init?.body)) as Record, + body: JSON.parse(init?.body instanceof Uint8Array ? Buffer.from(await Bun.zstdDecompress(init.body)).toString('utf8') : String(init?.body)) as Record, }) return Response.json({ id: 'resp_subagent_integration', diff --git a/src/services/openaiAuth/turnState.ts b/src/services/openaiAuth/turnState.ts new file mode 100644 index 00000000..aef7a835 --- /dev/null +++ b/src/services/openaiAuth/turnState.ts @@ -0,0 +1,28 @@ +/** One agentic user turn, never a whole conversation or shared agent session. */ +export class OpenAICodexTurnState { + #value: string | undefined + #closed = false + readonly #signal: AbortSignal + readonly #onAbort = () => this[Symbol.dispose]() + + constructor(signal: AbortSignal) { + this.#signal = signal + if (signal.aborted) this.#closed = true + else signal.addEventListener('abort', this.#onAbort, { once: true }) + } + + get(): string | undefined { + return this.#closed ? undefined : this.#value + } + + capture(value: string | null): void { + // Match Codex's OnceLock: the first response owns routing for this turn. + if (!this.#closed && this.#value === undefined && value) this.#value = value + } + + [Symbol.dispose](): void { + this.#closed = true + this.#value = undefined + this.#signal.removeEventListener('abort', this.#onAbort) + } +}