/** * Responses ⇄ Chat Completions 翻译层(纯函数 + SSE 状态机,零 IO)。 * * 规格依据:codex rust-v0.155.1 的线协议(codex-api/src/common.rs、 * codex-api/src/sse/responses.rs、protocol/src/models.rs): * - 请求侧白名单式构建:store/include/reasoning/prompt_cache_key/text/client_metadata * 等 Chat 上游不认识的字段一律不发,宁可少发不发错。 * - 响应侧只产 codex 实际消费的事件:output_item.done 是唯一驱动历史与工具执行的通道, * function_call 的 arguments 必须整体出现在 done(0.155.1 忽略 arguments.delta), * response.completed 的 response.id 必填,缺 usage 只影响统计不影响结束。 * - tool_call_id 链路:Responses 的 call_id 原样作为 Chat 的 tool_calls[].id 与 * role:"tool" 消息的 tool_call_id,codex 侧不透明字符串原样回传。 */ export type DiagnosticFn = (message: string) => void; /** codex POST /v1/responses 的请求体(只声明翻译用得到的字段,其余按白名单丢弃) */ export interface ResponsesCreateParams { model?: unknown; instructions?: unknown; input?: unknown; tools?: unknown; tool_choice?: unknown; parallel_tool_calls?: unknown; max_output_tokens?: unknown; temperature?: unknown; top_p?: unknown; stream?: unknown; [key: string]: unknown; } export interface ChatToolCall { id: string; type: 'function'; function: { name: string; arguments: string }; } export interface ChatMessage { role: 'developer' | 'system' | 'user' | 'assistant' | 'tool'; content?: unknown; tool_calls?: ChatToolCall[]; tool_call_id?: string; name?: string; } export interface ChatCompletionRequest { model: string; messages: ChatMessage[]; stream: boolean; stream_options?: { include_usage: boolean }; tools?: Array<{ type: 'function'; function: { name: string; description?: string; parameters?: unknown } }>; tool_choice?: 'auto' | 'none' | 'required' | { type: 'function'; function: { name: string } }; parallel_tool_calls?: boolean; max_tokens?: number; temperature?: number; top_p?: number; } interface ChatToolCallDelta { index?: number; id?: string; function?: { name?: string; arguments?: string }; } interface ChatChunkChoice { delta?: { content?: string | null; reasoning_content?: string | null; tool_calls?: ChatToolCallDelta[] }; finish_reason?: string | null; } interface ChatChunk { choices?: ChatChunkChoice[]; usage?: { prompt_tokens?: number; completion_tokens?: number; total_tokens?: number; } | null; } /** 翻译过程中丢弃的非空内容计数,进诊断,绝不静默有损 */ export interface TranslateDiagnostics { droppedInputItems: Map; droppedContentParts: Map; droppedTools: number; } function emptyDiagnostics(): TranslateDiagnostics { return { droppedInputItems: new Map(), droppedContentParts: new Map(), droppedTools: 0 }; } function bump(map: Map, key: string): void { map.set(key, (map.get(key) ?? 0) + 1); } function diagnosticsMessage(diagnostics: TranslateDiagnostics): string { const parts: string[] = []; for (const [key, count] of diagnostics.droppedInputItems) parts.push(`input.${key}×${count}`); for (const [key, count] of diagnostics.droppedContentParts) parts.push(`content.${key}×${count}`); if (diagnostics.droppedTools > 0) parts.push(`tools×${diagnostics.droppedTools}`); return parts.join(', '); } function asRecord(value: unknown): Record | null { return value !== null && typeof value === 'object' && !Array.isArray(value) ? (value as Record) : null; } function asString(value: unknown): string | null { return typeof value === 'string' ? value : null; } /** Responses 的 input_image → Chat image_url */ function imageContentPart(part: Record): ChatContentPart | null { const url = asString(part.image_url) ?? asString(asRecord(part.image_url)?.url); return url ? { type: 'image_url', image_url: { url } } : null; } type ChatContentPart = { type: 'text'; text: string } | { type: 'image_url'; image_url: { url: string } }; /** Responses message 的 content 数组 → Chat content(文本/图片;单文本压成字符串兼容性最好) */ function translateMessageContent(parts: unknown[], diagnostics: TranslateDiagnostics): string | ChatContentPart[] { const translated: ChatContentPart[] = []; for (const part of parts) { const record = asRecord(part); const type = asString(record?.type); if (type === 'input_text' || type === 'output_text') { const text = asString(record?.text) ?? ''; translated.push({ type: 'text', text }); } else if (type === 'input_image') { const image = imageContentPart(record ?? {}); if (image) translated.push(image); else bump(diagnostics.droppedContentParts, 'input_image'); } else if (type !== null) { bump(diagnostics.droppedContentParts, type ?? 'unknown'); } } const textsOnly = translated.every((part) => part.type === 'text'); if (textsOnly && translated.length <= 1) return (translated[0] as { text: string } | undefined)?.text ?? ''; return translated; } /** function_call_output.output:字符串或结构化数组,折叠成 Chat tool 消息的文本 content */ function translateOutputValue(value: unknown, diagnostics: TranslateDiagnostics): string { if (typeof value === 'string') return value; if (Array.isArray(value)) { const texts: string[] = []; for (const part of value) { const record = asRecord(part); const type = asString(record?.type); if (type === 'input_text' || type === 'output_text' || type === 'text') { texts.push(asString(record?.text) ?? ''); } else { bump(diagnostics.droppedContentParts, type ?? 'unknown'); } } return texts.join('\n'); } if (value === null || value === undefined) return ''; return String(value); } /** * Responses 请求体 → Chat Completions 请求体。 * stream 透传;stream:true 时补 stream_options.include_usage(不支持的端点会忽略)。 */ export function responsesToChatRequest( req: ResponsesCreateParams, onDiagnostic?: DiagnosticFn, ): ChatCompletionRequest { const diagnostics = emptyDiagnostics(); const messages: ChatMessage[] = []; const instructions = asString(req.instructions); if (instructions) messages.push({ role: 'developer', content: instructions }); const input = req.input; if (typeof input === 'string') { messages.push({ role: 'user', content: input }); } else if (Array.isArray(input)) { for (const raw of input) { const item = asRecord(raw); if (!item) continue; const type = asString(item.type); if (type === null || type === 'message') { const role = asString(item.role) ?? 'user'; const content = Array.isArray(item.content) ? translateMessageContent(item.content, diagnostics) : asString(item.content) ?? ''; messages.push({ role: role === 'developer' ? 'developer' : role as ChatMessage['role'], content }); continue; } if (type === 'function_call' || type === 'custom_tool_call') { const callId = asString(item.call_id) ?? ''; const name = asString(item.name) ?? ''; const args = type === 'custom_tool_call' ? asString(item.input) ?? JSON.stringify(item.input ?? null) : typeof item.arguments === 'string' ? item.arguments : JSON.stringify(item.arguments ?? null); messages.push({ role: 'assistant', content: null, tool_calls: [{ id: callId, type: 'function', function: { name, arguments: args } }], }); continue; } if (type === 'function_call_output' || type === 'custom_tool_call_output') { messages.push({ role: 'tool', tool_call_id: asString(item.call_id) ?? '', content: translateOutputValue(item.output, diagnostics), }); continue; } // reasoning / local_shell_call / web_search_call / item_reference 等对 Chat 上游无意义 bump(diagnostics.droppedInputItems, type ?? 'unknown'); } } const chat: ChatCompletionRequest = { model: asString(req.model) ?? '', messages, stream: req.stream === true, }; if (chat.stream) chat.stream_options = { include_usage: true }; const tools = Array.isArray(req.tools) ? req.tools : []; const chatTools: ChatCompletionRequest['tools'] = []; for (const raw of tools) { const tool = asRecord(raw); if (asString(tool?.type) !== 'function') { diagnostics.droppedTools += 1; continue; } const description = asString(tool?.description); chatTools.push({ type: 'function', function: { name: asString(tool?.name) ?? '', ...(description ? { description } : {}), ...(tool?.parameters !== undefined ? { parameters: tool?.parameters } : {}), }, }); } if (chatTools.length) chat.tools = chatTools; const toolChoice = req.tool_choice; if (toolChoice === 'auto' || toolChoice === 'none' || toolChoice === 'required') { chat.tool_choice = toolChoice; } else if (asString(asRecord(toolChoice)?.type) === 'function' && asRecord(toolChoice)?.name) { chat.tool_choice = { type: 'function', function: { name: asString(asRecord(toolChoice)?.name) ?? '' } }; } else if (toolChoice !== undefined && toolChoice !== null) { bump(diagnostics.droppedInputItems, `tool_choice:${typeof toolChoice}`); } if (typeof req.parallel_tool_calls === 'boolean') chat.parallel_tool_calls = req.parallel_tool_calls; if (typeof req.max_output_tokens === 'number') chat.max_tokens = req.max_output_tokens; if (typeof req.temperature === 'number') chat.temperature = req.temperature; if (typeof req.top_p === 'number') chat.top_p = req.top_p; const summary = diagnosticsMessage(diagnostics); if (summary && onDiagnostic) onDiagnostic(`有损转换丢弃:${summary}`); return chat; } /** Chat 非流式响应 → Responses 形状(防御路径:codex 恒为 stream:true) */ export function chatResponseToResponses(chat: Record, responseId: string): Record { const diagnostics = emptyDiagnostics(); const choices = Array.isArray(chat.choices) ? chat.choices : []; const choice = asRecord(choices[0]) ?? {}; const message = asRecord(choice.message) ?? {}; const finishReason = asString(choice.finish_reason); const output: Array> = []; const reasoning = asString(message.reasoning_content) ?? asString(message.reasoning); if (reasoning) { output.push({ type: 'reasoning', id: 'rs_1', summary: [{ type: 'summary_text', text: reasoning }], }); } const content = asString(message.content); if (content) { output.push({ type: 'message', id: 'msg_1', role: 'assistant', status: 'completed', content: [{ type: 'output_text', text: content }], }); } const toolCalls = Array.isArray(message.tool_calls) ? message.tool_calls : []; let callIndex = 0; for (const raw of toolCalls) { const call = asRecord(raw); const fn = asRecord(call?.function); callIndex += 1; output.push({ type: 'function_call', id: `fc_${callIndex}`, call_id: asString(call?.id) ?? `call_${callIndex}`, name: asString(fn?.name) ?? '', arguments: asString(fn?.arguments) ?? '', status: 'completed', }); } if (!reasoning && !content && toolCalls.length === 0) { bump(diagnostics.droppedInputItems, 'empty_choice'); } const usage = asRecord(chat.usage); const promptTokens = typeof usage?.prompt_tokens === 'number' ? usage.prompt_tokens : 0; const completionTokens = typeof usage?.completion_tokens === 'number' ? usage.completion_tokens : 0; const totalTokens = typeof usage?.total_tokens === 'number' ? usage.total_tokens : promptTokens + completionTokens; const incomplete = finishReason === 'length'; return { id: responseId, object: 'response', status: incomplete ? 'incomplete' : 'completed', ...(incomplete ? { incomplete_details: { reason: 'max_output_tokens' } } : {}), output, usage: { input_tokens: promptTokens, input_tokens_details: { cached_tokens: 0 }, output_tokens: completionTokens, output_tokens_details: { reasoning_tokens: 0 }, total_tokens: totalTokens, }, parallel_tool_calls: typeof chat.parallel_tool_calls === 'boolean' ? chat.parallel_tool_calls : undefined, }; } export type BridgeSseFrame = string; interface OpenToolCall { callId: string | null; name: string; arguments: string; itemId: string; } function randomId(): string { return Math.random().toString(36).slice(2, 10) + Date.now().toString(36).slice(-6); } function sseFrame(type: string, data: Record): string { return `event: ${type}\ndata: ${JSON.stringify({ type, ...data })}\n\n`; } /** * 上游 Chat SSE → codex Responses SSE 的流式状态机。 * push() 吃任意分包的文本(内部按行缓冲),finish()/fail() 收尾且幂等。 * 一个响应内文本/推理/工具调用的 item 顺序:按上游到达顺序开闭,工具调用按 index 归拢后 * 在 finish(或出现更大 index)时补发完整 done —— arguments 必须整体出现。 */ export class ChatSseTranslator { readonly #responseId: string; readonly #onDiagnostic?: DiagnosticFn; #buffer = ''; #itemSeq = 0; #closed = false; #textState: { itemId: string; text: string } | null = null; #reasoningState: { itemId: string; text: string } | null = null; /** index → 归拢中的工具调用;出现更大 index 时先收口前面的 */ #toolCalls = new Map(); #toolCallOrder: number[] = []; #toolCallsDone = false; #finishReason: string | null = null; #usage: { input_tokens: number; output_tokens: number; total_tokens: number } | null = null; #upstreamBytes = 0; #emittedEvents = 0; constructor(onDiagnostic?: DiagnosticFn) { this.#onDiagnostic = onDiagnostic; this.#responseId = `resp_${randomId()}`; } get responseId(): string { return this.#responseId; } get upstreamBytes(): number { return this.#upstreamBytes; } get emittedEvents(): number { return this.#emittedEvents; } #frame(type: string, data: Record): string { this.#emittedEvents += 1; return sseFrame(type, data); } #nextItemId(prefix: string): string { this.#itemSeq += 1; return `${prefix}_${this.#itemSeq}`; } /** 先建 SSE:不等上游响应头就发 created + in_progress,防单槽端点慢 prefill 掐线 */ begin(): string[] { if (this.#closed) return []; return [ this.#frame('response.created', { response: { id: this.#responseId } }), this.#frame('response.in_progress', { response: { id: this.#responseId } }), ]; } /** 关闭当前打开的文本/推理 item,产出 done 帧 */ #closeText(): string[] { if (!this.#textState) return []; const { itemId, text } = this.#textState; this.#textState = null; return [ this.#frame('response.output_item.done', { item: { type: 'message', id: itemId, role: 'assistant', status: 'completed', content: [{ type: 'output_text', text }], }, }), ]; } #closeReasoning(): string[] { if (!this.#reasoningState) return []; const { itemId, text } = this.#reasoningState; this.#reasoningState = null; return [ this.#frame('response.output_item.done', { item: { type: 'reasoning', id: itemId, summary: [{ type: 'summary_text', text }], }, }), ]; } /** 收口所有归拢中的工具调用(按 index 序),产出完整 function_call done 帧 */ #closeToolCalls(): string[] { if (this.#toolCallsDone) return []; this.#toolCallsDone = true; const frames: string[] = []; for (const index of this.#toolCallOrder) { const call = this.#toolCalls.get(index); if (!call) continue; frames.push( this.#frame('response.output_item.done', { item: { type: 'function_call', id: call.itemId, call_id: call.callId ?? `call_${index}`, name: call.name, arguments: call.arguments, status: 'completed', }, }), ); } return frames; } push(chunk: string): string[] { if (this.#closed) return []; this.#upstreamBytes += Buffer.byteLength(chunk); this.#buffer += chunk; const frames: string[] = []; let newlineIndex = this.#buffer.indexOf('\n'); while (newlineIndex >= 0) { const line = this.#buffer.slice(0, newlineIndex).replace(/\r$/, ''); this.#buffer = this.#buffer.slice(newlineIndex + 1); const payload = this.#dataPayload(line); if (payload) frames.push(...this.#handleData(payload)); newlineIndex = this.#buffer.indexOf('\n'); } return frames; } /** 提取 SSE data: 行的载荷;[DONE]/注释行/空行返回 null */ #dataPayload(line: string): string | null { if (!line.startsWith('data:')) return null; const payload = line.slice(5).trim(); if (!payload || payload === '[DONE]') return null; return payload; } #handleData(payload: string): string[] { let chunk: ChatChunk; try { chunk = JSON.parse(payload) as ChatChunk; } catch { // 坏 JSON 跳过不杀流 this.#onDiagnostic?.('上游 SSE 出现无法解析的 JSON 分片,已跳过'); return []; } const frames: string[] = []; if (chunk.usage && typeof chunk.usage === 'object') { const promptTokens = typeof chunk.usage.prompt_tokens === 'number' ? chunk.usage.prompt_tokens : 0; const completionTokens = typeof chunk.usage.completion_tokens === 'number' ? chunk.usage.completion_tokens : 0; const totalTokens = typeof chunk.usage.total_tokens === 'number' ? chunk.usage.total_tokens : promptTokens + completionTokens; this.#usage = { input_tokens: promptTokens, output_tokens: completionTokens, total_tokens: totalTokens }; } const choice = chunk.choices?.[0]; if (!choice) return frames; if (choice.finish_reason) this.#finishReason = choice.finish_reason; const delta = choice.delta; if (delta?.tool_calls?.length) { // 工具调用开跑:先收口打开的文本/推理 item frames.push(...this.#closeText(), ...this.#closeReasoning()); for (const callDelta of delta.tool_calls) { const index = typeof callDelta.index === 'number' ? callDelta.index : 0; let call = this.#toolCalls.get(index); if (!call) { // 出现新 index:说明更小的 index 已收口完成 frames.push(...this.#closeEarlierToolCalls(index)); call = { callId: null, name: '', arguments: '', itemId: this.#nextItemId('fc') }; this.#toolCalls.set(index, call); this.#toolCallOrder.push(index); } if (callDelta.id) call.callId = callDelta.id; if (callDelta.function?.name) call.name += callDelta.function.name; if (callDelta.function?.arguments) call.arguments += callDelta.function.arguments; } } if (typeof delta?.reasoning_content === 'string' && delta.reasoning_content.length > 0) { if (this.#textState) frames.push(...this.#closeText()); if (!this.#reasoningState) { this.#reasoningState = { itemId: this.#nextItemId('rs'), text: '' }; frames.push( this.#frame('response.output_item.added', { item: { type: 'reasoning', id: this.#reasoningState.itemId, summary: [] }, }), ); } this.#reasoningState.text += delta.reasoning_content; frames.push( this.#frame('response.reasoning_summary_text.delta', { delta: delta.reasoning_content, summary_index: 0, item_id: this.#reasoningState.itemId, }), ); } if (typeof delta?.content === 'string' && delta.content.length > 0) { if (this.#reasoningState) frames.push(...this.#closeReasoning()); if (!this.#textState) { this.#textState = { itemId: this.#nextItemId('msg'), text: '' }; frames.push( this.#frame('response.output_item.added', { item: { type: 'message', id: this.#textState.itemId, role: 'assistant', content: [] }, }), ); } this.#textState.text += delta.content; frames.push(this.#frame('response.output_text.delta', { delta: delta.content, item_id: this.#textState.itemId })); } return frames; } /** 出现 index N 时收口所有 < N 的工具调用 */ #closeEarlierToolCalls(index: number): string[] { const frames: string[] = []; for (const earlier of this.#toolCallOrder) { if (earlier >= index) break; const call = this.#toolCalls.get(earlier); if (!call) continue; frames.push( this.#frame('response.output_item.done', { item: { type: 'function_call', id: call.itemId, call_id: call.callId ?? `call_${earlier}`, name: call.name, arguments: call.arguments, status: 'completed', }, }), ); this.#toolCalls.delete(earlier); } return frames; } /** [DONE]/上游结束:收口全部 item 并补 response.completed。幂等 */ finish(): string[] { if (this.#closed) return []; this.#closed = true; const frames: string[] = [ ...this.#closeReasoning(), ...this.#closeText(), ...this.#closeToolCalls(), ]; const incomplete = this.#finishReason === 'length'; frames.push( this.#frame('response.completed', { response: { id: this.#responseId, status: incomplete ? 'incomplete' : 'completed', ...(incomplete ? { incomplete_details: { reason: 'max_output_tokens' } } : {}), ...(this.#usage ? { usage: { input_tokens: this.#usage.input_tokens, input_tokens_details: { cached_tokens: 0 }, output_tokens: this.#usage.output_tokens, output_tokens_details: { reasoning_tokens: 0 }, total_tokens: this.#usage.total_tokens, } } : {}), }, }), ); return frames; } /** 上游失败:补 response.failed 让 codex 走错误映射。幂等 */ fail(message: string): string[] { if (this.#closed) return []; this.#closed = true; return [ this.#frame('response.failed', { response: { id: this.#responseId, error: { code: 'upstream_error', message }, }, }), ]; } }