| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637 |
- /**
- * 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<string, number>;
- droppedContentParts: Map<string, number>;
- droppedTools: number;
- }
- function emptyDiagnostics(): TranslateDiagnostics {
- return { droppedInputItems: new Map(), droppedContentParts: new Map(), droppedTools: 0 };
- }
- function bump(map: Map<string, number>, 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<string, unknown> | null {
- return value !== null && typeof value === 'object' && !Array.isArray(value)
- ? (value as Record<string, unknown>)
- : null;
- }
- function asString(value: unknown): string | null {
- return typeof value === 'string' ? value : null;
- }
- /** Responses 的 input_image → Chat image_url */
- function imageContentPart(part: Record<string, unknown>): 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<string, unknown>, responseId: string): Record<string, unknown> {
- 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<Record<string, unknown>> = [];
- 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, unknown>): 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<number, OpenToolCall>();
- #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, unknown>): 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 },
- },
- }),
- ];
- }
- }
|