mirror of
https://github.com/NanmiCoder/claude-code-haha.git
synced 2026-10-10 11:53:10 +08:00
feat(api): align Codex request routing and tool payloads
Preserve optional function arguments with explicit non-strict schemas. Send stable session and thread identities, retain routing state through one user turn and dispose it on completion or cancellation. Compress Codex OAuth request bodies with zstd when supported, retaining readable request audit data and a local compatibility opt-out without replaying submitted requests. Validation: 352 focused offline API/provider tests passed. The real query loop regression covers tool continuation, next-turn reset and cleanup. Selected repository verification passed with mocked transports.
This commit is contained in:
@@ -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,
|
||||
|
||||
@@ -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<QueryParams['deps']>['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<typeof query>) {
|
||||
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<QueryParams['deps']>['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' })
|
||||
}
|
||||
})
|
||||
})
|
||||
@@ -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<string, unknown> }> = []
|
||||
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<string, unknown>
|
||||
}>
|
||||
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[] = []
|
||||
|
||||
@@ -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,
|
||||
})
|
||||
})
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -149,6 +149,7 @@ export type OpenAIResponsesRequest = {
|
||||
name: string
|
||||
description?: string
|
||||
parameters?: Record<string, unknown>
|
||||
strict?: boolean
|
||||
}>
|
||||
tool_choice?: unknown
|
||||
reasoning?: { effort?: OpenAIReasoningEffort }
|
||||
|
||||
@@ -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<ToolPermissionContext>;
|
||||
model: string;
|
||||
toolChoice?: BetaToolChoiceTool | BetaToolChoiceAuto | undefined;
|
||||
@@ -802,6 +804,8 @@ export async function queryModelWithoutStreaming({
|
||||
signal: AbortSignal;
|
||||
options: Options;
|
||||
}): Promise<AssistantMessage> {
|
||||
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,
|
||||
|
||||
@@ -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<Anthropic> {
|
||||
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)
|
||||
|
||||
@@ -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,
|
||||
},
|
||||
|
||||
@@ -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<Uint8Array, RequestBodyAudit>()
|
||||
|
||||
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
|
||||
}
|
||||
@@ -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<string, any> {
|
||||
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<string, unknown>,
|
||||
body: readWireBody(init) as Record<string, unknown>,
|
||||
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<string, unknown> }> = []
|
||||
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<string, unknown>
|
||||
}>
|
||||
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<Record<string, unknown>> = []
|
||||
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<string | null> = []
|
||||
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<string | null> = []
|
||||
const secondSent: Array<string | null> = []
|
||||
const upstream = (name: string, sent: Array<string | null>): 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<string | null>) => {
|
||||
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<string | null> = []
|
||||
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<string | null> = []
|
||||
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<string | null> = []
|
||||
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<string, unknown>
|
||||
}> = []
|
||||
const fetchOverride: typeof fetch = async (input, init) => {
|
||||
const body = JSON.parse(String(init?.body)) as Record<string, unknown>
|
||||
const body = readWireBody(init) as Record<string, unknown>
|
||||
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<Record<string, unknown>> = []
|
||||
const fetchOverride: typeof fetch = async (_input, init) => {
|
||||
upstreamBodies.push(JSON.parse(String(init?.body)) as Record<string, unknown>)
|
||||
upstreamBodies.push(readWireBody(init) as Record<string, unknown>)
|
||||
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<Record<string, unknown>> = []
|
||||
const fetchOverride: typeof fetch = async (_input, init) => {
|
||||
const body = JSON.parse(String(init?.body)) as Record<string, unknown>
|
||||
const body = readWireBody(init) as Record<string, unknown>
|
||||
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<Record<string, unknown>> = []
|
||||
const fetchOverride: typeof fetch = async (_input, init) => {
|
||||
const body = JSON.parse(String(init?.body)) as Record<string, unknown>
|
||||
const body = readWireBody(init) as Record<string, unknown>
|
||||
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<Record<string, unknown>> = []
|
||||
const fetchOverride: typeof fetch = async (_input, init) => {
|
||||
const body = JSON.parse(String(init?.body)) as Record<string, unknown>
|
||||
const body = readWireBody(init) as Record<string, unknown>
|
||||
// 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 =>
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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() }
|
||||
})
|
||||
})
|
||||
@@ -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<string | Uint8Array> {
|
||||
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
|
||||
}
|
||||
@@ -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()
|
||||
}
|
||||
})
|
||||
})
|
||||
@@ -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 }
|
||||
}
|
||||
@@ -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<string, unknown>,
|
||||
body: JSON.parse(init?.body instanceof Uint8Array ? Buffer.from(await Bun.zstdDecompress(init.body)).toString('utf8') : String(init?.body)) as Record<string, unknown>,
|
||||
})
|
||||
return Response.json({
|
||||
id: 'resp_subagent_integration',
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user