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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 8 additions & 6 deletions src/services/api/openaiShim.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3764,15 +3764,17 @@ test('the OpenAI shim façade creates independent client instances', () => {
})
// openaiShim test extraction seam 112 end

test('raw-text and XML fallback tool calls use one unique sequence', () => {
test('facade parseTextToolCalls and parseXmlToolCalls share adapter sequencing', () => {
const text = parseTextToolCalls('{"name":"from_text","arguments":{}}')
const xml = parseXmlToolCalls('<tool_call>{"name":"from_xml","arguments":{}}</tool_call>')
const xml = parseXmlToolCalls(
'<tool_call>{"name":"from_xml","arguments":{}}</tool_call>',
)

expect(text.calls[0]?.id).toMatch(/^ollama_tc_\d+$/)
expect(xml.calls[0]?.id).toMatch(/^xml_tc_\d+$/)
const textNum = Number(text.calls[0]?.id?.replace(/^\D+/, ''))
const xmlNum = Number(xml.calls[0]?.id?.replace(/^\D+/, ''))
// Same session counter: the second mint must be exactly one greater than the first.
expect(xmlNum).toBe(textNum + 1)
const textSequence = Number(text.calls[0]?.id?.replace(/^\D+/, ''))
const xmlSequence = Number(xml.calls[0]?.id?.replace(/^\D+/, ''))
expect(xmlSequence).toBe(textSequence + 1)
})

// ---------------------------------------------------------------------------
Expand Down
258 changes: 17 additions & 241 deletions src/services/api/openaiShim.ts
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,7 @@ import {
refreshCodexAccessTokenIfNeeded,
} from '../../utils/codexCredentials.js'
import { logForDebugging } from '../../utils/debug.js'
import { anthropicSsePassthrough as parseAnthropicSsePassthrough, createReaderCanceller, createStreamAbortError, getStreamIdleTimeoutMs, readWithIdleTimeout, StreamIdleTimeoutError, throwIfStreamAborted } from './openaiShim/streamControl.js'
import { createStreamAbortError, getStreamIdleTimeoutMs, readWithIdleTimeout, StreamIdleTimeoutError } from './openaiShim/streamControl.js'
export { getStreamIdleTimeoutMs } from './openaiShim/streamControl.js'
import { isBareMode, isEnvTruthy } from '../../utils/envUtils.js'
import {
Expand All @@ -68,10 +68,6 @@ import {
resolveRouteCredentialValue,
} from '../../integrations/routeMetadata.js'
import { getSessionId } from '../../bootstrap/state.js'
import {
createThinkTagFilter,
stripThinkTags,
} from './thinkTagSanitizer.js'
import {
codexStreamToAnthropic,
collectCodexCompletedResponse,
Expand All @@ -80,19 +76,21 @@ import {
convertToolsToResponsesTools,
performCodexRequest,
type AnthropicStreamEvent,
type AnthropicUsage,
type ShimCreateParams,
} from './codexShim.js'
import {
createRequestBodyPlanner,
hydrateOpenAIShimCompatibilityEnv as hydrateRequestPlanningEnv,
} from './openaiShim/requestPlanner.js'
import { buildAnthropicUsageFromRawUsage } from './cacheMetrics.js'
import {
convertOpenAIStreamUsage,
openaiStreamToAnthropic as convertOpenAIStream,
} from './openaiShim/streamConversion.js'
import { geminiSseToAnthropic as convertGeminiStream } from './openaiShim/geminiStreamConversion.js'
anthropicSsePassthrough,
convertGeminiToAnthropicResponse,
convertNonStreamingResponseToAnthropicMessage,
geminiSseToAnthropic,
makeMessageId,
openaiStreamToAnthropic as convertOpenAIResponseStream,
} from './openaiShim/responseAdapters.js'
export { parseTextToolCalls, parseXmlToolCalls } from './openaiShim/responseAdapters.js'
import { compressToolHistory } from './compressToolHistory.js'
import {
createClassifiedTransportError,
Expand Down Expand Up @@ -127,10 +125,6 @@ import {
markOpenAIRequestNonReplayable,
} from './openaiErrorClassification.js'
import { redactSecretValueForDisplay, type SecretValueSource } from '../../utils/providerProfile.js'
import {
normalizeToolArguments,
hasToolFieldMapping,
} from './toolArgumentNormalization.js'
import { logApiCallStart, logApiCallEnd } from '../../utils/requestLogging.js'
import {
createStreamState,
Expand All @@ -139,13 +133,6 @@ import {
} from '../../utils/streamingOptimizer.js'
import { stableStringifyJson } from '../../utils/stableStringify.js'
import {
findXmlToolCallOpener as findXmlToolCallOpenerModule,
isHy3Model as isHy3ModelModule,
parseXmlToolCalls as parseXmlToolCallsModule,
trailingXmlOpenerPrefixLen as trailingXmlOpenerPrefixLenModule,
} from './openaiShim/xmlToolCallParsing.js'
import {
convertNonStreamingResponseToAnthropicMessage as convertResponseToAnthropicMessage,
type NonStreamingOpenAIResponse,
} from './openaiShim/responseConversion.js'
import {
Expand Down Expand Up @@ -179,17 +166,6 @@ import {
convertMessages as convertAnthropicMessages,
convertSystemPrompt as convertSystemPromptImpl,
} from './openaiShim/messageConversion.js'
import {
JSON_REPAIR_SUFFIXES,
couldBeRawToolCallsRequestedPrefix,
extractBalancedJson,
parseRawToolCallsRequestedText,
parseTextToolCalls as parseTextToolCallsModule,
repairPossiblyTruncatedObjectJson,
stripRanges,
type ParsedRawToolCall,
type ParsedTextToolCall,
} from './openaiShim/rawToolCallParsing.js'
import {
convertTools as convertToolsModule,
normalizeSchemaForOpenAI as normalizeSchemaForOpenAIModule,
Expand Down Expand Up @@ -414,142 +390,6 @@ function convertTools(
// Streaming: OpenAI SSE → Anthropic stream events
// ---------------------------------------------------------------------------

interface OpenAIStreamChunk {
id: string
object: string
model: string
choices: Array<{
index: number
delta: {
role?: string
content?: string | null
reasoning_content?: string | null
extra_content?: Record<string, unknown>
tool_calls?: Array<{
index: number
id?: string
type?: string
function?: { name?: string; arguments?: string }
extra_content?: Record<string, unknown>
}>
}
finish_reason: string | null
}>
usage?: {
prompt_tokens?: number
completion_tokens?: number
total_tokens?: number
prompt_tokens_details?: {
cached_tokens?: number
}
}
}

function makeMessageId(): string {
return `msg_${crypto.randomUUID().replace(/-/g, '')}`
}

function convertChunkUsage(usage: OpenAIStreamChunk['usage'] | undefined): Partial<AnthropicUsage> | undefined {
return convertOpenAIStreamUsage(usage as Record<string, unknown> | undefined)
}

export function parseTextToolCalls(text: string): {
calls: ParsedTextToolCall[]
toolCallRanges: Array<[number, number]>
} {
return parseTextToolCallsModule(text, nextTextToolCallSequence)
}

// Shared façade state keeps raw-text and XML fallback IDs unique per session.
let textToolCallSequence = 0

function nextTextToolCallSequence(): number {
return ++textToolCallSequence
}

// ---------------------------------------------------------------------------
// XML tool parsing façade. Dialect handling lives in xmlToolCallParsing.ts.
// ---------------------------------------------------------------------------

function findXmlToolCallOpener(text: string, allowHy3: boolean): number {
return findXmlToolCallOpenerModule(text, allowHy3)
}

function isHy3Model(model: string): boolean {
return isHy3ModelModule(model)
}

export function parseXmlToolCalls(text: string, allowHy3 = false) {
return parseXmlToolCallsModule(text, allowHy3, nextTextToolCallSequence)
}

function trailingXmlOpenerPrefixLen(text: string, allowHy3: boolean): number {
return trailingXmlOpenerPrefixLenModule(text, allowHy3)
}

// The streaming finalize path buffers from this opener onward so the raw XML
// is never surfaced as text before extraction.
/**
* Async generator that transforms an OpenAI SSE stream into
* Anthropic-format BetaRawMessageStreamEvent objects.
*/
/**
* Passthrough for Anthropic Messages API SSE streams.
* The response events are already in AnthropicStreamEvent format —
* we just parse the SSE frames and yield them directly.
*/
async function* anthropicSsePassthrough(
response: Response,
_model: string,
signal?: AbortSignal,
): AsyncGenerator<AnthropicStreamEvent> {
yield* parseAnthropicSsePassthrough<AnthropicStreamEvent>(
response,
signal,
(message, options) => options?.level
? logForDebugging(message, { level: options.level })
: logForDebugging(message),
)
}

/**
* Transforms Google AI SDK SSE stream into Anthropic-format stream events.
* Google AI SDK yields frames with { candidates: [{ content: { role, parts } }] }.
*/
async function* geminiSseToAnthropic(
response: Response,
model: string,
signal?: AbortSignal,
): AsyncGenerator<AnthropicStreamEvent> {
yield* convertGeminiStream(response, model, signal, {
createReaderCanceller,
createStreamAbortError,
getStreamIdleTimeoutMs,
makeMessageId,
readWithIdleTimeout,
throwIfStreamAborted,
})
}
// Extraction seam: Gemini streaming | completed response conversion.

function convertNonStreamingResponseToAnthropicMessage(
data: NonStreamingOpenAIResponse,
model: string,
) {
return convertResponseToAnthropicMessage(data, model, {
makeMessageId,
buildUsage: usage => buildAnthropicUsageFromRawUsage(usage),
stripThinkTags,
parseXmlToolCalls,
isHy3Model,
stripRanges,
parseRawToolCalls: parseRawToolCallsRequestedText,
normalizeToolArguments,
getGeminiThoughtSignature: geminiThoughtSignatureFromExtraContent,
mergeGeminiThoughtSignature,
})
}

import { headersWithRequestUrl as buildHeadersWithRequestUrl } from './openaiShim/clientDispatch.js'

function headersWithRequestUrl(headers: Headers, requestUrl?: string): Headers {
Expand All @@ -565,31 +405,14 @@ async function* openaiStreamToAnthropic(
isOllama = false,
requestUrl?: string,
): AsyncGenerator<AnthropicStreamEvent> {
yield* convertOpenAIStream(response, model, signal, isOllama, requestUrl, {
convertNonStreamingResponseToAnthropicMessage: (data, streamModel) =>
convertNonStreamingResponseToAnthropicMessage(
data as NonStreamingOpenAIResponse,
streamModel,
),
couldBeRawToolCallsRequestedPrefix,
createReaderCanceller,
createStreamAbortError,
findXmlToolCallOpener,
geminiThoughtSignatureFromExtraContent,
getStreamIdleTimeoutMs,
yield* convertOpenAIResponseStream(
response,
model,
signal,
isOllama,
requestUrl,
headersWithRequestUrl,
isHy3Model,
makeMessageId,
mergeGeminiThoughtSignature,
parseRawToolCallsRequestedText,
parseTextToolCalls,
parseXmlToolCalls,
readWithIdleTimeout,
repairPossiblyTruncatedObjectJson,
stripRanges,
throwIfStreamAborted,
trailingXmlOpenerPrefixLen,
})
)
}


Expand Down Expand Up @@ -1113,54 +936,7 @@ class OpenAIShimMessages {
data: Record<string, unknown>,
model: string,
) {
const content: Array<Record<string, unknown>> = []
let hasToolUse = false
const candidates = data.candidates as Array<Record<string, unknown>> | undefined
const candidate = candidates?.[0]
const candidateContent = candidate?.content as { parts?: Array<Record<string, unknown>> } | undefined

if (candidateContent?.parts) {
for (const part of candidateContent.parts) {
const text = part.text as string | undefined
if (text) {
content.push({ type: 'text', text })
}
const fc = part.functionCall as { name?: string; args?: unknown } | undefined
if (fc?.name) {
hasToolUse = true
content.push({
type: 'tool_use',
id: `toolu_${crypto.randomUUID().replace(/-/g, '').slice(0, 24)}`,
name: fc.name,
input: fc.args ?? {},
})
}
}
}

const stopReason =
hasToolUse
? 'tool_use'
: candidate?.finishReason === 'MAX_TOKENS'
? 'max_tokens'
: 'end_turn'

const usageMetadata = data.usageMetadata as Record<string, number> | undefined
const usage = buildAnthropicUsageFromRawUsage({
input_tokens: usageMetadata?.promptTokenCount ?? 0,
output_tokens: (usageMetadata?.candidatesTokenCount ?? 0) + (usageMetadata?.thoughtsTokenCount ?? 0),
} as unknown as Record<string, unknown>)

return {
id: makeMessageId(),
type: 'message',
role: 'assistant',
content,
model,
stop_reason: stopReason,
stop_sequence: null,
usage,
}
return convertGeminiToAnthropicResponse(data, model)
}
}

Expand Down
Loading
Loading