/** * SSE 流式解析工具 * * 解析 OpenAI 兼容的 Server-Sent Events (SSE) 流式响应, * 产出 MetonaStreamEvent。DeepSeek 和 Agnes AI 共享此工具。 * * SSE 格式:data: {json}\n\n * 结束标记:data: [DONE] */ import { nanoid } from 'nanoid'; import type { MetonaStreamEvent, MetonaTokenUsage } from '../../types'; import { MetonaStreamEventType } from '../../types'; /** * 解析 OpenAI 兼容 SSE 流式响应 * * @param responseBody - fetch Response.body (ReadableStream) * @param requestId - 对应的请求 ID * @param sessionId - 会话 ID * @param iteration - 当前迭代轮次 * @yields MetonaStreamEvent */ export async function* parseSSEStream( responseBody: ReadableStream, requestId: string, sessionId: string, iteration: number, ): AsyncGenerator { const reader = responseBody.getReader(); const decoder = new TextDecoder(); let seq = 0; let buffer = ''; // 工具调用缓冲区:index → { name, argsBuffer } const toolCallsBuffer = new Map(); while (true) { const { done, value } = await reader.read(); if (done) break; buffer += decoder.decode(value, { stream: true }); const lines = buffer.split('\n'); buffer = lines.pop() ?? ''; for (const line of lines) { const trimmed = line.trim(); if (!trimmed || !trimmed.startsWith('data: ')) continue; const data = trimmed.slice(6); // 流结束 if (data === '[DONE]') { // 将缓冲区中未完成拼接的工具调用发送 for (const [index, buf] of toolCallsBuffer) { try { yield { type: MetonaStreamEventType.TOOL_CALL_COMPLETE, requestId, sessionId, iteration, seq: seq++, timestamp: Date.now(), toolCall: { id: `tc_${nanoid(8)}`, name: buf.name, args: buf.argsBuffer ? JSON.parse(buf.argsBuffer) : {}, iteration, timestamp: Date.now(), }, }; } catch { // JSON 解析失败,跳过 } } toolCallsBuffer.clear(); yield { type: MetonaStreamEventType.DONE, requestId, sessionId, iteration, seq: seq++, timestamp: Date.now(), }; return; } try { const chunk = JSON.parse(data); const delta = chunk.choices?.[0]?.delta; // 文本内容增量 if (delta?.content) { yield { type: MetonaStreamEventType.TEXT_DELTA, requestId, sessionId, iteration, seq: seq++, timestamp: Date.now(), delta: delta.content, }; } // 推理内容增量(Thinking 模式) if (delta?.reasoning_content) { yield { type: MetonaStreamEventType.REASONING_DELTA, requestId, sessionId, iteration, seq: seq++, timestamp: Date.now(), delta: delta.reasoning_content, }; } // 工具调用增量 — 缓冲拼接 if (delta?.tool_calls) { for (const tc of delta.tool_calls) { const idx = tc.index ?? 0; if (!toolCallsBuffer.has(idx)) { toolCallsBuffer.set(idx, { name: tc.function?.name ?? '', argsBuffer: '' }); } const buf = toolCallsBuffer.get(idx)!; if (tc.function?.name) buf.name = tc.function.name; if (tc.function?.arguments) buf.argsBuffer += tc.function.arguments; yield { type: MetonaStreamEventType.TOOL_CALL_DELTA, requestId, sessionId, iteration, seq: seq++, timestamp: Date.now(), toolCallDelta: { index: idx, name: tc.function?.name, argsDelta: tc.function?.arguments, }, }; } } // Token 使用统计 if (chunk.usage) { const usage: MetonaTokenUsage = { inputTokens: chunk.usage.prompt_tokens ?? 0, outputTokens: chunk.usage.completion_tokens ?? 0, totalTokens: chunk.usage.total_tokens ?? 0, reasoningTokens: chunk.usage.completion_tokens_details?.reasoning_tokens, cacheHitTokens: chunk.usage.prompt_cache_hit_tokens, cacheMissTokens: chunk.usage.prompt_cache_miss_tokens, }; yield { type: MetonaStreamEventType.USAGE, requestId, sessionId, iteration, seq: seq++, timestamp: Date.now(), usage, }; } } catch { // 跳过解析失败的行 } } } } /** * 解析 OpenAI 兼容的非流式 JSON 响应 → MetonaResponse */ export function parseOpenAICompatibleResponse( data: Record, requestId: string, provider: string, defaultModel: string, ): { content: string; reasoningContent?: string; toolCalls?: Array<{ id: string; name: string; args: Record; iteration: number; timestamp: number }>; finishReason: string; usage: MetonaTokenUsage; } { const choice = (data.choices as Array>)?.[0]; const message = choice?.message as Record | undefined; const usage = data.usage as Record | undefined; const rawToolCalls = message?.tool_calls as Array> | undefined; return { content: (message?.content as string) ?? '', reasoningContent: message?.reasoning_content as string | undefined, toolCalls: rawToolCalls?.map((tc) => { const fn = tc.function as Record; return { id: tc.id as string, name: fn.name as string, args: JSON.parse(fn.arguments as string), iteration: 0, timestamp: Date.now(), }; }), finishReason: mapOpenAIFinishReason(choice?.finish_reason as string), usage: { inputTokens: (usage?.prompt_tokens as number) ?? 0, outputTokens: (usage?.completion_tokens as number) ?? 0, totalTokens: (usage?.total_tokens as number) ?? 0, reasoningTokens: (usage?.completion_tokens_details as Record)?.reasoning_tokens as number | undefined, cacheHitTokens: usage?.prompt_cache_hit_tokens as number | undefined, cacheMissTokens: usage?.prompt_cache_miss_tokens as number | undefined, }, }; } function mapOpenAIFinishReason(reason: string): string { switch (reason) { case 'stop': return 'stop'; case 'length': return 'length'; case 'tool_calls': return 'tool_calls'; case 'content_filter': return 'content_filter'; default: return 'stop'; } }