|
|
@@ -0,0 +1,637 @@
|
|
|
+/**
|
|
|
+ * 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 },
|
|
|
+ },
|
|
|
+ }),
|
|
|
+ ];
|
|
|
+ }
|
|
|
+}
|